8
9 Implementation
lead edited this page 2026-09-22 12:48:50 +07:00

9. Implementation Architecture (as-built)

Scope. This is the as-built code map of the fleetd module: the real packages, classes, flows, and state machines in the source tree. It is a companion to the design-level 1. Architecture and 2. Message Server pages. Every class name, enum value, and route below was checked against the source at main a49e968 (2026-08-31). Where this page and the code ever disagree again, trust the code, not this page.

fleetd is a single-host Java 25 / Maven daemon. It is the one gateway between an on-subscription primary (a lead, running Opus) and off-subscription members (workers and architects), speaking to the herdr PTY manager over a Unix-domain socket. It presents two equivalent faces — a REST server (rest.FleetApp) and an MCP server (mcp.FleetMcp) — over one shared service core (msg). The Java package root is dev.ltms.fleet.

Component map

flowchart TB
    lead["Lead (primary)<br/>on-subscription"]
    peerlead["Peer lead<br/>another daemon"]
    worker["Worker / architect<br/>off-subscription"]

    subgraph fleetd["fleetd daemon"]
        direction TB
        subgraph faces["North faces (thin adapters)"]
            rest["rest.FleetApp<br/>REST · Javalin"]
            mcp["mcp.FleetMcp<br/>MCP · /mcp servlet"]
        end
        authz["auth.CallerResolver + auth.Authz<br/>identity + the one decision table"]
        msg["msg.MessageService + Rendezvous<br/>+ ReplyInbox + ReplyPushLoop<br/>service core"]
        inject["inject.Injector + StatusPoller<br/>+ CompletionResolver<br/>single-writer, status-gated delivery"]
        health["health.FleetHealthMonitor<br/>slow whole-fleet evidence loop"]
        session["session.SessionManager<br/>+ GitWorktrees<br/>member registry + worktrees"]
        peerspi["peer.PeerLauncher SPI<br/>member.CompositePeerLauncher"]
        adapters["member.ClaudeCodeLauncher<br/>member.OpenCodeLauncher"]
        guard["guard.SubscriptionGuard<br/>boundary enforcement"]
        herdr["herdr.AgentControl / WorkspaceControl<br/>via herdr.HerdrRouter"]
        leadmsg["msg.LeadChannel / LeadMailbox<br/>lead ↔ lead over a shared broker"]
    end

    herdrd["herdr daemon<br/>(PTY manager)"]

    lead -->|"fleet_send / fleet_spawn / answer"| rest
    lead -->|"MCP tool calls"| mcp
    worker -->|"fleet_reply / fleet_ask"| mcp
    rest --> authz
    mcp --> authz
    authz --> msg
    mcp --> session
    session --> peerspi
    peerspi --> adapters
    adapters --> guard
    adapters --> herdr
    msg --> inject
    health --> session
    health --> msg
    inject --> herdr
    herdr --> herdrd
    herdrd -.->|"drives PTYs"| worker
    mcp --> leadmsg
    leadmsg <-.->|"shared broker, cross-host"| peerlead

    classDef core fill:#2b6cb0,stroke:#1a365d,color:#ffffff;
    classDef gate fill:#b7791f,stroke:#7b341e,color:#ffffff;
    class msg,inject core
    class guard,authz gate

Figure 1 — component layers. Both north faces resolve identity through auth before reaching the service core (msg); delivery to a member reaches it only via inject → herdr. Spawning a member is delegated through the peer.PeerLauncher SPI, today implemented by member.CompositePeerLauncher routing to member.ClaudeCodeLauncher or member.OpenCodeLauncher; the subscription guard (amber) lives inside those adapters and gates spawn. health watches the same state the other layers produce but never drives delivery itself. msg.LeadChannel is the separate cross-host path a lead uses to reach a peer lead directly (see 10. Cross-Host Messaging).

Package & class reference

Sixteen packages under dev.ltms.fleet: auth, config, guard, health, herdr, inject, lead, logging, mcp, member, metrics, msg, peer, placement, rest, session. Below, per package: what it owns and its few most important classes, each with a file:line you can check yourself. This is not every class in every package — it is the ones a maintainer needs to find first.

herdr — the wire layer

The only speaker of the herdr protocol. Encodes newline-delimited JSON-RPC and projects herdr's workspace/tab/pane/agent nodes into typed Java records.

Class Role
HerdrClient (herdr/HerdrClient.java:1) The client face — one call(method, params).
UnixSocketHerdrClient (herdr/UnixSocketHerdrClient.java:1) JDK Unix-socket transport; fresh connection per call.
AgentControl (herdr/AgentControl.java:25) Domain wrapper over agent.*: start/send/read/status/close a member's agent.
WorkspaceControl (herdr/WorkspaceControl.java:1) Placement over workspace.*/tab.*.
AgentStatus (herdr/AgentStatus.java:8) A member's lifecycle state — see the state machine below.
HerdrRouter (herdr/HerdrRouter.java:7) Routes a call to the lead herdr daemon or the separate member herdr daemon (CB-185's memberHerdrSocket), by target id.
LeadTabScanner (herdr/LeadTabScanner.java:16) Discovers which panes host a lead by scanning herdr for operator-labelled tabs, feeding auth.CallerResolver.
PaneLocator (herdr/PaneLocator.java:1) Resolves a PID to a herdr terminal_id.

AgentStatus.injectable() is only ever IDLE, BLOCKED, or DONE — never WORKING.

inject — status-gated delivery and turn lifecycle

The single-writer delivery layer. Polls a member's status, queues per target, injects only when safe, and synthesizes turn boundaries so a blocked fleet_send resolves even when a member never calls fleet_reply.

Class Role
Injector (inject/Injector.java:48) Queues and delivers at most one message per turn, serialized per target.
StatusPoller (inject/StatusPoller.java:21) Virtual-thread loop sampling active members; feeds Injector.onStatus.
StatusRefiner (inject/StatusRefiner.java:26) Reclassifies UNKNOWN by reading the pane's real content.
CompletionResolver (inject/CompletionResolver.java:1) TurnListener that scrapes the pane to resolve Rendezvous on turn completion, failure, or usage-limit exhaustion.
TurnListener (inject/TurnListener.java:13) Turn-lifecycle callbacks: onDelivered, onTurnComplete, onTurnFailed.
MemberPresence (inject/MemberPresence.java:18) Tracks which members have connected the bridge MCP — the real readiness signal, not herdr's idle. Imported by FleetMcp (mcp/FleetMcp.java:12, used at 465-469). This is the class the wiki used to call WorkerPresence; that name is gone.
ExhaustedPatternLookup, ExhaustionSink Seams for the CB-578 usage-limit classification (see Outcome.BACKEND_EXHAUSTED below).

msg — the service core

Owns the forward rendezvous (fleet_send → fleet_reply), the reverse rendezvous (fleet_ask → answer), async fire-and-poll, the reply inbox that holds a member's reply when no send is open, the push loop that nudges the primary to drain it, and — new since the wiki last described this page — the lead-to-lead channel.

Class Role
MessageService (msg/MessageService.java:46) Orchestrates send/reply/ask/answer + async dispatch.
Rendezvous (msg/Rendezvous.java:26) Low-level registry of forward waiters and reverse-ask futures.
ReplyInbox (msg/ReplyInbox.java:19) The port: publish, peek, ack, plus own/release for federation.
InMemoryReplyInbox / AmqpReplyInbox The two adapters — soft-state, and cross-restart durable over AMQP.
ReplyPushLoop (msg/ReplyPushLoop.java:50) Status-gated loop that nudges a lead's own pane about queued replies, finished tickets, and open questions.
LeadChannel (msg/LeadChannel.java:22) This daemon's lead-to-lead mailbox port: publish, peek, ack, selfCoordId.
LeadMailbox The production LeadChannel implementation, over a shared AMQP broker.
LeadCoordLoop (msg/LeadCoordLoop.java:15) Takes what arrived in this daemon's own mailbox and types it into the local lead's pane.
LeadHeartbeatLoop (msg/LeadHeartbeatLoop.java:20) Opt-in nudge for a lead that has sat idle too long with nothing driving it (CB-551).

MessageService.Outcome (msg/MessageService.java:65): REPLIED, COMPLETED_UNREPLIED, WORKER_FAILED, BACKEND_EXHAUSTED, QUESTION, TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY, STALE_TURN. This is nine values, not the six the wiki used to list — BACKEND_EXHAUSTED (a backend refused on a usage limit, CB-578) and QUESTION/STALE_TURN (the fleet_ask reverse rendezvous) were missing before.

MessageService.Phase (an async ticket's lifecycle, msg/MessageService.java:141): PENDING, ASKING, DONE, FAILED. ASKING is new: it means the member paused mid-turn in fleet_ask and fleet_poll shows the open question instead of a stale "still pending".

Rendezvous.Kind (msg/Rendezvous.java:29): REPLY, COMPLETION, FAILED, BACKEND_EXHAUSTED, QUESTION.

MessageService.AskOutcome: ANSWERED, NO_WAITER, TIMED_OUT.

Ticket retention (issue #197, fixed 2026-08-31). MessageService.pruneTerminalTickets (msg/MessageService.java) drops a finished ticket once TICKET_TTL_NANOS (msg/MessageService.java:62, 10 minutes) has passed since it completed. Task stamps completedNanos from a whenComplete hook registered in its constructor, so every completion path stamps it — a reply, the completion fallback, a timeout, a failure, or an abandon on teardown.

Until 2026-08-31 the age was measured from createdNanos instead, so the window to collect was ten minutes minus however long the task ran. Any delegation lasting longer than the TTL had its reply destroyed on the first sweep after it finished, and the reply lives only in Task.future, so nothing could recover it. Worth knowing when reading older behaviour reports, and when running a daemon that has not been redeployed since the fix.

auth — who is calling, and what they may do

Every request arrives on a connection, and this package turns that connection into a Principal and then decides. Identity is never taken from a tool argument.

Class Role
Role (auth/Role.java:12) PRIMARY, WORKER, ARCHITECT, ANONYMOUS.
Principal (auth/Principal.java:16) (role, terminal, pid, name) plus factories anonymous, primary, leader, worker, architect.
CallerResolver (auth/CallerResolver.java:42) Connection → Principal. Resolution order: a pane named by leaders:/the tab scanner ⇒ PRIMARY; a pane bound to an architect slot ⇒ ARCHITECT; any other on-host pane ⇒ WORKER; otherwise a bearer token (token mode) or loopback trust ⇒ PRIMARY; otherwise ANONYMOUS.
Authz (auth/Authz.java:12) The one decision table, permits(caller, action, targetSession).
MemberRegistry (auth/MemberRegistry.java:35) The architect-slot registry: configured slots plus the live terminal → slot bindings an architect spawn creates.
AuditLog (auth/AuditLog.java:22) One JSON line per allow/deny decision — never the message content.

Authz.Action (auth/Authz.java:18): SPAWN, STOP, SEND, REPLY, ASK, DRAIN, READ, METRICS — eight actions, not the five the wiki used to group them into.

The decision table (auth/Authz.java:44):

Action Permitted to
SPAWN, STOP, DRAIN the primary only
SEND the primary or an architect
REPLY, ASK the caller that owns the target session — only ever itself
READ, METRICS primary, worker, or architect

An architect gets SEND but never SPAWN/STOP/DRAIN: it can delegate turns to workers, but fleet lifecycle stays the primary's alone.

peer — the launcher SPI

The seam that keeps the core peer-neutral. The core delegates spawn/teardown to a PeerLauncher implementation and knows nothing about how a peer's process is actually built.

Class Role
PeerLauncher (peer/PeerLauncher.java:22) SPI: spawn, stop, effectiveCwd, parityOverlay, profiles, defaultProfile, list, reapOrphanWorkers, capabilities, capabilitiesFor, clearContext.
PeerHandle (peer/PeerHandle.java:12) Returned by spawn: id() is the registry/routing key, terminalId() the transport session id.
SpawnRequest (peer/SpawnRequest.java:16) Spawn parameters, including role (CB-557: which contract) and resumeSessionId.
Capability (peer/Capability.java:9) Declared adapter capabilities: MID_TURN_ASK, SELF_PR, WORKTREE, CONTEXT_RESET, ORPHAN_REAP, SESSION_NAME, SESSION_RESUME.
MemberRole (peer/MemberRole.java:24) What a member is for: ARCHITECT, DEV, REVIEWER — independent of which backend (profile) it runs on.
CharterReceipt (peer/CharterReceipt.java:21) A SHA-256 fingerprint of the exact charter text a member was launched with, for audit — never the charter text itself.

Shipped kinds are claude-code and opencode only (config/FleetConfig.java:333,335). AgentAPI was never built — no adapter for it exists in this source tree.

member — the peer adapters

Class Role
HerdrPeerLauncher (member/HerdrPeerLauncher.java:63) Abstract base for a peer materialized as a herdr agent — tab/pane placement, the CB-306 spawn-readiness gate, unique naming, orphan reap, teardown, listing, cwd resolution. Everything transport-shared lives here.
ClaudeCodeLauncher (member/ClaudeCodeLauncher.java:43) The Claude Code adapter: builds the worker env (ANTHROPIC_BASE_URL), asserts the subscription guard, mounts the bridge MCP inline.
OpenCodeLauncher (member/OpenCodeLauncher.java:54) The opencode adapter: no subscription guard (opencode is not on the subscription), writes an ephemeral opencode.json instead of inline flags, model as a -m flag.
CompositePeerLauncher (member/CompositePeerLauncher.java:69) The PeerLauncher the core actually holds when more than one adapter is configured — a router in front of one HerdrPeerLauncher per peer kind, routing by profile, by pane id, or fanning out fleet-wide.
MemberEnvAllowList / EnvAllowListScrub The credential-guard allow-list a spawned member's environment is scrubbed against (CB-596/CB-611 family).

mcp — the MCP north face

Exposes fleetd as a Streamable-HTTP MCP endpoint and resolves caller identity from the connection, never from tool arguments.

Class Role
FleetMcp (mcp/FleetMcp.java:67) Builds the MCP server, registers the fleet_* tools, holds the thin tool adapters.
ConnectionIdentity (mcp/ConnectionIdentity.java:16) Resolves who is calling from the OS peer PID and herdr's pane map — unforgeable.
PrimaryRegistry (mcp/PrimaryRegistry.java:22) Tracks the primary's terminal id, and (CB-532) which lead delegated to which worker, so a reply nudge reaches the right lead.
LsofPeerPidLookup / LsofProcessCwdLookup lsof-based OS lookups behind the identity resolution.

The registered tool set is exactly eleven tools (mcp/FleetMcp.java:301-326): fleet_send, fleet_reply, fleet_ask, fleet_status, fleet_poll, fleet_ack, fleet_spawn, fleet_list, fleet_stop, fleet_profiles, fleet_whoami. There is no fleet_read tool. A fleet_send carrying a coordId (instead of a sessionId/turnId) routes to a peer lead on another daemon through msg.LeadChannel, not to a worker session — see sendToLead (mcp/FleetMcp.java:616-642).

session — the member registry and worktrees

Class Role
SessionManager (session/SessionManager.java:42) Authoritative in-daemon registry of every member this process spawned. Implements TurnListener so the injector's turn boundaries drive SPAWNING → READY → BUSY → DONE/FAILED. One-shot: a finished member is torn down, never reused.
MemberSession (session/MemberSession.java:35) The immutable record: pane id, terminal id, profile, MemberRole, cwd, owner, state, worktree/branch, CharterReceipt, agentSessionId.
GitWorktrees (session/GitWorktrees.java:34) Production Worktrees implementation, shelling git. Neutralizes .mcp.json and opencode's repo-level config in every worktree so a member never inherits the primary's credentialed MCP mounts.
Worktrees (session/Worktrees.java:7) The seam: add, remove, hasUncommitted, overlayParity, repoRoot, plus snapshot/wipRefs/pruneWipRefs for the CB-586 refs/wip/* retention sweep.
SessionReaper (session/SessionReaper.java:13) Virtual-thread loop tearing down idle READY/DONE sessions past their TTL, and sweeping stale refs/wip/* snapshots.

health — fleet health evidence (new package)

Not in the daemon the wiki last described this page against. A pure classifier plus a slow collection loop, deliberately separate from the delivery poller in inject.

Class Role
FleetHealth (health/FleetHealth.java:11) Pure classifier: HealthSnapshot + prior tick's HealthPrior → HealthDecision. No I/O.
HealthState (health/HealthState.java:4) STARTING, IDLE, WORKING, WORK_PENDING, BLOCKED_AMBIGUOUS, NEVER_READY, GONE, TURN_BOUNDARY_LOST, ERROR_ON_SCREEN, STALL_SUSPECTED, MUTE, REPLY_STRANDED, DELEGATION_ORPHANED, CONTROL_LINK_DOWN.
FleetHealthMonitor (health/FleetHealthMonitor.java:23) The scheduled loop that collects a HealthSnapshot per member and calls FleetHealth.decide.
MuteCounter (health/MuteCounter.java:11) Counts turns that ended via the completion fallback instead of a structured fleet_reply — an observation, never a classifier state on its own.

ERROR_ON_SCREEN is declared but not decided yet — the class javadoc says it needs a bounded pane read and an adapter-specific fatal signature that do not exist yet.

lead — starting configured leads

Class Role
LeadLauncher (lead/LeadLauncher.java:50) Starts the leads fleet.leaders: declares, when none is already running. Deliberately does not reuse HerdrPeerLauncher: a lead must never get the reply charter, the idle reaper, or ANTHROPIC_BASE_URL — all three are true of every member spawn path.

placement, metrics, guard, config, logging, rest — the edge

Class Role
placement.PlacementPolicy (placement/PlacementPolicy.java:7) How an unqualified spawn picks a profile: fixed (the historical default-profile behaviour) or weighted (smooth weighted round-robin with maxLoad gating).
placement.BackendQuarantine (placement/BackendQuarantine.java:25) Keyed by credential id, not profile: puts a credential on cooldown after a BACKEND_EXHAUSTED classification (CB-578 stage B), so a fresh spawn does not walk straight back onto an account that just refused.
metrics.Metrics (metrics/Metrics.java:25) The registry and Prometheus text renderer — dependency-free by design.
metrics.FleetMetrics (metrics/FleetMetrics.java:21) The named series: fleet_sends_total, fleet_replies_total, fleet_push_nudges_total, fleet_lead_heartbeat_nudges_total, fleet_spawns_total, fleet_herdr_calls_total, fleet_auth_failures_total, fleet_sessions, fleet_inbox_depth.
guard.SubscriptionGuard (guard/SubscriptionGuard.java:23) assertWorker (the member's base_url must be on the off-subscription allowlist) and assertPrimaryClean (the primary's own env must carry none).
config.FleetConfig (config/FleetConfig.java:81) The YAML config record: memberHerdrSocket (CB-185, optional second herdr daemon for members), Coordinator (CB-637 lead-to-lead broker), Leader, Fleet (the architects/developers/reviewers pools).
logging.McpCancelledNotificationFilter Suppresses one noisy SDK warning for a normal MCP cancellation notice — cosmetic, not a control.
rest.FleetApp (rest/FleetApp.java:46) Javalin routes; validates bodies, maps Outcome to an HTTP status.
Fleetd (Fleetd.java) Static main — wires every real collaborator above and starts both faces.

REST routes

The acceptance surface — every capability MCP exposes is reachable here too, without Claude in the loop (rest/FleetApp.java:143-158). There is no /events route, and the worker endpoints are /members, not /workers.

Method + path Purpose
GET /healthz Liveness. Checks both herdr daemons when memberHerdrSocket is configured, and reports protocolMismatch if their protocol versions differ.
GET /metrics Prometheus text exposition (present only when a metrics registry is wired).
GET /sessions herdr workspace list, merged across both daemons.
GET /agents herdr-tracked agents.
GET /members Registry roster + live herdr status (CB-304).
GET /profiles Configured backend profiles + the default.
POST /members Spawn a member (?role=&profile=&cwd=&worktree=&ticket=, or a JSON body).
DELETE /members/{paneId} Stop a member.
POST /sessions/{id}/message fleet_send — blocking, wait:false, or an answer via turnId.
POST /sessions/{id}/reply fleet_reply (member).
GET /sessions/{id}/replies Drain a member's held replies from the inbox (peek + ack).
POST /sessions/{id}/ask fleet_ask (member → primary).
GET /sessions/{id}/status fleet_status, plus ready (the injector's own readiness gate).
GET /tasks/{ticket} Poll an async (wait:false) send.

Bootstrap wiring

Fleetd.main loads the config, reports it and runs cfg.validateAll(). Everything after that — the whole object graph, in dependency order, then starting both faces — lives in FleetdAssembly.assembleAndStart(AssemblyInputs, ResourcePorts) (fleetd #612). main is now three statements and one call.

flowchart LR
    cfg["load config<br/>+ assertPrimaryClean"] --> router["herdr.HerdrRouter<br/>(lead daemon + optional member daemon)"]
    router --> adapters["member.ClaudeCodeLauncher<br/>member.OpenCodeLauncher<br/>→ member.CompositePeerLauncher"]
    adapters --> reap["reapOrphanWorkers()"]
    reap --> sessions["session.SessionManager<br/>+ session.GitWorktrees"]
    sessions --> leads["lead.LeadLauncher<br/>.ensureLeads()"]
    leads --> rv["msg.Rendezvous → inject.Injector<br/>+ inject.StatusPoller"]
    rv --> push["msg.ReplyPushLoop<br/>+ msg.LeadHeartbeatLoop"]
    push --> ms["msg.MessageService"]
    ms --> fh["health.FleetHealthMonitor"]
    fh --> auth["auth.MemberRegistry<br/>+ auth.CallerResolver.withLeadsAndMembers"]
    auth --> mcp["mcp.FleetMcp"]
    mcp --> app["rest.FleetApp<br/>(start Javalin)"]

Figure 2 — startup wiring, now inside FleetdAssembly.assembleAndStart. The guard asserts the primary's own environment is clean before anything else; the two adapters plug into one CompositePeerLauncher; reapOrphanWorkers() clears stale panes left by a prior daemon restart; CallerResolver is built last among the core services because it needs the live lead and architect bindings that everything above it produces.

Why the assembly is a separate method

Nothing used to call Fleetd.main far enough to observe what it passed. So any injected wiring could be swapped for its inert variant — X.none(), _ -> null, () -> Map.of() — and the whole suite stayed green. That is not a theory: #602/#606 shipped exactly that defect, and a later sweep found 16 more call sites with the same hole.

Three types carry the fix:

  • FleetdAssembly.assembleAndStart — the same statements main used to run inline, in the same order. Construction and start order is preserved deliberately; it was not rebuilt into "construct everything, then start everything", because that would change boot timing.
  • FleetdRuntime — owns the assembled objects and their single ordered close(). Its package-private accessors return the identical instances the running daemon uses, never a copy. A test that inspected a snapshot built alongside the real objects could pass while production silently received something else, which is the defect being closed.
  • ResourcePorts — every boot-time side effect a real daemon must do for real and a test must not: read the environment, connect a herdr client, open a broker, read a clock, start a scheduler, register the shutdown hook, bind the HTTP server. ResourcePorts.system() is the one production implementation.

There is deliberately no ResourcePorts.none() and no smaller overload of assembleAndStart. Adding an inert default, even for tests, would hand a future edit the exact compiling substitute this work exists to rule out. A test writes its own fake and owns that choice.

Two consequences worth knowing before you write a boot test:

  • One statement could not move without reordering startup. main registered its shutdown hook before building the Javalin app, so FleetdRuntime is built first and the app attached afterwards via attachApp, still before HTTP listens. See that field's javadoc.
  • A wiring that differs between the lead and member herdr daemons cannot be pinned by a fixture with one FakeHerdr. When memberHerdrSocket is unset the assembly falls back to memberHerdr = herdr, so router.leadAgents() and router.memberAgents() become the same object and a test cannot tell them apart. Configure two sockets. This cost two rounds of rework on #612 alone.

Flow: forward rendezvous (fleet_send → fleet_reply)

The main path. A primary's send blocks on a per-session waiter that resolves on the member's explicit reply — or, as a fallback, on a confirmed working → injectable turn boundary the injector observes, so a member that finishes without calling fleet_reply still unblocks the caller.

sequenceDiagram
    autonumber
    participant P as Primary
    participant MS as MessageService
    participant RV as Rendezvous
    participant IJ as Injector
    participant W as Member

    P->>MS: send(target, content)
    Note over MS: acquire per-session ReentrantLock
    MS->>RV: open(target) — register the waiter first
    MS->>IJ: enqueue(target, content)
    IJ->>W: inject when injectable (one msg/turn)
    IJ-->>MS: onDelivered (capture waiter + baseline)
    alt member replies explicitly
        W->>RV: fleet_reply → resolve(REPLY)
    else turn completes without reply
        IJ-->>RV: onTurnComplete → resolveCompletion(COMPLETION)
    else member wedges (CB-109)
        IJ-->>RV: onTurnFailed → resolveFailure(FAILED)
    else scrape matches a usage-limit refusal
        IJ-->>RV: resolveExhausted(BACKEND_EXHAUSTED)
    end
    RV-->>MS: Resolution
    MS-->>P: Reply (REPLIED / COMPLETED_UNREPLIED / WORKER_FAILED / BACKEND_EXHAUSTED / TIMED_OUT_*)

Figure 3 — forward send. If the timeout fires first, the outcome is TIMED_OUT_WORKING (delivered) or TIMED_OUT_QUEUED (never delivered). BACKEND_EXHAUSTED (CB-578) is a healthy pane whose account refused on a usage limit — kept apart from WORKER_FAILED so the primary sees the real cause.

The stranded-reply case. The figure is the happy path — a live waiter exists. The case that motivated the reply inbox is the reverse: a member finishes and calls fleet_reply after its send has already timed out, or one was never opened, so Rendezvous.resolve finds no waiter. Before that inbox existed, the reply was silently lost. Now MessageService.reply (msg/MessageService.java:370) publishes it into the ReplyInbox under a fresh id, and the primary collects it later — fleet_poll(target=…) or GET /sessions/{id}/replies. With broker: configured, AmqpReplyInbox gives this cross-restart durability; without it, the in-memory adapter is soft-state and undrained replies are lost on a restart.

Delivery is not purely pull. When a reply lands with no waiter, reply also calls ReplyPushLoop.onReplyQueued(target), which nudges the primary's own herdr pane to run fleet_poll — but only when the primary is injectable() (never mid-turn), bounded to a configured reminder count, and stopping the instant a drain empties the inbox.

Flow: reverse rendezvous (fleet_ask → answer)

A member pauses its own turn to ask the primary a question; the question surfaces on the primary's open forward send, and the answer resumes the same member turn. Duplicate asks from one session coalesce onto a single turnId.

sequenceDiagram
    autonumber
    participant W as Member
    participant MS as MessageService
    participant RV as Rendezvous
    participant P as Primary

    W->>MS: ask(session, question)
    MS->>RV: openAsk(session)
    Note over RV: coalesce by session — mint turnId (fresh=true)<br/>or ride the open one (fresh=false)
    alt fresh owner, a forward send is waiting
        RV->>RV: resolveQuestion(session, question, turnId)
        RV-->>P: forward send returns QUESTION + turnId
        P->>MS: answer(turnId, content)
        MS->>RV: answerAsk(turnId, content) — unblock the member
        MS->>RV: open(memberSession) — a new forward waiter
        W-->>RV: resumes the turn → fleet_reply
        RV-->>P: Reply
    else no forward send open
        RV-->>W: AskOutcome.NO_WAITER (closeAsk)
    end

Figure 4 — reverse ask. The member session is derived from turnId (Rendezvous.askSession), never a caller argument. If the turnId is unknown, answer returns Outcome.STALE_TURN.

State machine: AgentStatus

The member's herdr-reported lifecycle (herdr/AgentStatus.java:8). fromWire maps idle|working|blocked|done and folds any unrecognized string to UNKNOWN. A status is injectable only when IDLE, BLOCKED, or DONE — never mid-turn (WORKING). UNKNOWN both wedges delivery and is reclassified by StatusRefiner from the pane's real content.

stateDiagram-v2
    [*] --> IDLE
    IDLE --> WORKING: message injected, turn starts
    WORKING --> IDLE: turn completes
    WORKING --> BLOCKED: member asks / awaits input
    WORKING --> DONE: member exits the turn
    BLOCKED --> WORKING: unblocked
    DONE --> IDLE: next turn
    IDLE --> UNKNOWN: herdr can't classify
    WORKING --> UNKNOWN: herdr can't classify
    UNKNOWN --> IDLE: StatusRefiner reads the pane
    UNKNOWN --> WORKING: StatusRefiner reads the pane

    note right of WORKING
        not injectable
    end note
    note right of DONE
        injectable; a turn-boundary
        equivalent to IDLE
    end note

Figure 5 — AgentStatus. Injectable states: IDLE, BLOCKED, DONE.

State machine: Injector per-target turn lifecycle

The heart of delivery reliability (inject/Injector.java:48). A turn is judged complete only from a confirmed WORKING → injectable boundary. Three grace counters bound the failure paths, all still PICKUP_GRACE_POLLS = 8, TURN_STALL_GRACE_POLLS = 120, READINESS_GRACE_POLLS = 240 (inject/Injector.java:58,68,81).

stateDiagram-v2
    [*] --> Quiescent
    Quiescent --> Queued: enqueue(text)
    Queued --> AwaitingPickup: injectable and ready and send ok
    Queued --> Failed: 240 polls not ready

    AwaitingPickup --> AwaitingCompletion: WORKING observed
    AwaitingPickup --> Quiescent: 8 injectable polls, no WORKING (missed fast turn)

    AwaitingCompletion --> Complete: next injectable after WORKING
    AwaitingCompletion --> Failed: 120 UNKNOWN polls (stall)

    Complete --> Quiescent: onTurnComplete fires
    Failed --> Quiescent: onTurnFailed fires
    Queued --> Failed: drop(cause)
    AwaitingCompletion --> Failed: drop(cause)

    Complete --> [*]
    Failed --> [*]

*Figure 6 — per-target turn lifecycle in Injector. Complete fires TurnListener.onTurnComplete (→ CompletionResolver scrapes the pane and resolves the captured waiter, or classifies it BACKEND_EXHAUSTED when the scrape matches a configured usage-limit pattern); Failed fires onTurnFailed. CompletionResolver captures the exact waiter at delivery time and suppresses any scrape matching the pre-turn baseline, so a late completion for turn N can never resolve turn N+1's waiter (CB-116 safety). Not shown: an optional post-turn action a TurnListener can request after Complete, which the injector waits to settle before the next delivery.
Complete/Failed state names are a naming device for this diagram — the code tracks the same lifecycle as boolean flags and counters on a per-target Track record, not a literal Java enum.