Documentation
¶
Overview ¶
Package peers is a thin NATS JetStream client for the claude-peers v3 network. The durable event log (stream PEERS, subjects peers.>) is the source of truth; presence (KV PEERS_PRESENCE) is a projection. Every mutating op publishes an event, so any consumer subscribing peers.> sees the whole network. NATS does the queueing, durability and presence-TTL — this package only wires to it.
Index ¶
- Variables
- func ResolveIdentity(cwd, explicit string) (name, source string)
- func SanitizeName(s string) string
- func URLFromEnv() string
- type Client
- func (c *Client) Claim(ctx context.Context, p Peer) (*Peer, error)
- func (c *Client) ClaimWithFallback(ctx context.Context, p Peer) (Peer, error)
- func (c *Client) Close()
- func (c *Client) Consumers(ctx context.Context) ([]ConsumerStatus, error)
- func (c *Client) DeleteInbox(ctx context.Context, name string) error
- func (c *Client) Deregister(ctx context.Context, agent string) error
- func (c *Client) Heartbeat(ctx context.Context, agent, session string) error
- func (c *Client) NATS() *nats.Conn
- func (c *Client) OnReconnect(f func())
- func (c *Client) Peers(ctx context.Context) ([]Peer, error)
- func (c *Client) Register(ctx context.Context, p Peer) error
- func (c *Client) Send(ctx context.Context, m Message) (DeliveryStatus, error)
- func (c *Client) SetSummary(ctx context.Context, agent, summary string) error
- func (c *Client) Setup(ctx context.Context) error
- func (c *Client) Subscribe(ctx context.Context, agent string, h func(Message)) error
- func (c *Client) Watch(ctx context.Context, fromStart bool, h func(Envelope)) error
- type ConsumerStatus
- type DeliveryStatus
- type Envelope
- type EventType
- type Message
- type Peer
Constants ¶
This section is empty.
Variables ¶
var ErrInvalidName = errors.New("invalid agent name")
ErrInvalidName is returned for a recipient that cannot be a real agent.
var ErrNameLost = errors.New("agent name now held by another session")
ErrNameLost is returned by Heartbeat when the name is now held by a different session — e.g. this machine slept past the TTL and someone else claimed it. The caller must stop heartbeating this name: fighting over the key makes presence flap between machines and splits the durable inbox across two live sessions.
var ErrNameTaken = errors.New("agent name already held")
ErrNameTaken is returned by Claim when a live session already holds the name.
Functions ¶
func ResolveIdentity ¶
ResolveIdentity resolves this session's agent name, most explicit wins: explicit (--as flag) > CLAUDE_PEERS_AGENT > .claude-peers-agent file in cwd > sanitized basename of cwd. The dir default makes the common case zero-config: the agent working in ~/projects/pith is "pith" — which is how people already think about their sessions. Source is "flag", "env", "file", or "default".
func SanitizeName ¶
SanitizeName makes a string safe as a NATS subject token (names ride in peers.msg.<name>): lowercase, [a-z0-9_-] only, everything else collapses to a single '-', trimmed, capped at 32 chars. Returns "" if nothing survives.
func URLFromEnv ¶
func URLFromEnv() string
URLFromEnv resolves the server url the same way ConnectFromEnv does: NATS_URL env, then ~/.config/cp3/url, then localhost.
Types ¶
type Client ¶
type Client struct {
// contains filtered or unexported fields
}
Client wraps a NATS connection + JetStream + the presence KV.
func Connect ¶
Connect dials NATS. creds is a path to a .creds file and token is a plain auth token; either may be "" (empty both = no auth). creds wins if both are set.
func ConnectFromEnv ¶
ConnectFromEnv dials using NATS_URL (default nats://127.0.0.1:4222), NATS_CREDS (a .creds file) and a token — the single place auth is resolved so every binary authenticates the same way. Token resolution: NATS_TOKEN env, then the file named by NATS_TOKEN_FILE, then ~/.config/cp3/token. The file paths keep the secret out of argv, JSON config, and dotfile-synced shell rc.
func (*Client) Claim ¶
Claim registers p only if the name is free (or already held by p's own session). Returns ErrNameTaken + the current holder otherwise. ponytail: best-effort uniqueness via the presence KV — a dead holder's key expires in one TTL (30s), then the name frees. A tighter guarantee would need a lock stream; not worth it for a fleet of agents.
func (*Client) ClaimWithFallback ¶ added in v0.2.0
ClaimWithFallback claims p.Agent, falling back to a machine-qualified name when another live session holds it: "astrobot" taken -> "astrobot-macbook1". Deterministic and self-describing — unlike ordinal suffixes, the fallback says WHERE the twin lives (the usual cause: the same synced project dir open on two machines). If even that is taken (same dir twice on one machine), a short session tag breaks the tie. Returns the Peer actually registered.
func (*Client) Consumers ¶
func (c *Client) Consumers(ctx context.Context) ([]ConsumerStatus, error)
Consumers lists every consumer attached to the PEERS stream.
func (*Client) DeleteInbox ¶ added in v0.2.0
DeleteInbox removes a durable inbox consumer. Callers must check it holds no undelivered messages first — deleting a consumer discards its backlog.
func (*Client) Deregister ¶
Deregister removes presence and emits a deregister event.
func (*Client) Heartbeat ¶
Heartbeat refreshes presence ONLY while this session still owns the name: the record's session must match, and the write is a revision-CAS update so a racing claim can't be silently overwritten. session "" skips the ownership check (bare CLI use).
func (*Client) NATS ¶
NATS exposes the raw connection for advanced uses (e.g. the fleet-compat projector publishing legacy fleet.* events). Prefer the typed methods.
func (*Client) OnReconnect ¶ added in v0.2.0
func (c *Client) OnReconnect(f func())
OnReconnect registers f to run each time the underlying NATS connection is re-established (laptop wake, server bounce). Presence holders use it to re-claim immediately instead of waiting for the next heartbeat tick.
func (*Client) Peers ¶
Peers returns everyone currently present (KV projection = live view). A missing bucket means the network has never been set up here — that is an empty network, not an error.
func (*Client) Register ¶
Register records presence in the KV and emits a register event.
INVARIANT: on the roster means reachable. The inbox is created with the presence record, never after it, because the alternative shipped: cp3 register published presence with no consumer, so every message to that name landed in the log addressed to nobody while the sender was told "sent". Presence without a mailbox is not a peer, it is a decoy.
func (*Client) Send ¶
Send appends a message event to peers.msg.<to> and reports what will actually happen to it.
func (*Client) SetSummary ¶
SetSummary updates an agent's presence summary (re-puts the KV record and emits a presence event so consumers see the change).
func (*Client) Setup ¶
Setup ensures the PEERS stream (the log) and PEERS_PRESENCE KV exist. Idempotent — safe to call on every start.
func (*Client) Subscribe ¶
Subscribe delivers messages addressed to agent via a DURABLE consumer, so messages sent while offline drain on reconnect. Blocks until ctx is done; h is called for each message (auto-acked).
type ConsumerStatus ¶
type ConsumerStatus struct {
Name string
Pending uint64 // not yet delivered
AckPending int // delivered, unacked
Waiting int // outstanding pull requests: >0 means a session is attached right now
LastDelivery *time.Time
}
ConsumerStatus is a liveness snapshot of one consumer on the log. Pending piling up with no recent delivery = an abandoned durable (the v1 graveyard).
type DeliveryStatus ¶ added in v0.2.0
type DeliveryStatus string
DeliveryStatus is what actually happened to a sent message. Publishing always "succeeds" — the log accepts the write regardless — so reporting that as delivery is how a message can vanish while the sender is told it worked. These are three different facts and callers must see which one they got.
const ( // DeliveredLive: the recipient's inbox exists and a consumer is actively // waiting on it — the message goes to a running session now. DeliveredLive DeliveryStatus = "delivered" // Queued: the inbox exists but nothing is draining it. The message waits // durably and lands when that agent reconnects. This is the offline path // working as designed, not a failure. Queued DeliveryStatus = "queued" // NoInbox: NOTHING will ever receive this. No durable consumer exists for // the name, so the message sits in the log addressed to nobody. Presence // can still list the agent (see cmdRegister), which is exactly how this // stayed invisible. NoInbox DeliveryStatus = "no-inbox" )
func (DeliveryStatus) Human ¶ added in v0.2.0
func (d DeliveryStatus) Human(to string) string
Human returns a one-line description suitable for a CLI or a tool result.
type Envelope ¶
type Envelope struct {
V int `json:"v"`
ID string `json:"id"`
Type EventType `json:"type"`
TS int64 `json:"ts"` // unix millis
Actor string `json:"actor"`
Data json.RawMessage `json:"data"`
}
Envelope is the versioned wire contract every consumer builds against.
Directories
¶
| Path | Synopsis |
|---|---|
|
cmd
|
|
|
cp3
command
cp3 — CLI for the claude-peers v3 NATS-native network.
|
cp3 — CLI for the claude-peers v3 NATS-native network. |
|
internal
|
|
|
boot
Package boot connects to the peers network, bringing a local network up on demand: when the target is localhost and nothing is listening, it spawns a detached `cp3 serve` and retries.
|
Package boot connects to the peers network, bringing a local network up on demand: when the target is localhost and nothing is listening, it spawns a detached `cp3 serve` and retries. |
|
bridge
Package bridge is the shared skeleton for external-runtime adapters: connect to the network, claim presence, and inject each inbound peer message into the target runtime via a runtime-specific inject func.
|
Package bridge is the shared skeleton for external-runtime adapters: connect to the network, claim presence, and inject each inbound peer message into the target runtime via a runtime-specific inject func. |
|
codex
Package codex bridges the peer network into a codex session via the codex app-server (JSON-RPC over stdio, newline-delimited).
|
Package codex bridges the peer network into a codex session via the codex app-server (JSON-RPC over stdio, newline-delimited). |
|
mcp
cp3-mcp — the Claude Code injection adapter for claude-peers v3.
|
cp3-mcp — the Claude Code injection adapter for claude-peers v3. |
|
opencode
cp3-opencode — bridges the claude-peers NATS network to a running opencode server: each inbound peer message is injected as a steered prompt into one opencode session.
|
cp3-opencode — bridges the claude-peers NATS network to a running opencode server: each inbound peer message is injected as a steered prompt into one opencode session. |