relay

package
v0.2.5 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 3, 2026 License: Apache-2.0 Imports: 14 Imported by: 0

README

relay — agent messaging across tiers

services/relay delivers messages between agents. An agent is a named worker — versioned instructions plus a live conversation session — that holds assignments and receives messages at a typed agent address. A message from an agent is a proposal; it is never permission to execute anything.

Every address is a provider's type under ipc/model/<provider>, and the type alone fixes the tier the message travels:

Tier Meaning Address types today
in_process same runtime: a direct call, no capability copilot.InProcessAddress — a Workbench session on the adapter mounted in this runtime
host another process on this machine, through a unix socket or loopback claudecode.HostAddress — a Claude Code session and its local socket
network another machine, through ipc/net copilot.NetworkAddress — a Workbench session on another adapter

Each envelope records the widest tier between sender and recipient, and the relay exports csf_relay_envelopes_sent_total{tier} and csf_relay_envelopes_delivered_total{tier} so avoidable host and network crossings are measurable.

Mounting it (bare binary first)

The relay is a runtime service and owns no listener. A binary builds it, mounts it into its runtime.HostRuntime (or starts it on a runtime.Scope), and hands it to whatever needs it:

core, err := relay.NewRelay[*agentv1.AgentMessage](relay.WithMetrics(registry))
host.Mount("relay", core)                          // before anything that sends
service, err := csf.New(csf.WithAgentMessaging(core)) // MCP tools for host/network agents
messenger, err := relay.NewMessenger(core)            // the same contract, in-process

A relay carries exactly one message type, its contract: every sender, inbox, receiver and MCP tool on it agrees on that type at compile time, and nothing is erased or asserted between them. CSF's relay carries candace.agent.v1.AgentMessage (proto/candace/agent/v1/agent.proto): a kind, the text, an optional in_reply_to and an optional reference, each refined by Liquid Proto. relaytest.DescribeRelayConformance[Body] is the messaging contract any relay instantiation runs.

app/csf does exactly this, and binds every Workbench session's register_agent / send_agent_message tools through copilotadapter.SessionMessaging; a message to a Workbench agent becomes a queued prompt in its session.

Using it from a Claude Code session (host tier)

A Claude Code session is a host-tier agent: it calls CSF's authenticated agent MCP endpoint (csf.Service.AgentMCPHandler, mounted by app/csf at /mcp/agent when started with --agent-mcp-key-file). Every request carries three headers minted for that agent and session by the host that holds the signing key (csf.AgentMCPAuthenticator.AgentMCPHeaders): Authorization: Bearer v1.…, X-CSF-Agent-ID, X-CSF-Session-ID.

  1. Add the endpoint to the session's MCP configuration (for example a project .mcp.json entry of type http whose url is the loopback address and /mcp/agent path, with the three headers above).
  2. Register once per session, and again after every restart: RegisterAgentAddress {"provider": "claudecode", "endpoint": "127.0.0.1:14111"} (or "unix:/absolute/socket/path" when CSF serves a unix socket). Only a unix socket or a loopback address is accepted, so the registration is host-tier by construction. Messages queued while the session was away are kept.
  3. Send: SendAgentMessage {"to": "reviewer", "message": {"kind": 2, "text": "…"}}. Kind is 1 note, 2 request, 3 reply, 4 status; the message is refused unless it satisfies the contract. A Workbench agent reads it as a queued prompt.
  4. Poll: FetchAgentInbox {"limit": 16} returns unacknowledged envelopes oldest first, each with id, from, tier, message, sent_at; the same envelopes come back until AcknowledgeAgentInbox {"ids": ["…"]} removes them.

Limits

  • Inboxes and registrations are in memory: a process restart loses them. The durable registry, claims and checkpoints are slice A1; any IRegistry implementation must keep the semantics written on that interface.
  • Host and network agents pull; nothing pushes to them yet.
  • In-process delivery is at most once: an envelope taken by Receive whose prompt submission then fails is logged, not retried.
  • There is no command yet that mints agent MCP headers for an out-of-process agent; the host that holds the signing key must call AgentMCPHeaders.

Documentation

Overview

Package relay delivers messages between agents. An agent is a named worker (versioned instructions plus a live conversation session) that receives messages at a typed agent address; the relay holds one inbox per registered agent and moves typed envelopes into it.

It is named for what it does: it relays an envelope from a sender's address to a recipient's inbox and records how far the envelope travelled. It does not decide what an agent does with a message — a message from an agent is a proposal, never permission to execute anything — and it is not a durable queue: inboxes live in memory until the durable registry, claims and checkpoints arrive (slice A1, stored through csfpg after S4).

Tiers are decided by address types

Every agent address is a provider's type from ipc/model/<provider>, and the type alone fixes its tier (see model.IAgentAddress): in_process for a session in this runtime, host for another process on this machine, network for another machine. Every Envelope records the widest tier its two addresses reach, and the relay counts envelopes per tier so avoidable crossings are measurable.

The relay itself holds no capability and opens no socket. In-process agents exchange envelopes by direct calls into their inboxes (Messenger.Send and Messenger.Receive); a host or network agent reaches the same inbox through CSF's per-agent authenticated MCP endpoint, whose listener the binary grants from candace/ipc/net, and pulls it with Messenger.Fetch and Messenger.Acknowledge.

Goroutines

The relay is a runtime service: Relay.Start receives its scope, and every goroutine it needs — the directory's owner and one owner per inbox, each a candace/pkg/mailbox — starts through that scope. Nothing here calls go.

Index

Constants

View Source
const (
	// MetricEnvelopesSent counts envelopes accepted into an inbox, by tier.
	MetricEnvelopesSent = "csf_relay_envelopes_sent_total"
	// MetricEnvelopesDelivered counts envelopes handed to an in-process
	// receiver or acknowledged by a host or network agent, by tier.
	MetricEnvelopesDelivered = "csf_relay_envelopes_delivered_total"
	// MetricTierLabel is the label both counters carry: in_process, host or
	// network, as ipc.Tier names them.
	MetricTierLabel = "tier"
)

Variables

View Source
var (
	// ErrInvalidRegistration reports a registration without a valid agent
	// identifier or address.
	ErrInvalidRegistration = errors.New("relay: invalid registration")
	// ErrUnknownAgent reports an agent that has never registered.
	ErrUnknownAgent = errors.New("relay: agent is not registered")
	// ErrUnknownAddress reports an address no agent currently holds — never
	// registered, or replaced when its agent re-registered after a restart.
	ErrUnknownAddress = errors.New("relay: no agent is registered at that address")
	// ErrAddressTaken reports an address another agent already holds.
	ErrAddressTaken = errors.New("relay: address is registered to another agent")
	// ErrNotStarted reports use of a relay whose runtime has not started it.
	ErrNotStarted = errors.New("relay: not started")
	// ErrStopped reports use of a relay whose scope has ended.
	ErrStopped = errors.New("relay: stopped")
)

Functions

This section is empty.

Types

type AgentID

type AgentID string

AgentID names an agent: the named worker, not its current session.

func (AgentID) Validate

func (agent AgentID) Validate() error

Validate reports whether the identifier matches the agent grammar.

type Envelope

type Envelope[Body any] struct {
	ID        uuid.UUID
	FromAgent AgentID
	From      model.IAgentAddress
	ToAgent   AgentID
	To        model.IAgentAddress
	// Tier is the widest tier the two addresses reach: the boundary this
	// envelope crosses between sender and recipient.
	Tier   ipc.Tier
	Body   Body
	SentAt time.Time
}

Envelope is one message from one agent to another. Body is the sender's payload; for an in-process recipient it is handed over by value without encoding, so a pointer inside it is shared rather than copied — ownership of anything it references passes to the recipient on delivery.

type IMessenger

type IMessenger[Body any] interface {
	// Send delivers body from the agent registered at from to the agent
	// registered at to, and returns the envelope as queued.
	Send(ctx context.Context, from model.IAgentAddress, to model.IAgentAddress, body Body) (Envelope[Body], error)
	// Receive takes the agent's oldest envelope, waiting until one arrives.
	Receive(ctx context.Context, agent AgentID) (Envelope[Body], error)
}

IMessenger is the agent messaging contract: send a typed envelope to an agent address, and receive from an inbox. The addresses' types decide the tier; Body is the message type both sides agree on.

type IRegistry

type IRegistry interface {
	Register(ctx context.Context, registration Registration) error
	Resolve(ctx context.Context, agent AgentID) (Registration, error)
	Locate(ctx context.Context, addressKey string) (Registration, error)
}

IRegistry records which address each agent's live session is at. The relay consults it to route; it never decides placement or ownership.

MemoryRegistry is the implementation until slice A1 provides the durable one (stored through csfpg after S4). A durable implementation must keep these semantics:

  • Register is an upsert keyed by agent: re-registering replaces the agent's address, and the old address stops resolving at once.
  • One address belongs to at most one agent; a second agent registering a held address fails with ErrAddressTaken.
  • Resolve and Locate read the latest committed registration and report ErrUnknownAgent and ErrUnknownAddress respectively.
  • Each address round-trips: a registration read back must carry a value of the same provider address type, equal to the one registered, so its tier is unchanged by storage.

A registry needs no coordination with inboxes: the relay opens an agent's inbox before it calls Register, so every registration a registry commits already has an inbox. A failed Register leaves no registration behind.

type MemoryRegistry

type MemoryRegistry struct {
	// contains filtered or unexported fields
}

MemoryRegistry keeps registrations in this process. Reads take the current immutable snapshot; a write copies it and publishes the copy, so readers never wait and registration — rare next to routing — pays for the copy.

func NewMemoryRegistry

func NewMemoryRegistry() *MemoryRegistry

NewMemoryRegistry returns an empty in-memory registry.

func (*MemoryRegistry) Locate

func (registry *MemoryRegistry) Locate(ctx context.Context, addressKey string) (Registration, error)

Locate returns the registration currently holding addressKey.

func (*MemoryRegistry) Register

func (registry *MemoryRegistry) Register(ctx context.Context, registration Registration) error

Register upserts the agent's registration.

func (*MemoryRegistry) Resolve

func (registry *MemoryRegistry) Resolve(ctx context.Context, agent AgentID) (Registration, error)

Resolve returns the agent's current registration.

type Messenger

type Messenger[Body any] struct {
	// contains filtered or unexported fields
}

Messenger is a Relay's agent-facing view. It shares its relay's message type: a relay carries one contract, and every messenger over it speaks it.

func NewMessenger

func NewMessenger[Body any](relay *Relay[Body]) (*Messenger[Body], error)

NewMessenger returns a messenger over relay.

func (*Messenger[Body]) Acknowledge

func (messenger *Messenger[Body]) Acknowledge(ctx context.Context, agent AgentID, identifiers []uuid.UUID) (int, error)

Acknowledge removes fetched envelopes from the agent's inbox and reports how many it removed; identifiers already acknowledged are ignored.

func (*Messenger[Body]) Fetch

func (messenger *Messenger[Body]) Fetch(ctx context.Context, agent AgentID, limit int) ([]Envelope[Body], error)

Fetch returns up to limit of the agent's unacknowledged envelopes, oldest first, without removing them. It is the host or network agent's side: what it fetched is redelivered until it acknowledges it.

func (*Messenger[Body]) Receive

func (messenger *Messenger[Body]) Receive(ctx context.Context, agent AgentID) (Envelope[Body], error)

Receive takes the oldest envelope; see IMessenger.Receive. It is the in-process agent's side: the handoff is the delivery.

func (*Messenger[Body]) Send

func (messenger *Messenger[Body]) Send(ctx context.Context, from model.IAgentAddress, to model.IAgentAddress, body Body) (Envelope[Body], error)

Send delivers body; see IMessenger.Send.

type Option

type Option func(configuration *relayConfiguration) error

Option configures a Relay before NewRelay builds anything.

func WithClock

func WithClock(now func() time.Time) Option

WithClock replaces the clock that stamps envelopes.

func WithInboxIdleTimeout

func WithInboxIdleTimeout(idle time.Duration) Option

WithInboxIdleTimeout is how long an inbox with nothing queued, no receiver waiting and no call in progress keeps its goroutine. The inbox and everything it holds outlive the goroutine; the next message starts it again. Optional; the default is one minute.

func WithMetrics

func WithMetrics(registerer prometheus.Registerer) Option

WithMetrics registers the per-tier envelope counters with registerer.

func WithRegistry

func WithRegistry(registry IRegistry) Option

WithRegistry supplies the registration store. Optional; the default is a MemoryRegistry.

type Registration

type Registration struct {
	Agent   AgentID
	Address model.IAgentAddress
}

Registration binds an agent to the address its live session receives messages at. An agent re-registers after a restart: its new session's address replaces the old one and its inbox, with everything still queued in it, carries over.

func (Registration) Kind

func (registration Registration) Kind() ipc.Tier

Kind is the registration's tier, which its address type decides.

func (Registration) Validate

func (registration Registration) Validate() error

Validate reports whether the registration names a valid agent and address.

type Relay

type Relay[Body any] struct {
	// contains filtered or unexported fields
}

Relay is the core every Messenger shares, typed by the one message contract it carries: Body is the type every sender, inbox and receiver on this relay agrees on, so nothing is erased or asserted between them. It is a runtime service: mount it into the host runtime before anything that sends.

func NewRelay

func NewRelay[Body any](options ...Option) (*Relay[Body], error)

NewRelay validates its options and builds a relay that has not started.

func (*Relay[Body]) Register

func (relay *Relay[Body]) Register(ctx context.Context, registration Registration) error

Register opens the agent's inbox if it has none and then binds the agent to its address. The inbox comes first so a committed registration always has an inbox to deliver to: if the registry write fails, the agent is left with an inbox and no registration, which nothing routes to and which the next successful Register reuses. Re-registering after a restart keeps the inbox and everything queued in it.

func (*Relay[Body]) Resolve

func (relay *Relay[Body]) Resolve(ctx context.Context, agent AgentID) (Registration, error)

Resolve returns the agent's current registration.

func (*Relay[Body]) Start

func (relay *Relay[Body]) Start(scope *runtime.Scope) error

Start runs the directory on the relay's scope, and a stopper that retires the directory when the scope is canceled. Each inbox is a lazy service mounted on the same scope when its agent first registers: its goroutine starts on a child scope when the inbox is first used and is retired after the inbox idle timeout, and the scope's cancellation stops every inbox.

Directories

Path Synopsis
Package relaytest holds the conformance specs every relay.IRegistry implementation must pass, and the messaging specs every relay.Relay[Body] must pass over such a registry.
Package relaytest holds the conformance specs every relay.IRegistry implementation must pass, and the messaging specs every relay.Relay[Body] must pass over such a registry.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL