Table of Contents
- 9. Implementation Architecture (as-built)
- Component map
- Package & class reference
- herdr — the wire layer
- inject — status-gated delivery and turn lifecycle
- msg — the service core
- auth — who is calling, and what they may do
- peer — the launcher SPI
- member — the peer adapters
- mcp — the MCP north face
- session — the member registry and worktrees
- health — fleet health evidence (new package)
- lead — starting configured leads
- placement, metrics, guard, config, logging, rest — the edge
- REST routes
- Bootstrap wiring
- Flow: forward rendezvous (fleet_send → fleet_reply)
- Flow: reverse rendezvous (fleet_ask → answer)
- State machine: AgentStatus
- State machine: Injector per-target turn lifecycle
9. Implementation Architecture (as-built)
Scope. This is the as-built code map of the
fleetdmodule: 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 atmaina49e968(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 statementsmainused 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 orderedclose(). 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.
mainregistered its shutdown hook before building the Javalin app, soFleetdRuntimeis built first and the app attached afterwards viaattachApp, 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. WhenmemberHerdrSocketis unset the assembly falls back tomemberHerdr = herdr, sorouter.leadAgents()androuter.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.
📖 fleet
Home — overview & the decision
Chapters
- Architecture — system · 2 invariants · 2 modes
- Message Server — the
fleetddesign - Approaches — transports compared, why herdr
- Setup — ⚫ superseded by 13
- Operations — ⚫ superseded by 13
- Team — orchestrating a mixed fleet
- Use Cases — the review scenario + mechanisms
- Roadmap — delivery record: what is live, what is off, what was dropped
- Implementation — as-built code map · classes · flows · state machines
- Cross-Host Messaging — broker topology · exchanges · queues per entity
- Features — what it can do · the knob that turns it on · why · the gotcha
- Claude → OpenCode — porting a workspace to a second host
- User Guide — 🟢 install · configure · run · delegate · the traps
- Fleet Manager — many fleets on one host, over REST
- REST API Reference — all 14 routes, roles, and bodies
- Security & Trust Boundary — the guard · authz · what a member inherits
Design proposals (not built)
- CB-548 Lead Quorum — a deterministic decision procedure around a lead's judgment
🟢 herdr-centric fleetd · AgentAPI = research, never built