What it is
The nats package implements the
NATS Agent Protocol v0.3
for Phero agents. It provides two top-level types:
-
nats.Server: registers anyagent.Agentas a NATS micro service. It handles discovery (via$SRV.PING/INFO), streaming prompt responses, periodic heartbeats, and on-demand status replies. The agent is wire-compatible with the TypeScript and Python SDKs in the synadia-agents repository. -
nats.Client: discovers compliant agents on the NATS bus and sends them prompts.Client.AsTool()wraps a remote agent as anllm.Toolthat any local Phero agent can call. -
nats.Resolver: combines discovery with heartbeat tracking for callers that stay running. It caches each agent's handle, reuses it only while the agent keeps sending heartbeats, and looks the agent up again after a failed prompt.
Unlike the A2A transport (HTTP-based), the NATS transport uses pub/sub messaging, which makes it well-suited for service meshes and environments where NATS is already running.
Example: registering an agent as a NATS server
See the full code in examples/nats-agent/server.
package main
import (
"context"
"log"
"os"
"os/signal"
"syscall"
natsgo "github.com/nats-io/nats.go"
"github.com/henomis/phero/agent"
"github.com/henomis/phero/llm/openai"
natsagent "github.com/henomis/phero/nats"
)
func main() {
nc, _ := natsgo.Connect(natsgo.DefaultURL)
defer nc.Drain()
llmClient := openai.New(os.Getenv("OPENAI_API_KEY"))
a, _ := agent.New(llmClient, "my-agent", "A helpful assistant.")
srv, _ := natsagent.New(nc, a, "alice", "demo")
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
// Serves until ctx is cancelled, then drains in-flight prompts.
if err := srv.Start(ctx); err != nil {
log.Println(err)
}
}
New(nc, agent, owner, name) registers the agent as a NATS micro service.
The owner and name values become the 4th and 5th tokens of the prompt subject
(agents.prompt.phero.<owner>.<name>).
Start(ctx) serves until the context is cancelled, then drains (see
Graceful shutdown). A Server serves once: calling
Start again after it has drained returns ErrServerStopped.
Graceful shutdown
Cancelling the context passed to Start shuts the server down. It does not abandon prompts
that are already running. The shutdown happens in this order:
- The prompt endpoint is unsubscribed, so NATS sends new prompts to another replica in the
agentsqueue group. A prompt already buffered on this instance gets a503 server is drainingreply instead of being served. - Heartbeats stop, so this instance no longer shows as live.
- Handlers that are already running keep their own context and finish their LLM call.
- If they are still running when most of the drain budget is used up, their context is cancelled. The last fifth of the budget is left for them to return.
WithDrainTimeout(d) sets the total budget for all of this (default 30 s). Set it to cover a
typical completion: those calls are already paid for. A zero or negative value cancels in-flight handlers
immediately. If handlers still have not returned when the budget runs out, Start returns
ErrDrainIncomplete.
To trigger the same shutdown from elsewhere, for example from a signal handler, call
srv.Drain(ctx). It is safe to call concurrently and more than once, and it also makes a
blocked Start return. The ctx you pass limits only how long that caller waits;
it does not limit the drain itself.
srv, _ := natsagent.New(nc, a, "alice", "demo",
natsagent.WithDrainTimeout(2*time.Minute),
)
sigCtx, stop := signal.NotifyContext(context.Background(), syscall.SIGTERM)
defer stop()
go func() {
<-sigCtx.Done()
shutdownCtx, cancel := context.WithTimeout(context.Background(), 3*time.Minute)
defer cancel()
if err := srv.Drain(shutdownCtx); err != nil {
log.Printf("drain: %v", err)
}
}()
_ = srv.Start(context.Background())
Identifiers, sessions and payload limits
The agent id, owner, name and session are inserted into NATS subjects, so New checks each
one with ValidateSubjectToken. Only [A-Za-z0-9_-] is allowed, up to 64
characters. ., * and > are rejected because they would change
which subjects the server subscribes to. For example, an owner of * would create a wildcard
subscription in the shared queue group. Invalid values return ErrInvalidSubjectToken.
ValidateSubjectToken is exported, so you can check identifiers when you load your config.
With WithSession(s), the agent registers under the instance name
InstanceName(name, s) (name-s). It advertises s itself as
metadata.session in its service info, heartbeats and status replies. To address agents by
session, use FilterBySession(s), which saves you from building the combined name yourself.
srv, _ := natsagent.New(nc, a, "acme", "worker", natsagent.WithSession("prod"))
// registers as agents.prompt.phero.acme.worker-prod, metadata.session = "prod"
agents, _ := c.Discover(ctx,
natsagent.FilterByOwner("acme"),
natsagent.FilterBySession("prod"),
)
WithMaxPayload sets the prompt size limit that the agent advertises (default "1MB").
New compares it with the max_payload of the NATS connection. If you did not set
a limit, the default is lowered to the connection's limit. If you set one higher than the connection
allows, New returns ErrInvalidMaxPayload. On the client side,
Prompt checks the prompt against both limits and returns ErrPayloadTooLarge
before publishing. Without that check, the failure would be the less clear
nats: maximum payload exceeded error.
Example: discovering and prompting a remote agent
See the full code in examples/nats-agent/client.
nc, _ := natsgo.Connect(natsgo.DefaultURL)
defer nc.Drain()
c := natsagent.NewClient(nc)
// Discover all compliant agents on the bus
agents, err := c.Discover(ctx)
if err != nil {
panic(err)
}
// Stream a prompt to the first discovered agent
stream, err := c.Prompt(ctx, agents[0], "What is NATS?")
if err != nil {
panic(err)
}
defer stream.Close()
text, _ := stream.Text(ctx)
fmt.Println(text)
Discover uses a stall strategy: it collects responses to a $SRV.INFO.agents
fan-out request for up to 750 ms of silence (configurable), capped by a 2-second absolute deadline.
Discovery options such as FilterByOwner, FilterByName and
FilterBySession narrow the results client-side.
Both calls respect ctx. If ctx has an earlier deadline than the discovery
timeout, Discover stops at that deadline. A cancelled ctx returns its own error
instead of ErrNoAgentsFound. Prompt returns immediately if ctx has
already ended, so no agent starts a prompt that nobody will read.
Wrapping a remote agent as a tool
Any discovered agent can be turned into an llm.Tool that a local Phero agent can call.
This is the building block for NATS-based multi-agent pipelines.
c := natsagent.NewClient(nc)
agents, _ := c.Discover(ctx, natsagent.FilterByName("researcher"))
tool, err := agents[0].AsTool("researcher", "Research a topic and return key facts.")
if err != nil {
panic(err)
}
orchestrator.AddTool(tool)
AgentHandle.AsTool(toolName, toolDesc) is shorthand for Client.AsTool(info, toolName, toolDesc).
The tool sends a JSON-encoded prompt to the remote agent's prompt subject and collects the streamed
response chunks before returning the full text to the LLM.
This tool keeps the one handle that discovery returned for as long as the tool exists. That works for a script that discovers, prompts and exits. For a process that stays up, use a Resolver instead: when the target restarts, it gets a new instance, and a tool holding the old handle keeps prompting an address that no longer exists.
Resolver: long-lived callers
Resolver looks up an agent by (owner, name) and returns a live
AgentHandle without running discovery on every call. It adds a
HeartbeatTracker to the Client:
- A cached handle is reused only while the agent is still sending heartbeats and the cache TTL (default
5 min, see
WithResolverCacheTTL) has not expired. - When several callers miss the cache for the same agent at the same time, only one discovery runs.
Invalidate(owner, name)removes a cached handle, so the nextResolveruns discovery again.Resolver.AsToollooks the agent up on every call and invalidates the handle when a prompt fails. If the target restarts between calls, the next call finds the new instance.
c := natsagent.NewClient(nc)
resolver, err := natsagent.NewResolver(c)
if err != nil {
panic(err)
}
defer resolver.Close()
// Optional: wait until the agents have sent a heartbeat. This uses no NATS traffic.
err = resolver.WaitReady(ctx, []natsagent.AgentKey{
{Owner: "newsroom", Name: "researcher"},
}, 30*time.Second, 0)
tool, err := resolver.AsTool("newsroom", "researcher",
"researcher", "Research a topic and return key facts.")
if err != nil {
panic(err)
}
orchestrator.AddTool(tool)
The heartbeat tracker starts when NewResolver returns, so create the resolver
before starting the agents whose first heartbeats WaitReady has to see.
Retrying: permanent vs transient errors
nats.Permanent(err) reports whether retrying cannot fix an error. It returns true for
requests the agent will always reject: an empty or oversized prompt, attachments the agent does not
accept, a malformed envelope, or a 4xx ServiceError. It also returns true for configuration
errors, such as a nil dependency or an empty or invalid identifier. It returns false for everything else,
including ErrNoAgentsFound, ErrStreamTimeout and context errors. Use it instead
of listing error kinds yourself, so new error kinds are classified correctly without changes on your side.
text, err := stream.Text(ctx)
if err != nil && !natsagent.Permanent(err) {
// worth retrying, perhaps against another replica
}
Multi-agent NATS pipeline
For a production-style example, see
examples/nats-agent/multi-agent.
It shows three specialised agents (researcher, writer, editor) each running as a NATS micro service,
coordinated by a local orchestrator that finds them through a Resolver
and calls each one as an llm.Tool. If a worker restarts during a run, the orchestrator finds
the new instance.
docker run --rm -p 4222:4222 nats
# Terminal 1
OPENAI_API_KEY=<key> go run ./examples/nats-agent/multi-agent/researcher
# Terminal 2
OPENAI_API_KEY=<key> go run ./examples/nats-agent/multi-agent/writer
# Terminal 3
OPENAI_API_KEY=<key> go run ./examples/nats-agent/multi-agent/editor
# Terminal 4 — orchestrator
OPENAI_API_KEY=<key> go run ./examples/nats-agent/multi-agent/orchestrator -topic "quantum computing"
API reference
Server
nats.New(nc, handler, owner, name, opts...) (*Server, error)— create a server;handleris anyHandler, e.g.*agent.Agentsrv.Start(ctx) error— serve until ctx is cancelled, then drain; returnsErrServerStoppedif the server has already drainedsrv.Drain(ctx) error— stop accepting prompts and wait for in-flight ones; safe to call concurrently and more than once
Server options
WithAgentID(id string)— themetadata.agentvalue and 3rd subject token (default:"phero")WithSession(session string)— register asInstanceName(name, session)and advertisemetadata.sessionWithVersion(v string)— the micro service version (default:"0.1.0")WithMaxPayload(s string)— advertised max prompt size, e.g."512KB"(default:"1MB", lowered to the connection's limit)WithAttachmentsOk(ok bool)— advertise attachment support in endpoint metadataWithHeartbeatInterval(d time.Duration)— how often to publish a heartbeat (default: 30 s)WithKeepaliveInterval(d time.Duration)— how often to send a keepalive chunk during a long response (default: 30 s)WithDrainTimeout(d time.Duration)— total shutdown budget for in-flight prompts (default: 30 s)
Client
nats.NewClient(nc, opts...) *Client— create a clientclient.Discover(ctx, opts...) ([]*AgentHandle, error)— discover agents on the busclient.Prompt(ctx, info, text) (*Stream, error)— send a prompt and get a streamclient.AsTool(info, toolName, toolDesc) (*llm.Tool, error)— wrap a remote agent as a tool
Client options
WithDiscoveryTimeout(d time.Duration)— absolute deadline for discovery (default: 2 s; discovery also stops after 750 ms with no new replies)WithInactivityTimeout(d time.Duration)— how long a stream waits for the next chunk (default: 60 s)
Discovery options
FilterByAgent(agent string)— keep only agents with matchingmetadata.agentFilterByOwner(owner string)— keep only agents with matchingmetadata.ownerFilterByName(name string)— keep only agents with a matching instance name (for sessioned agents, theInstanceName)FilterBySession(session string)— keep only agents with matchingmetadata.session
AgentHandle
handle.Prompt(ctx, text) (*Stream, error)— shorthand forClient.Prompthandle.AsTool(toolName, toolDesc) (*llm.Tool, error)— shorthand forClient.AsToolhandle.AgentInfo— embeddedAgentInfowith name, owner, session, subjects, etc.
Stream
stream.Text(ctx) (string, error)— collect all chunks into a single stringstream.Close()— unsubscribe from the reply subject
Resolver
nats.NewResolver(client, opts...) (*Resolver, error)— create a resolver and start its heartbeat trackerresolver.Resolve(ctx, owner, name) (*AgentHandle, error)— return the cached handle while the agent is live, otherwise run discoveryresolver.Invalidate(owner, name)— remove a cached handleresolver.AsTool(owner, name, toolName, toolDesc) (*llm.Tool, error)— a tool that looks the agent up on every callresolver.WaitReady(ctx, keys, timeout, poll) error— wait until everyAgentKeyhas sent a heartbeatresolver.Close() error— stop the heartbeat tracker (theClientstays open)WithResolverCacheTTL(d time.Duration)— how long a handle stays cached before discovery runs again (default: 5 min)
HeartbeatTracker
nats.NewHeartbeatTracker(nc) (*HeartbeatTracker, error)— subscribe toagents.hb.*.*.*tracker.IsAgentOnline(owner, name) bool/tracker.IsOnline(instanceID) bool— true if a heartbeat arrived within 3× its intervaltracker.Stop() error— unsubscribe
Helpers
nats.InstanceName(name, session) string— the instance name an agent registers under withWithSessionnats.ValidateSubjectToken(token) error— check that an identifier is usable as a subject tokennats.Permanent(err) bool— whether retrying cannot fix the error
Errors
ErrNilConn/ErrNilHandler/ErrNilClient— nil dependency passed toNeworNewResolverErrEmptyOwner/ErrEmptyName— blank owner or name passed toNewErrInvalidSubjectToken— an identifier contains characters that are not allowed in a subject tokenErrInvalidMaxPayload—max_payloadcannot be parsed, or is above the connection's limitErrEmptyPrompt— empty prompt passed toPromptErrPayloadTooLarge— prompt exceeds the agent's advertised or the connection's max payloadErrAttachmentsNotAllowed— the agent does not accept attachmentsErrMalformedEnvelope— the request payload is not a valid envelopeErrNoAgentsFound—Discoverfound no matching agentsErrStreamTimeout— no chunk arrived within the inactivity timeoutErrServiceError— the agent returned an error; useerrors.Aswith*ServiceErrorto readCodeandDescriptionErrDrainIncomplete— handlers were still running when the drain budget ran outErrServerStopped—Startcalled on a server that has already drained
Run the simple example
Start a NATS server first (plain core NATS — no JetStream needed):
docker run --rm -p 4222:4222 nats
# Terminal 1 — start the agent server
OPENAI_API_KEY=<key> go run ./examples/nats-agent/server -owner=alice -name=demo
# Terminal 2 — interactive client
go run ./examples/nats-agent/client
Interoperability
Phero agents are fully wire-compatible with the TypeScript and Python SDKs from synadia-agents. Any client speaking the NATS Agent Protocol can discover and prompt a Phero agent.
# Discover the running Phero agent with the Python SDK
uv run python examples/01-discover.py --url nats://127.0.0.1:4222
# Stream a single prompt
uv run python examples/02-prompt-text.py --url nats://127.0.0.1:4222 "What is 2+2?"
Related packages
- a2a: HTTP-based agent-to-agent transport — use when NATS is not available
- memory:
memory/natsstores conversation history in NATS JetStream KV;memory/nats.Open(nc, bucket, sessionID)creates the bucket if it does not exist, or binds the existing one - tool:
tool/kvgives an agentkv_read/kv_writetools over a JetStream KV bucket (kv.Open(nc, bucket), thenstore.Tools()) - agent: the agent type wrapped by both
ServerandClient.AsTool - llm:
AsTool()returns an*llm.Tool