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
- Variables
- type AgentID
- type Envelope
- type IMessenger
- type IRegistry
- type MemoryRegistry
- type Messenger
- func (messenger *Messenger[Body]) Acknowledge(ctx context.Context, agent AgentID, identifiers []uuid.UUID) (int, error)
- func (messenger *Messenger[Body]) Fetch(ctx context.Context, agent AgentID, limit int) ([]Envelope[Body], error)
- func (messenger *Messenger[Body]) Receive(ctx context.Context, agent AgentID) (Envelope[Body], error)
- func (messenger *Messenger[Body]) Send(ctx context.Context, from model.IAgentAddress, to model.IAgentAddress, ...) (Envelope[Body], error)
- type Option
- type Registration
- type Relay
Constants ¶
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 ¶
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.
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 ¶
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 WithInboxIdleTimeout ¶
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 ¶
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 (*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]) Start ¶
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.
Source Files
¶
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. |