NATS

Connect agents over NATS pub/sub with the NATS Agent Protocol.

What it is

The nats package implements the NATS Agent Protocol v0.3 for Phero agents. It provides two top-level types:

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:

  1. The prompt endpoint is unsubscribed, so NATS sends new prompts to another replica in the agents queue group. A prompt already buffered on this instance gets a 503 server is draining reply instead of being served.
  2. Heartbeats stop, so this instance no longer shows as live.
  3. Handlers that are already running keep their own context and finish their LLM call.
  4. 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:

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

Server options

Client

Client options

Discovery options

AgentHandle

Stream

Resolver

HeartbeatTracker

Helpers

Errors

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