Compare commits

...

32 Commits

Author SHA1 Message Date
Dai Ha 0ff0d6cb02 fleetd #808: a late fleet_reply claims its turn's waiter so the completion scrape cannot double-publish
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 2m19s
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 46s
CI / build (push) Failing after 2m3s
A send that times out leaves its captured Rendezvous waiter open for the CB-106
completion fallback (strandLateResolution). If the member then calls fleet_reply,
reply() found no live waiter and queued the answer to the inbox, but never touched
that captured waiter — so the eventual completion fallback still resolved it and
published a second, redundant scrape of the same turn.

reply() now looks up the exact waiter strandLateResolution is tracking for the
target (lateWaiters, a new per-target index onto a per-turn CompletableFuture) and
claims it with Rendezvous.resolveLateReply before falling back to the inbox. The
CompletableFuture's own single-winner complete() is the discriminator: whichever
side resolves it first decides what the other sees. A later completion scrape
against an already-claimed waiter is a no-op in CompletionResolver.resolve's
existing waiter.isDone() check, so nothing is scraped or published a second time.
A turn that never gets an explicit reply is unaffected — its own waiter is still
open when the completion fallback fires, exactly as #801 already fixed.
2026-10-07 07:02:42 +02:00
Dai Ha a2cbf50961 Bridge block: a blocking send creates no ticket (#801)
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 58s
CI / build (push) Failing after 2m19s
The paragraph that already warns a blocking fleet_send is capped by the caller's own
MCP timeout now also says it creates no ticket, so a timeout leaves no id to poll and
a resend can deliver the same brief twice, and that the one recovery route --
draining the target's inbox -- is DRAIN, so only a lead may take it.

Kept byte-identical with the wiki template (fleetd.wiki f812689).
2026-10-07 06:17:58 +02:00
Dai Ha 11352d0d06 fleetd #801: a blocking send's timeout receipt names a route the caller may use
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 50s
CI / build (push) Failing after 2m6s
A blocking fleet_send creates no ticket -- new Task( is constructed only in
sendAsync -- so its TIMED_OUT_WORKING / TIMED_OUT_QUEUED / BUSY receipt had no id
to hand back and said "retry or poll status". A retry on that path can deliver the
brief twice, and the one thing that does recover the answer went unnamed.

The receipt now carries NO_TICKET_NO_RESEND and names the target's inbox, on both
the MCP and REST surfaces. It names it only when the caller's own DRAIN is granted:
fleet_poll{target} and GET /sessions/{id}/replies both resolve to DRAIN, which is
primary-only, while SEND is granted to an architect, a collaborator and an observer
as well. A refused caller is told the reply cannot be recovered on its channel
instead of being pointed at a call it cannot make.

MessageService.strandLateResolution routes a late turn-completion to the target's
inbox. send's TimeoutException branch never completes its waiter, and
Rendezvous.close only deregisters it from session lookups, so the CB-106 fallback
can still resolve that exact future after the caller gave up -- previously into
nothing. Only Kind.COMPLETION is routed: an explicit fleet_reply reaches the inbox
through reply()'s own lookup once close() has run, so routing it here as well would
publish it twice. The publish is wrapped, because a throw inside a whenComplete
action is captured by the discarded dependent stage and reaches no caller.

Verified on the merge result, not the branch: mvn clean install, BUILD SUCCESS,
Tests run: 2244, Failures: 0, Errors: 0, Skipped: 0 -- tallied from 183 surefire
XML files, and the branch adds exactly 8 @Test methods. Reviewed in two rounds:
the first receipt named fleet_poll{target} unconditionally and the late publish was
unguarded; the second added send/answer overloads defaulting that grant to true
with no production caller, which c8e31e1 removed.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-07 06:16:21 +02:00
Dai Ha c8e31e1195 fleetd #801: drop the two mayDrainPoll-defaulting send/answer overloads
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m57s
They had no production caller — only test call sites — and the default
pointed the unsafe way: true means naming fleet_poll{target}, which is
the exact receipt this ticket is fixing. Keep one signature each and
pass mayDrainPoll explicitly at every test call site instead.
2026-10-07 06:12:57 +02:00
Dai Ha 24782fe329 fleetd #801: gate the timeout receipt's recovery route on the caller's own DRAIN grant
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 2m10s
A TIMED_OUT/BUSY receipt named fleet_poll{target}/GET /sessions/{id}/replies
unconditionally, but that route needs Authz.Action.DRAIN, which is primary-only,
while SEND and ANSWER are also granted to an architect, a collaborator, and an
observer. Those callers were told to run a call they are refused, and the late
reply they were pointed at is published somewhere they can never read.

Thread each caller's own DRAIN grant into formatReply (MCP) and writeReply
(REST); name the route only when the grant holds, otherwise say the reply
cannot be recovered on that channel.

Also wrap strandLateResolution's whenComplete callback in try/catch: an
exception thrown there lands in the discarded dependent stage and is never
rethrown, so an unguarded inbox.publish failure would lose a late reply with
no log line at all.
2026-10-07 06:05:28 +02:00
Dai Ha a49d4dcce9 Bridge block: a collaborator may read its own ticket (#804)
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 54s
CI / build (push) Failing after 2m9s
Invariant 3 said a collaborator is still refused every ticket, including its
own, and the Collaborator section said you cannot read a ticket at all. Both are
false after #804: a collaborator may poll a ticket its own fleet_send{wait:false}
created, and no other.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-07 05:40:04 +02:00
Dai Ha 2aacf07841 fleetd #803/#804: a null ticket is refused, and a collaborator may read its own
#803: MessageService.ownsTicket(String, String) guards a null ticket instead of
handing it to ConcurrentHashMap.get, which rejects a null key. fleet_poll{} with
no ticket resolves to TASK_READ with a null id, so the crash sat on the path a
caller reaches by calling the tool wrong.

#804: TASK_READ now grants a collaborator a ticket its own fleet_send{wait:false}
created, on the same ownership classifier the observer grant uses. A collaborator
could already create a ticket -- FleetMcp.sendAsync records any caller as the
creator with no role filter -- so the reply had an owner that was refused it.

Renames NO_OBSERVER_OWNED_TICKET to NO_OWNED_TICKET and observerOwnsTicket to
ticketOwnedByCaller, since neither is observer-specific any more.

The first version also threaded a mayTaskReadNudge boolean through FleetMcp,
PrimaryRegistry and ReplyPushLoop to suppress a ticket nudge to a role Authz
would refuse. With the collaborator granted, that condition is true for every
role reaching the call site, so 820e6f2 removes it.

Verified on the merge result, not the branch: mvn clean install, BUILD SUCCESS,
Tests run: 2236, Failures: 0, Errors: 0, Skipped: 0 -- tallied again from 183
surefire XML files.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-07 05:39:00 +02:00
Dai Ha f2e025cb13 fleetd #801: name the inbox route in a blocking-send timeout, and stop dropping a late completion
CI / shell-tests (pull_request) Failing after 11s
CI / contract (pull_request) Successful in 57s
CI / build (pull_request) Failing after 2m10s
The TIMED_OUT_WORKING/TIMED_OUT_QUEUED/BUSY receipt (MCP and REST) named no
recovery route and invited a bare retry, which can duplicate the delivery. It
now names fleet_poll{target=...} (or GET /sessions/{id}/replies) and states
that a blocking send is never tracked by a ticket.

A blocking send's TimeoutException path leaves its rendezvous waiter in place:
Rendezvous.close only deregisters it from future session lookups, while the
turn-completion fallback resolves it through a captured future reference and
can still succeed after the caller gave up. With nobody left listening, that
late completion is now published to the target's inbox, the same route an
unstructured fleet_reply already uses when it arrives with no caller waiting.
2026-10-07 05:37:55 +02:00
Dai Ha 820e6f2620 fleetd #803/#804: remove dead mayTaskReadNudge plumbing
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 56s
CI / build (pull_request) Failing after 2m6s
The #804 grant gave a collaborator TASK_READ for its own ticket, which
dissolved the only case the mayTaskReadNudge gate existed for. With the
real ownership classifier forced to t -> true at its one call site, the
grant reduced to "every non-anonymous role", so the filter could never
suppress a ticket nudge in production. Revert ReplyPushLoop, PrimaryRegistry
and FleetMcp to their pre-plumbing shape, and drop the two tests that only
proved the dead seam, not that any caller could reach it.
2026-10-07 05:36:13 +02:00
Dai Ha 72cc7560f1 fleetd #803/#804: fix null-ticket NPE in ownsTicket, grant collaborator TASK_READ
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Failing after 2m11s
#803: MessageService.ownsTicket(String, String) called tasks.get(null) when
fleet_poll{} supplied no ticket, throwing NPE instead of refusing. Guard the
lookup so a null ticket returns false, same as an unknown one.

#804: a collaborator could create a ticket via fleet_send to another
collaborator but could never read it back, since TASK_READ's grant only
covered primary/worker/architect/observer. Extend the grant to a collaborator
confined by the same ownership classifier observer already uses, rename the
classifier parameter now that it serves both roles, and gate the ticket-nudge
in ReplyPushLoop the same way #778 gated the reply/question nudges so no
caller is nudged toward a fleet_poll its role would refuse.
2026-10-07 05:26:34 +02:00
Dai Ha 50c7674b98 Bridge block: invariant 4 — the prompt-box gate is off by default (#797)
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 51s
CI / build (push) Failing after 1m59s
Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-07 05:20:03 +02:00
Dai Ha fcce44386b fleetd #797: make the prompt-box gate configurable, default off
PromptBox.clearToSubmit returns true at once when the gate is disabled,
before any pane read, hold log or streak update, so every call site behaves
as if the box were empty whatever its own polarity. promptBoxGateEnabled is
a new top-level config key; null and false both leave the gate off.

The gate held deliveries on the pane's own autocomplete suggestion, which
detection reads as the operator's unsubmitted typing. Detection is unchanged
and still wrong (#802); turn the gate back on once that is fixed.

Merge of PR #805 (branch worker/797-disable-box-gate-5b5478-12).

Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-07 05:18:43 +02:00
Dai Ha f3f067f46f fleetd #797: add a config switch to disable the prompt-box delivery gate
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 2m13s
PromptBox currently misreads a pane's own autocomplete suggestion at the
caret as the operator's unsubmitted typing, holding deliveries that were
never actually at risk. Add promptBoxGateEnabled (default off) so the gate
can be switched off without deleting its classification logic — detection
stays fixable separately (fleetd #802).

When disabled, PromptBox.clearToSubmit returns true for every target
without reading the pane, logging a hold, or counting a streak. The three
call sites (ReplyPushLoop, LeadCoordLoop, LeadHeartbeatLoop) each gained a
boolean constructor parameter threaded straight into their PromptBox,
with the previous constructor demoted to a delegator passing true so
every existing caller keeps the gate on by default.
2026-10-07 05:13:50 +02:00
Dai Ha 9f4b736fb9 Bridge block: an observer may read a ticket it created (#778)
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 1m59s
Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-07 05:10:58 +02:00
Dai Ha 52cd0470b6 fleetd #778: an observer may read a ticket it created
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 52s
CI / build (push) Failing after 1m59s
Authz grants TASK_READ to an observer only for a ticket whose creating
terminal is its own, checked by a classifier the poll call sites thread in.
ReplyPushLoop no longer nudges a recipient to run a DRAIN or ANSWER its own
role would be refused.

Merge of PR #795 (branch worker/778-e301b0-1).

Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-07 05:05:06 +02:00
Dai Ha a3eeace447 handover skill: BOOTSTRAP_NEVER_SENT, and relaunchReadySeconds bounds three waits
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 1m2s
CI / build (push) Failing after 1m55s
Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-06 19:28:21 +02:00
Dai Ha b696c31756 fleetd #796: retry bootstrap delivery to a fresh lead pane
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m5s
CI / build (push) Failing after 2m17s
A roll killed the old pane, launched the new one, then lost bootstrapText to a
herdr agent_not_ready refusal. The successor woke with no handover and no way to
tell a failed roll from a cold start.

sendBootstrapWithRetry retries only that refusal, bounded by
relaunchReadySeconds. A persistent refusal now ends the roll with the new
BOOTSTRAP_NEVER_SENT outcome, so fleet_handover{action:"status"} can report it.

Verified at 027d413 in a throwaway worktree: 2210 tests, 0 failures, 0 errors,
clean install.

Co-Authored-By: Claude <noreply@anthropic.com>
2026-10-06 19:27:36 +02:00
Dai Ha 96ebae28a6 reviewer skill: send fleet_reply before writing the reasoning out
Five reviewer turns were lost in one session. Each wrote a long, correct
analysis to its own pane and ended the turn with no fleet_reply, so the
bridge scraped the pane and the lead received a clipped fragment.

The briefs are half the cause: each carried a five-item checklist of things
to hunt alongside a ~90-word capped output format, which reads as two
contradictory output contracts. The skill now says the checklist is where to
look, not the shape of the answer, and that the reply goes out as soon as
the answer is known.

Measured: tickets task-7da785-29 and -30 resolved at 18:28 via the
turn-completion fallback, 1233 and 1403 chars scraped.
2026-10-06 19:01:38 +02:00
Dai Ha b41aa663f7 Bridge block: invariant 6 — an operator outranks the confirm rule
Invariant 6 as first written told every receiver to confirm, with no
mention of who may override it. A peer cannot impose a rule on another
operator's session.

Measured: the trinotes pane (observer, work config dir) answered the
communication test and then said its operator's standing rule is not to
answer fleet messages, that a peer cannot change that rule, and that it
would ask its operator before following mine. That is the correct reading
and the rule now says so.

Also records the symmetric error: a sender must not read silence as
agreement or as a dead session.

wiki/7-Use-Cases.md synced byte-identical.
2026-10-06 18:48:12 +02:00
Dai Ha 027d413ce9 fleetd #796: document bootstrap retry readiness bound
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Failing after 1m59s
2026-10-06 18:26:01 +02:00
Dai Ha f544906621 Bridge block: a receiver confirms what it receives
Adds invariant 6 to the canonical block. Nothing in it said a received
message must be answered, so a sender could not tell a handled message
from one that never arrived.

Measured today: the anki pane (observer) sent to vms, which reads
deliverable:false because it has not contacted the daemon since the 08:05
restart. The send was accepted, held 83s, and failed with the body 'nst' --
three characters scraped off the vms screen. The sender read that as an
inconclusive result and asked the lead what happened. fleetd #757 covers
the accept-time refusal; this rule covers the half the protocol owns.

Operator's words: 'when a message sent, at least the receiver should
confirm unless explicit told to not reply'.

wiki/7-Use-Cases.md template synced byte-identical (wiki 3ad606f).
2026-10-06 18:18:10 +02:00
Dai Ha c88f01ecb8 fleetd #796: retry bootstrap delivery after agent readiness refusal
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Failing after 2m4s
2026-10-06 18:16:47 +02:00
Dai Ha 6022aef612 fleetd #778: let an observer read its own ticket, never nudge a role to run a call it is refused
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 59s
CI / build (pull_request) Failing after 2m15s
Part 1: Authz.TASK_READ now grants an observer the ticket it created (a new
observerOwnsTicket classifier), via MessageService.ownsTicket(String, String)
(side-effect-free) and a new Authz.permits overload. fleet_status and
fleet_poll{target} (DRAIN) stay refused for an observer. REST and MCP share
the same gate.

Part 2: ReplyPushLoop no longer nudges a reply/question delegator to run
fleet_poll(target=...) or fleet_send(turnId=...) unless Authz would actually
grant that call to its role. FleetMcp.sendHandler asks Authz.permits(DRAIN)/
permits(ANSWER) once, eagerly, while the real Principal is in scope, and
threads the booleans through PrimaryRegistry.recordDelegation into the new
Delegation fields and nudgeMayDrainFor/nudgeMayAnswerFor accessors --
plain booleans, not a Role reference, so no new or widened edge crosses the
frozen auth<->mcp<->msg package-cycle baseline. ReplyPushLoop filters a
reply/question out of its pending set when the grant is false, falling back
to the durable inbox / the ask window lapsing, same as when no nudge target
is known at all. Log lines that called a non-lead nudge target "lead" now
say "nudge target".
2026-10-06 13:55:26 +02:00
Dai Ha 66ab9cab24 fleetd #791: addendum — plugin 0.3.0 has no mount and is installed in four instances
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m59s
2026-10-06 08:03:50 +02:00
Dai Ha f611488c9b fleetd #791: README — the mod still reads FLEETD_MCP_URL as an optional override
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 56s
CI / build (push) Failing after 2m4s
2026-10-06 08:03:06 +02:00
Dai Ha 44d16e7639 fleetd #791: merge PR #794 — fleet plugin 0.3.0, no MCP mount, mod off for members
The plugin no longer carries .mcp.json; the instance or the project mounts
fleet. The mod asks fleet_whoami once and skips the fleet_inbox poll for a
worker or an architect, so members keep paste delivery. The mod reads
FLEETD_MCP_URL with $.env.get (documented in the mods API), defaulting to
http://127.0.0.1:8765/mcp.

Lead verification on 87f7c03: claude plugin validate plugin passed;
claude plugin test plugin 21 pass, 0 fail. With the role gate removed, the
worker and architect tests fail (19 pass, 2 fail), so they are not vacuous.
2026-10-06 08:03:01 +02:00
Dai Ha 87f7c03535 fleetd #791: fleet plugin 0.3.0 — drop the MCP mount, gate the mod off for members
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Failing after 2m7s
- Delete plugin/.mcp.json; mounting fleet is now the instance's or the
  project's job, not the plugin's. The setup skill still writes the
  project .mcp.json entry.
- register.js: read FLEETD_MCP_URL via $.env.get, defaulting to
  http://127.0.0.1:8765/mcp, documented at
  https://code.claude.com/docs/en/plugins/mods/api.md ("$.env — get
  and set environment variables").
- Gate the fleet_inbox poll by role: a worker or architect learns its
  role once from fleet_whoami and then never polls fleet_inbox; every
  other role polls as before. An unreachable daemon retries whoami on
  a later tick rather than deciding.
- Bump 0.2.0 -> 0.3.0 in plugin.json and marketplace.json, and update
  README/SKILL.md to drop every claim that the plugin mounts fleetd.
- Add 6 tests for the role gating; all 21 tests and
  `claude plugin validate` pass.
2026-10-06 06:46:26 +02:00
Dai Ha 562354a58d fleetd #790: document observer -> lead send in the bridge block
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 51s
CI / build (push) Failing after 1m59s
Invariant 3, the observer paragraph and the unconfigured-pane row now say an
observer may send to a lead. Invariant 4 now says a direct send to a lead pane
is not box-gated yet (#793). Wiki template synced (check printed in sync: True).
2026-10-06 06:25:31 +02:00
Dai Ha eaa1f0d170 fleetd #790: merge PR #792 — an observer pane may fleet_send to a lead
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 56s
CI / build (push) Failing after 2m5s
An observer's SEND classifier (CallerResolver.observerSendTarget, renamed from
sendableObserverTarget) now accepts a configured lead terminal as well as
another observer pane. Collaborators, architect slots and spawned members stay
refused. The [fleet_send from observer term_…] prefix is kept. An observer
learns a lead's sessionId from its filtered fleet_list panes rows (role "lead").

Lead verification: mvn clean install on 72ea7de in a throwaway worktree,
exit 0, surefire totals 2208 run / 0 failures / 0 errors / 0 skipped.

Follow-up found in review: #793 (the Injector has no prompt-box gate for a
direct send to a lead pane).
2026-10-06 06:24:33 +02:00
Dai Ha 72ea7def0a fleetd #790: let an observer pane fleet_send to a lead
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Failing after 2m7s
An observer could only reach another observer pane, so a hand-opened tab
could answer a lead with fleet_reply but never start a conversation with
one.

CallerResolver.sendableObserverTarget() is renamed observerSendTarget()
and now accepts a configured lead terminal as well as a pane that falls
to the observer floor. It reads the same maps resolve() reads, in the
same order: a live spawned member is refused first, a lead is then
accepted even when the same pane is also bound to an architect slot,
and a collaborator or an architect-slot pane is refused. A collaborator,
an architect and a spawned member stay unreachable.

Authz keeps its one SEND case; only the classifier it is handed widened,
so the MCP gate (FleetMcp.denyFor) and the REST gate (FleetApp.allow)
give the same answer from the same predicate.

fleet_list shows the lead's address in the observer's filtered `panes`
rows rather than by adding the `leads` array: the panes filter already
reads the same classifier as the SEND gate, so reachability and
visibility cannot drift, and the reduced row carries no lead name,
context window or config dir.

Gates traced end to end for an observer -> lead send: denyFor, the
sendTool argument validation (profileTargetError rejects profile names
only), MessageService (records the caller's owner key, no role gate),
the injector's readiness gate (Fleetd.deliverableTo accepts a lead
through the leads map, never a presence entry) and its status gate.
Unchanged: an observer still holds no TASK_READ, so it cannot poll the
ticket a wait:false send returns.

Tests: observer -> lead allowed and observer -> collaborator / architect
slot / spawned member refused over denyFor, each with a positive control
in the same test; the same matrix over the REST route; a real MCP client
through the real MessageService and Injector to the real herdr call,
proving the [fleet_send from observer term_...] prefix reaches a lead's
pane and that the lead passes the production readiness gate; and the
observer's fleet_list panes rows now carrying the lead.

mvn clean install in fleetd/: Tests run: 2208, Failures: 0, Errors: 0,
Skipped: 0 — BUILD SUCCESS. Reverting the widening in CallerResolver
alone turns 7 of the new assertions red across all five touched test
classes, so none of them passes vacuously.
2026-10-06 06:20:43 +02:00
Dai Ha 98f3cd5aa7 fleetd #788: document fleet_inbox in the bridge block and the addendum
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 1m57s
Block: observer INBOX right, invariant 3 (inbox is only-as-itself),
invariant 4 (a pane that collects its own mail skips the paste and the
box gate, not the status gate), and an intent-table row. The wiki
template carries the same block; the sync check printed in sync: True.
2026-10-06 05:54:13 +02:00
Dai Ha aec5ff2c2f fleetd #788: merge PR #789 — fleet_inbox, so a pane can collect its own mail
A pane that calls fleet_inbox within 15 s is mod-served: the Injector
offers its next message for collection instead of pasting it, and keeps
it queued until the pane takes it. A pane that stops polling gets its
mail pasted again. The fleet mod polls fleet_inbox every 3 s on one
reused MCP session and hands each message to Claude with
$.prompt.submit.

Lead proof at d9fedfb, in a throwaway worktree:
mvn clean install exit 0; surefire XML 182 files, 2202 tests,
0 failures, 0 errors, 0 skipped. claude plugin validate plugin: passed.
claude plugin test plugin: 15 pass, 0 fail.
2026-10-06 05:53:16 +02:00
43 changed files with 2176 additions and 326 deletions
+2 -2
View File
@@ -8,8 +8,8 @@
{
"name": "fleet",
"source": "./plugin",
"description": "Mount the fleetd MCP gateway and apply standard Claude Code settings so a session can orchestrate delegated workers. Ships no credentials.",
"version": "0.2.0",
"description": "Apply standard Claude Code settings so a session can orchestrate delegated workers, and run the fleet mod for cross-session messaging. Mounting the fleetd MCP gateway is the instance's or the project's job, not this plugin's. Ships no credentials.",
"version": "0.3.0",
"author": {
"name": "LTMS"
}
+4 -2
View File
@@ -189,8 +189,9 @@ If the first number has grown past 20, somebody has rolled under the restart pat
should be replaced with what they measured.
**Three separate timeouts bound a roll.** `leadRollover.relaunchReadySeconds` (default 45,
`FleetConfig.java:1483`) bounds **each** of two waits that run after the relaunch, so the worst case
there is about twice that number, not 45 seconds in total. A third bound gives your old pane 10
`FleetConfig.java:1483`) bounds **each** of three waits that run after the relaunch, so the worst
case there is about three times that number, not 45 seconds in total. The third wait retries
`bootstrapText` while herdr answers `agent_not_ready`. A separate bound gives your old pane 10
seconds to die (`LeadRollover.PANE_DEATH_TIMEOUT_SECONDS`).
`fleet_handover{action: "status", token}` answers with one of these:
@@ -204,6 +205,7 @@ seconds to die (`LeadRollover.PANE_DEATH_TIMEOUT_SECONDS`).
| `RELAUNCH_FAILED` | launching the fresh pane failed |
| `RELAUNCH_NEVER_READY` | the fresh pane never became ready within `relaunchReadySeconds` |
| `RELAUNCH_NOT_RECOGNISED` | the fresh terminal never resolved as a lead |
| `BOOTSTRAP_NEVER_SENT` | the fresh pane was ready, but herdr refused `bootstrapText` with `agent_not_ready` for the whole bound, so the successor never learned where the handover file is |
| `FAILED` | the roll threw; `runRollover`'s catch records this rather than leaving it stuck |
Only `TURN_NEVER_SETTLED` guarantees your context is intact. The other failures can leave you
+9
View File
@@ -37,6 +37,15 @@ confident-but-wrong finding. Anything you could settle by reading more code is y
## 4. The finding — what goes in `fleet_reply`
**Call `fleet_reply` as soon as you know your answer, before you write the reasoning out.** The
four lines below are the whole deliverable, and they are short on purpose. Analysis you type into
your terminal reaches nobody: when a turn ends with no `fleet_reply`, the bridge scrapes the pane
and the lead receives a clipped fragment instead of a finding. A long, correct analysis and no
`fleet_reply` is a failed turn, and it is the most common way this role fails.
If the delegation also handed you a list of things to check, that list is where to *look*. It is
not the shape of the answer. Work the list, then still send these four lines.
Report the **single most important** real issue in the scope, in these four lines, under
~90 words:
+77 -27
View File
@@ -33,10 +33,13 @@ gate uses. A worker also carries its `sessionId`, `profile`, `worktree` and `bra
carries the slot name it was bound to; a collaborator carries its registry name and its own
`sessionId`, and **no `leader` key** — a collaborator is a named peer, not a primary. An **observer**
carries only its own `sessionId`: a pane the daemon could not place as any of the above, authorized
to `READ`/`METRICS`, to `REPLY`/`ASK` on its own pane, and to `SEND` only to a target that resolves
as an observer too — never to a lead, a collaborator, or a spawned member, and never a ticket. It
finds such a target in `fleet_list`'s `panes` array, which for an observer is filtered to exactly
what it may send to and reduced to `sessionId`, `label`, `status`, `role` and `deliverable`.
to `READ`/`METRICS`, to `REPLY`/`ASK`/`INBOX` on its own pane, to `SEND` only to a lead or to a target
that resolves as an observer too — never to a collaborator, an architect, or a spawned member — and to
`TASK_READ` only a ticket its own `fleet_send{wait:false}` created. Nothing wider: not another
caller's ticket, and not a session's status, which stays refused. It finds such a target in
`fleet_list`'s `panes` array, which for an observer is filtered to
exactly what it may send to and reduced to `sessionId`, `label`, `status`, `role` and `deliverable`;
a lead's row reads `role: "lead"`.
Don't infer what you can ask.
Only if that call is unavailable, fall back to these — each is one-way, so keep reading until one
@@ -67,23 +70,54 @@ and the sender silently receives nothing. Fail toward the recoverable error.
3. **Identity comes from the connection, never an argument.** Workers never pass a target; you
cannot act as another session. Spawn/stop/drain are lead-only; **send is lead, architect,
collaborator, or observer** — and a collaborator may send only to a lead or another collaborator,
never to a spawned member's terminal, while an observer may send only to another observer pane;
reply/ask are only-as-itself — any peer may answer for its own pane,
and for no other. A call outside your role is refused, not queued.
never to a spawned member's terminal, while an observer may send only to a lead or another observer
pane;
reply/ask/inbox are only-as-itself — any peer may answer, or collect its mail, for its own
pane, and for no other. **A ticket is owned by the caller that created it, not by a role**, so an
observer or a collaborator may poll a ticket its own `fleet_send{wait:false}` created and no
other — never another caller's, and never a session's status, which stays refused for both. A
call outside your role is refused, not queued.
4. **Delivery is status-gated: one message per turn.** Don't busy-poll a peer's terminal and don't
re-send because a call looks slow — the bridge delivers when the peer is `idle`, `blocked` or
`done`. A spawned member must **also** have mounted the bridge MCP: until it has, it is not
deliverable, and a send waits on that gate for ~60s and then fails without ever reaching its pane.
**A lead's own pane has a second gate: its input box must be empty.** The multiplexer pastes and
submits in one step, so a delivery that lands while the operator is typing submits their
half-written line. A heartbeat, a ticket nudge and lead-to-lead mail therefore wait until the box
is clear, and a pane the daemon cannot read as a box waits too. Nothing is lost — every one of
those paths retries — but a lead that leaves text sitting in its box receives nothing until it
clears, and the only sign is one warning in `fleetd.out` after 20 held checks in a row. Delivery
to a *member* is not gated this way, because nobody types in a member's pane.
**A lead's own pane has a second gate, and it is off by default.** The multiplexer pastes and
submits in one step, so a delivery that lands while the operator is typing would submit their
half-written line. A heartbeat, a ticket nudge and lead-to-lead mail can therefore wait until the
box is clear. That check is the config key `promptBoxGateEnabled`; both an absent key and `false`
leave it off, so **right now nothing waits on the box**. It is off because detection read the
pane's own autocomplete suggestion as the operator's typing, and so held deliveries that were
never at risk — 2231 holds in 11 hours across four panes, measured 2026-10-06. Turn it on only
once detection is fixed (#797 made it switchable, #802 is the detection bug). While it is on, a
lead that leaves text in its box receives nothing until it clears, and the only sign is one
warning in `fleetd.out` after 20 held checks in a row; nothing is lost, because every one of
those paths retries. A direct `fleet_send` to a lead's pane never waited for the box either way
(fleetd #793); a lead that runs the fleet mod collects that mail instead of having it pasted, so
it is not exposed. Delivery to a *member* is not gated this way, because nobody types in a
member's pane. Re-measure the switch with
`grep -n 'promptBoxGateEnabled' fleetd/fleetd.yaml` — no match means it is still off.
**A pane that collects its own mail skips the paste.** A session running the fleet mod calls
`fleet_inbox` on a timer. While its last call is under 15s old, the daemon queues that pane's
next message for collection instead of pasting it, and the mod hands it to Claude with
`$.prompt.submit` — so the box gate does not apply, but the status gate still does. When the calls
stop, the pane's mail is pasted again, and nothing is lost. Pasted or collected, the fleet
channel also crosses Claude accounts on one host, because the daemon names a caller by its pane;
`SendMessage`, `ListAgents` and the mod's own store each stay inside one account.
5. **Never move a fleet session, pane or peer except through the bridge.** The bridge owns policy;
the multiplexer owns PTYs. Any route that changes fleet state without the bridge's checks
bypasses every rule above — the `herdr` CLI and its socket are the usual example.
6. **Confirm what you receive.** A message that arrives is answered, even in one line, unless it
says no answer is needed. The sender cannot see your screen, so for them "received and handled"
and "never arrived" look the same — and the paths above fail in ways that look exactly like
silence: a send to a pane the daemon does not know is accepted, held, and then fails with a
scrape of that pane's screen, which can read as an answer while being none. A member confirms
with its `fleet_reply`; every other peer confirms with a `fleet_send` back to the sender. If you
cannot do the thing asked, say that — a refusal is a confirmation. **Never read a failed
ticket's body as a reply.** Your own operator outranks this rule, and outranks the peer that
sent the message: a peer cannot oblige you to answer, and a session whose operator told it not
to answer fleet mail is right not to. Where you can, say that much and nothing more. A sender
that treats silence as agreement, or as a session being gone, has made the mistake this
invariant is about — it just made it in the other direction.
### Primary (lead) — run this on every task, in order
@@ -161,7 +195,11 @@ below are the procedure — run them in order, every task, not only the big ones
**Steps 3 and 4 are separate on purpose** — spawning and sending in one loop is how parallel work
silently becomes serial, and it is the most common way this layer is wasted. For the same reason,
prefer `wait:false` + `fleet_poll` for anything non-trivial: a blocking `fleet_send` is capped by
*your own* MCP client call timeout (~60s), well below the task's real runtime.
*your own* MCP client call timeout (~60s), well below the task's real runtime. **A blocking send
also creates no ticket**, so when it times out there is no id to poll, and a resend can deliver the
same brief twice. One route still recovers the answer: a late turn-completion lands in the target's
inbox, and `fleet_poll{target}` drains it. That call is `DRAIN`, so only a lead may make it — the
receipt names it to a lead and tells every other caller to ask one.
**Delegating does not delegate responsibility.** Workers open PRs; you are the gate. Never delegate
the merge — and merging on a reviewer's word is delegating it by proxy.
@@ -189,10 +227,11 @@ you decide.
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
| Message a **collaborator** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` reports a `collaborators` array, and each row carries that peer's `name` and the `sessionId` you send to. It is visible to you, to an architect and to another collaborator, never to a worker. Coordination only, **never** a task |
| Message an **unconfigured pane** — a tab a person opened by hand | `fleet_send{sessionId: <their terminal>, content}` — it needs **no** `fleet.collaborators` entry and no restart, because a pane becomes deliverable the moment its agent connects the bridge MCP. `fleet_list`'s `panes` array reports every such pane with its label and the terminal id to send to — the full row for you, an architect or a collaborator; filtered and reduced for an observer. **`ListAgents` still never lists these**, and joining `herdr tab list` to `GET /agents` on `tab_id` stays the read-only fallback if the array is missing. Such a pane resolves as an `observer`: it can answer you with `fleet_reply`, and it can `fleet_send` to another observer pane, but never to you. Coordination only, **never** a task |
| Message an **unconfigured pane** — a tab a person opened by hand | `fleet_send{sessionId: <their terminal>, content}` — it needs **no** `fleet.collaborators` entry and no restart, because a pane becomes deliverable the moment its agent connects the bridge MCP. `fleet_list`'s `panes` array reports every such pane with its label and the terminal id to send to — the full row for you, an architect or a collaborator; filtered and reduced for an observer. **`ListAgents` still never lists these**, and joining `herdr tab list` to `GET /agents` on `tab_id` stays the read-only fallback if the array is missing. Such a pane resolves as an `observer`: it can answer you with `fleet_reply`, and it can `fleet_send` to you or to another observer pane, but never to a collaborator or a member. Coordination only, **never** a task |
| Answer a peer lead that messaged you | `fleet_send{coordId}` — or `{sessionId}` if they are on this host. **Not** `fleet_reply`: it has no peer route and the publish is refused |
| Read your own held lead-to-lead mail (no ack) | `fleet_poll{coordId: <your own coord-id, from fleet_list's coordinator.selfId>}` — primary-only; never acks, so `fleet_list`'s `held[]` still shows it after. `fleet_list`'s `held[]` gives only a truncated preview — this is the only way to read the full body |
| Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` |
| Collect the mail queued for your OWN pane, instead of having it pasted | `fleet_inbox` — no arguments, any role, own pane only. The fleet mod calls it on a timer; you rarely call it by hand |
| Tear down a member | `fleet_stop{paneId}` |
| Replace your OWN lead session when its context is full | `fleet_handover{action:"open", reason?}` → write the handover file it names → `fleet_handover{action:"confirm", token, operatorConfirmed}`. Primary-only. **In that order**: the file must be modified *after* `open`, or `confirm` refuses it as stale. There is no terminal parameter — the pane is always your own, so you can never roll another lead. `{action:"cancel", token}` drops a pending request |
@@ -267,19 +306,21 @@ simply complies has thrown away the reason there are two of you.
`fleet_whoami` answered `collaborator`, so your pane's tab matches a `fleet.collaborators.<name>.tab`
entry. You are **not** a member: nothing delegates to you, you have no brief, no worktree and no
ticket, and **you owe no `fleet_reply`** — the turn contract above is for a session a lead spawned,
and it does not apply to you. Read it only to understand what the members around you are doing.
ticket raised against you, and **you owe no `fleet_reply`** — the turn contract above is for a
session a lead spawned, and it does not apply to you. Read it only to understand what the members around you are doing.
What you may do: observe the fleet (`fleet_list`, `fleet_profiles`, `fleet_whoami`), and send to a
lead or to another collaborator. What you may not: spawn, stop or drain anything, roll a lead's
session, answer a member's `fleet_ask`, poll a ticket, or send to a spawned member's terminal. Each
of those is refused at the gate, not queued.
What you may do: observe the fleet (`fleet_list`, `fleet_profiles`, `fleet_whoami`), send to a lead
or to another collaborator, and poll a ticket your own `fleet_send{wait:false}` created. What you may
not: spawn, stop or drain anything, roll a lead's session, answer a member's `fleet_ask`, read
another caller's ticket or any session's status, or send to a spawned member's terminal. Each of
those is refused at the gate, not queued.
Two limits worth knowing before you hit them. **You cannot reach a worker** — not even to help one —
because a worker belongs to the lead that spawned it, and routing around that would make you a
second orchestrator with no plan. Send to the lead instead. And **you cannot read a ticket**, so you
cannot collect a delegation's reply: `fleet_poll` refuses you at the role gate, and a ticket also
records the terminal that created it, so even a leaked id reads nothing.
second orchestrator with no plan. Send to the lead instead. And **you cannot read another caller's
ticket** — only one your own `fleet_send{wait:false}` created. A ticket records the terminal that
created it, so a leaked id reads nothing, and a lead's delegation stays the lead's: you cannot
collect its reply.
Being named buys you a channel, not authority. Your `fleet_send` to a lead is coordination between
peers: the lead owes you no obedience, and you owe it none.
@@ -345,8 +386,17 @@ must obey belongs in the charter, not here.
`handover` (write the file a fresh lead session inherits when the outgoing one hands off,
fleetd #480).
- **This repo is also a Claude Code marketplace, and ships a plugin.** `.claude-plugin/marketplace.json`
points at `plugin/`, which carries the MCP mount and the `setup` skill
(`/claude-bridge:setup` — make any project bridge-ready). It was added in CB-527 and then went
points at `plugin/`, which carries the `setup` skill (`/fleet:setup` — make any project
bridge-ready) and no MCP mount: the instance or the project mounts `fleet`. `plugin/` is also a Claude Code mod (`plugin/hooks/`):
it polls `fleet_inbox` and hands each message to Claude with `$.prompt.submit`. Measured
2026-10-06: `fleet@fleetd` 0.3.0 is installed at user scope in the `gx10`, `ltms`, `ollama` and
`work` instances, so every session started from them runs the mod, and a spawned member on those
config dirs loads it too — but the mod skips the inbox poll for a `worker` or an `architect`, so a
member still gets its brief pasted. The install made a copy under each instance's
`plugins/cache/fleetd/fleet/0.3.0`, so assume an edit to `plugin/` reaches sessions only after a
version bump and `claude plugin update fleet@fleetd` per instance. Re-measure with
`grep -l '"fleet@fleetd"' ~/.ccs/instances/*/plugins/installed_plugins.json`; delete this sentence
if the plugin is uninstalled. It was added in CB-527 and then went
unmentioned by every instruction file, so it drifted and a later session planned it from scratch
(#362). **Read `plugin/` before designing anything about onboarding a project.** Two limits are
structural, not bugs: a plugin cannot carry the role agent files, because
+6
View File
@@ -174,6 +174,12 @@ bind:
# idleSleepGuard:
# enabled: false
# Whether PromptBox checks a lead's input box before a delivery pastes into it. Default off —
# omitting this key (or setting it false) leaves every delivery proceeding without reading the
# pane. Detection currently misreads the pane's own autocomplete suggestion at the caret as the
# operator's unsubmitted typing, so the check holds deliveries that were never actually at risk.
# promptBoxGateEnabled: false
# herdr Unix socket. Omit to use the client default
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
herdrSocket: ~/.config/herdr/herdr.sock
@@ -421,7 +421,7 @@ final class FleetdAssembly {
// are counted at their single funnel rather than at each of the two caller-facing surfaces.
Metrics metrics = FleetMetrics.create(sessions, replyInbox);
var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox,
pushScheduler, maxReminders, backoffMs, metrics);
pushScheduler, maxReminders, backoffMs, metrics, cfg.promptBoxGateEnabled());
// fleetd #201 Unit 5: point the forwarding holder captured by the backendErrorSink lambda
// above at the real push loop, now that it exists.
pushLoopRef.set(pushLoop);
@@ -448,7 +448,8 @@ final class FleetdAssembly {
Fleetd.leadContextSource(leadContextGauge, router.leadAgents(), leads,
Fleetd.leadConfigDirLookup(() -> config.get().profiles(), leaders),
Fleetd.leadContextWindowLookup(() -> config.get().profiles(), leaders)),
Boolean.TRUE.equals(hb.contextHighNudge()), requireOperatorConfirm);
Boolean.TRUE.equals(hb.contextHighNudge()), requireOperatorConfirm,
cfg.promptBoxGateEnabled());
heartbeat.start();
} else {
heartbeat = null;
@@ -545,7 +546,7 @@ final class FleetdAssembly {
if (leadMailbox != null) {
var leadCoordScheduler = ports.newScheduler("bridge-leadcoord-");
leadCoordLoop = new LeadCoordLoop(leadMailbox, router.leadAgents(), leads, leadCoordScheduler,
LEAD_COORD_INTERVAL_MS);
LEAD_COORD_INTERVAL_MS, cfg.promptBoxGateEnabled());
leadCoordLoop.start(); // FIFTH of the recurring background loops to start (optional).
leadCoordSchedulerRef = leadCoordScheduler;
} else {
@@ -80,34 +80,60 @@ public final class Authz {
/**
* The fail-closed classifier for an observer's {@code SEND}: answers no for every target, so
* the grant is refused unless a caller supplies a real one. {@code
* CallerResolver#sendableObserverTarget()} is the real one, read from the same maps {@code
* CallerResolver#resolve} consults, so a target that classifier calls known is one {@code
* resolve} would actually resolve as {@link Role#OBSERVER}.
* CallerResolver#observerSendTarget()} is the real one, read from the same maps {@code
* CallerResolver#resolve} consults, so a target that classifier accepts is one {@code resolve}
* would actually resolve as a lead ({@link Role#PRIMARY}) or as {@link Role#OBSERVER}.
*/
public static final Predicate<String> NO_KNOWN_OBSERVER_TARGET = target -> false;
public static final Predicate<String> NO_OBSERVER_SEND_TARGET = target -> false;
/**
* The fail-closed classifier for a ticket-scoped {@code TASK_READ} (an observer's or a
* collaborator's): answers no for every ticket, so the grant is refused unless a caller
* supplies a real one. {@code MessageService#ownsTicket(String, String)} is the real one, so
* a ticket this classifier accepts is one that caller's own {@code fleet_send(wait:false)}
* actually created.
*/
public static final Predicate<String> NO_OWNED_TICKET = ticket -> false;
/**
* Convenience form for a caller with no classifier to supply. Fails closed: a collaborator's
* or an observer's {@code SEND} is refused, as if no terminal were a configured lead,
* collaborator, or observer target — the same decision {@link #NO_KNOWN_LEAD_OR_COLLABORATOR}
* and {@link #NO_KNOWN_OBSERVER_TARGET} give explicitly. Every other action's result is
* identical to the five-argument form's, since none of them consult either classifier.
* or an observer's {@code SEND}, and an observer's or a collaborator's {@code TASK_READ}, are
* refused, as if no terminal were a configured lead, collaborator, or observer-reachable
* target, and no ticket were the caller's own — the same decision
* {@link #NO_KNOWN_LEAD_OR_COLLABORATOR}, {@link #NO_OBSERVER_SEND_TARGET} and
* {@link #NO_OWNED_TICKET} give explicitly. Every other action's result is identical to the
* full form's, since none of them consult any of the three classifiers.
*
* <p>Its default classifiers deny every collaborator and every observer, so a caller
* enforcing authorization must use the five-argument form instead.
* enforcing authorization must use the full form instead.
*/
public static boolean permits(Principal caller, Action action, String targetSession) {
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR, NO_KNOWN_OBSERVER_TARGET);
return permits(caller, action, targetSession, NO_KNOWN_LEAD_OR_COLLABORATOR,
NO_OBSERVER_SEND_TARGET, NO_OWNED_TICKET);
}
/**
* As {@link #permits(Principal, Action, String)}, with a real classifier for a collaborator's
* {@code SEND}. An observer's {@code SEND} still fails closed ({@link #NO_KNOWN_OBSERVER_TARGET}) —
* a caller enforcing both grants must use the five-argument form.
* {@code SEND}. An observer's {@code SEND}, and an observer's or a collaborator's
* {@code TASK_READ}, still fail closed — a caller enforcing all three grants must use the
* full form.
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator) {
return permits(caller, action, targetSession, knownLeadOrCollaborator, NO_KNOWN_OBSERVER_TARGET);
return permits(caller, action, targetSession, knownLeadOrCollaborator,
NO_OBSERVER_SEND_TARGET, NO_OWNED_TICKET);
}
/**
* As {@link #permits(Principal, Action, String, Predicate)}, with a real classifier for an
* observer's {@code SEND} too. An observer's or a collaborator's {@code TASK_READ} still
* fails closed — a caller enforcing all three grants must use the full form.
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> observerSendTarget) {
return permits(caller, action, targetSession, knownLeadOrCollaborator, observerSendTarget,
NO_OWNED_TICKET);
}
/**
@@ -115,21 +141,29 @@ public final class Authz {
*
* @param targetSession the session id in the request path; only consulted for the
* worker-scoped actions ({@code REPLY}, {@code ASK},
* {@code INBOX}), for a collaborator's {@code SEND}, and for an
* observer's {@code SEND}, ignored otherwise, may be
* {@code null}
* {@code INBOX}), for a collaborator's {@code SEND}, for an
* observer's {@code SEND}, and — carrying a ticket id instead
* of a session id — for an observer's or a collaborator's
* {@code TASK_READ}, ignored otherwise, may be {@code null}
* @param knownLeadOrCollaborator whether a terminal is a configured lead or collaborator —
* consulted only for a collaborator's {@code SEND}, to confine
* it to another named peer and never a spawned member's
* terminal
* @param knownObserverTarget whether a terminal is one this daemon would itself resolve as
* {@link Role#OBSERVER} — consulted only for an observer's
* {@code SEND}, to confine it to another observer pane and never
* a lead, a collaborator, or a spawned member
* @param observerSendTarget whether a terminal is one this daemon would itself resolve as
* a lead ({@link Role#PRIMARY}) or as {@link Role#OBSERVER} —
* consulted only for an observer's {@code SEND}, to confine it
* to a lead or another observer pane and never a collaborator,
* an architect, or a spawned member
* @param ticketOwnedByCaller whether {@code targetSession} (here, a ticket id) was created
* by this same caller — consulted only for an observer's or a
* collaborator's {@code TASK_READ}, to confine it to a ticket
* its own {@code fleet_send(wait:false)} created, never another
* caller's
*/
public static boolean permits(Principal caller, Action action, String targetSession,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> knownObserverTarget) {
Predicate<String> observerSendTarget,
Predicate<String> ticketOwnedByCaller) {
if (caller == null || caller.isAnonymous()) {
return false; // authenticated as nothing ⇒ authorized for nothing
}
@@ -143,12 +177,12 @@ public final class Authz {
// Delivering a turn to a local session is open to the primary and the architect
// unconditionally. A collaborator may reach only a target that is itself a configured
// lead or collaborator, never a spawned member's terminal. An observer may reach only
// a target that would itself resolve as an observer, never a lead, a collaborator, or
// a spawned member. A worker is excluded from every case — sending would be it
// escalating into the orchestrator role.
// a target that would itself resolve as a lead or as another observer, never a
// collaborator, an architect, or a spawned member. A worker is excluded from every
// case — sending would be it escalating into the orchestrator role.
case SEND -> caller.isPrimary() || caller.isArchitect()
|| (caller.isCollaborator() && knownLeadOrCollaborator.test(targetSession))
|| (caller.isObserver() && knownObserverTarget.test(targetSession));
|| (caller.isObserver() && observerSendTarget.test(targetSession));
// Resolving a worker's blocked question is part of delegating to it, open to the same
// two roles that may stand up that delegation in the first place. Not a collaborator:
@@ -178,11 +212,17 @@ public final class Authz {
|| caller.isCollaborator() || caller.isObserver();
// Ticket polling and session status, open to every role READ is open to except a
// collaborator or an observer. MessageService compares a ticket's creator to the
// caller on every read as well, so dropping this gate would not expose another
// session's reply — it would move the refusal later and widen what a caller that
// never orchestrates can probe.
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect();
// collaborator or an observer — with one exception: an observer or a collaborator may
// poll a ticket its own fleet_send(wait:false) created, confined by
// ticketOwnedByCaller. Session status is not ticket-scoped, so that call site
// supplies no real classifier here and an observer's or a collaborator's TASK_READ on
// it stays refused. MessageService compares a ticket's creator to the caller on every
// read as well, so dropping this gate would not expose another session's reply — it
// would move the refusal later and widen what a caller that never orchestrates can
// probe.
case TASK_READ -> caller.isPrimary() || caller.isWorker() || caller.isArchitect()
|| ((caller.isObserver() || caller.isCollaborator())
&& ticketOwnedByCaller.test(targetSession));
// fleetd #421: reading held lead-to-lead mail is the primary's alone. An architect
// holds READ today (CB-548), so "not primary" must mean not-architect here too — this
@@ -271,22 +271,27 @@ public final class CallerResolver {
}
/**
* Whether {@code target} names a terminal this resolver would itself resolve as {@link
* Role#OBSERVER} — the classifier an observer's {@code SEND} is checked against, read from the
* same maps and functions {@link #resolve} consults so a target this accepts is exactly one
* {@code resolve} would hand back {@link Role#OBSERVER} for, and the reverse.
* Whether {@code target} names a terminal an observer may {@code SEND} to: one this resolver
* would itself resolve as a lead ({@link Role#PRIMARY}) or as {@link Role#OBSERVER}. Read from
* the same maps and functions {@link #resolve} consults, and in the same order, so a target
* this accepts is exactly one {@code resolve} would hand back one of those two roles for, and
* the reverse.
*
* <p>A live spawned member is refused first, whatever a tab map says about its terminal — the
* order {@link #resolve} itself uses. A pane named as a lead is then accepted even when it is
* also bound to an architect slot, because that is the role {@code resolve} gives it.
*/
public Predicate<String> sendableObserverTarget() {
public Predicate<String> observerSendTarget() {
return target -> target != null
&& spawnedMemberRole.apply(target) == null
&& !leadTerminals.get().containsKey(target)
&& !boundToArchitectSlot(target)
&& !collaboratorTerminals.get().containsKey(target);
&& (leadTerminals.get().containsKey(target)
|| (!boundToArchitectSlot(target)
&& !collaboratorTerminals.get().containsKey(target)));
}
/**
* Whether {@code terminal} is bound to a configured slot the live roster still confirms as an
* architect — the one classifier {@link #sendableObserverTarget()} and {@code FleetMcp}'s
* architect — the one classifier {@link #observerSendTarget()} and {@code FleetMcp}'s
* {@code panes} row both read, so a pane's reported role and its {@code SEND} reachability can
* never drift apart.
*/
@@ -40,8 +40,9 @@ public enum Role {
* fleet.collaborators.<name>.tab} registry). Never spawned — identity comes from the
* connection, never a request argument, exactly like {@link #WORKER} and {@link #ARCHITECT}.
* May {@code SEND} only to a configured lead or collaborator, {@code REPLY}/{@code ASK} only
* as its own pane, and {@code READ}/{@code METRICS}; may not {@code SPAWN}/{@code STOP}/
* {@code DRAIN}/{@code HANDOVER}, poll a ticket ({@code TASK_READ}), or reach the
* as its own pane, {@code READ}/{@code METRICS}, and {@code TASK_READ} only a ticket its own
* {@code fleet_send(wait:false)} created; may not {@code SPAWN}/{@code STOP}/{@code DRAIN}/
* {@code HANDOVER}, read a session's status or another caller's ticket, or reach the
* coordination broker ({@code COORD_SEND}/{@code COORD_READ}).
*/
COLLABORATOR,
@@ -51,10 +52,11 @@ public enum Role {
* configured lead, not a bound architect slot, not a configured collaborator tab. Unforgeable
* like a worker's — derived from the connection's pane, never from a request argument, and
* honoured regardless of auth mode. May {@code READ} and {@code METRICS}, {@code REPLY}/
* {@code ASK} only as its own pane, and {@code SEND} only to a target that would itself
* resolve as {@code OBSERVER}; may not {@code SPAWN}/{@code STOP}/{@code DRAIN}/
* {@code HANDOVER}, poll a ticket ({@code TASK_READ}), or reach the coordination broker
* ({@code COORD_SEND}/{@code COORD_READ}).
* {@code ASK} only as its own pane, {@code SEND} only to a target that would itself resolve
* as a lead ({@link #PRIMARY}) or as {@code OBSERVER}, and {@code TASK_READ} only a ticket its
* own {@code fleet_send(wait:false)} created; may not {@code SPAWN}/{@code STOP}/
* {@code DRAIN}/{@code HANDOVER}, read a session's status or another caller's ticket, or reach
* the coordination broker ({@code COORD_SEND}/{@code COORD_READ}).
*/
OBSERVER,
@@ -183,8 +183,9 @@ import java.util.function.Supplier;
* recounted again for fleetd #362, again after {@code idleSleepGuard:} was added, again after
* {@code models:} was added as deferred, again for fleetd #422, which moved {@code models:}
* from deferred to hot-excluded once its on/off half was read live everywhere, and again after
* {@code leadRollover:} was added (fleetd #480).</strong>
* {@code FleetConfig} has 26 top-level record components: 5 cold, 13 deferred, 3 split, 5
* {@code leadRollover:} was added (fleetd #480), and again after {@code promptBoxGateEnabled:}
* was added as deferred.</strong>
* {@code FleetConfig} has 27 top-level record components: 5 cold, 14 deferred, 3 split, 5
* hot-excluded. Five of them are named nowhere in this file, and the reason is the same for all
* five: {@code placement}, {@code memberCredentials}, {@code memberLoginShell}, {@code models} and
* {@code leadRollover} are <strong>hot</strong> and correctly absent — all five are read live off
@@ -288,7 +289,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
static final Set<String> DEFERRED_KEYS = Set.of(
"guard", "worktreeRoot", "worktreeGroup", "memberSkills", "primary", "configReload",
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
"quarantineCooldownSeconds", "profiles", "idleSleepGuard");
"quarantineCooldownSeconds", "profiles", "idleSleepGuard", "promptBoxGateEnabled");
private final Path path;
private final AtomicReference<FleetConfig> current;
@@ -531,6 +532,12 @@ public final class ConfigRef implements Supplier<FleetConfig> {
if (!Objects.equals(old.idleSleepGuard(), fresh.idleSleepGuard())) {
changed.add("idleSleepGuard");
}
// ReplyPushLoop, LeadCoordLoop and LeadHeartbeatLoop each build their own PromptBox once, at
// startup, off this value — a running loop keeps whichever gate setting it was built with
// regardless of a later edit here.
if (!Objects.equals(old.promptBoxGateEnabled(), fresh.promptBoxGateEnabled())) {
changed.add("promptBoxGateEnabled");
}
if (!Objects.equals(old.spawnReadyTimeoutMs(), fresh.spawnReadyTimeoutMs())
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
changed.add("spawnReady*");
@@ -143,6 +143,13 @@ import java.util.regex.PatternSyntaxException;
* @param leadRollover opt-in lead rollover (fleetd #480): {@code null} ⇒ off, and no
* {@code dev.ltms.fleet.lead.LeadRollover} is constructed at all — an upgraded
* daemon never clears a lead's pane on its own initiative. See {@link LeadRollover}.
* @param promptBoxGateEnabled whether {@link dev.ltms.fleet.herdr.PromptBox} checks a lead's input
* box before a delivery pastes into it. {@code null} (the default) and
* {@code false} both leave the check off: a delivery proceeds without reading
* the pane, and the box is never classified. Detection currently misreads the
* pane's own autocomplete suggestion at the caret as the operator's
* unsubmitted typing, so the check holds deliveries that were never actually
* at risk; set this to {@code true} only once that is fixed.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record FleetConfig(
@@ -171,7 +178,23 @@ public record FleetConfig(
String memberSkills,
IdleSleepGuard idleSleepGuard,
Models models,
LeadRollover leadRollover) {
LeadRollover leadRollover,
Boolean promptBoxGateEnabled) {
/** Back-compat form before the {@code promptBoxGateEnabled} key was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
Guard guard, String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials, Coordinator coordinator, String worktreeGroup,
String memberLoginShell, String memberSkills, IdleSleepGuard idleSleepGuard,
Models models, LeadRollover leadRollover) {
this(bind, herdrSocket, memberHerdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, coordinator, worktreeGroup,
memberLoginShell, memberSkills, idleSleepGuard, models, leadRollover, null);
}
/** Back-compat form before the {@code leadRollover:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, String memberHerdrSocket, Map<String, Profile> profiles,
@@ -331,6 +354,7 @@ public record FleetConfig(
} else {
profiles = Map.of();
}
promptBoxGateEnabled = promptBoxGateEnabled != null && promptBoxGateEnabled;
}
@JsonIgnoreProperties(ignoreUnknown = true)
@@ -1466,7 +1490,7 @@ public record FleetConfig(
* CALLING lead's own turn to end (its pane to report {@code IDLE} or
* {@code DONE}) before ending that pane's process at all. See the
* paragraph above.
* @param relaunchReadySeconds default 45 — bound on EACH of two separate waits that run after
* @param relaunchReadySeconds default 45 — bound on EACH of three separate waits that run after
* the old lead's pane has been torn down and a fresh one launched: first,
* for the fresh pane itself to reach a real turn boundary ({@code IDLE} or
* {@code DONE}, never merely {@code BLOCKED}) — the safety gate, since
@@ -1479,8 +1503,9 @@ public record FleetConfig(
* 10s live), so a budget has to clear more than one scan interval to leave
* any real margin for the CLI's own boot time; 20 was rejected for exactly
* that reason — at a 10s scan interval it only buys two scans. 45 buys
* roughly four. Only a timeout on the FIRST wait (the pane never becomes
* ready) withholds {@code bootstrapText}.
* roughly four; third, to retry sending {@code bootstrapText} while herdr
* reports {@code agent_not_ready}. A timeout on the first wait or the third
* one withholds {@code bootstrapText}.
* @param bootstrapText default a sentence naming the RESOLVED handover path — sent to the
* fresh lead's pane once it reaches a real turn boundary after relaunch,
* telling the fresh session where to read the handover and carry on. Left
@@ -1936,7 +1961,7 @@ public record FleetConfig(
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials", "coordinator", "worktreeGroup", "memberLoginShell", "memberSkills",
"idleSleepGuard", "models", "leadRollover");
"idleSleepGuard", "models", "leadRollover", "promptBoxGateEnabled");
/** Load and validate config from {@code path}. */
public static FleetConfig load(Path path) {
@@ -2733,10 +2758,13 @@ public record FleetConfig(
// own compact constructor defaults the fields of a block that IS present. Defaulting it
// here would construct a LeadRollover object (via Fleetd.java's presence gate) for every
// config that never mentioned it.
// promptBoxGateEnabled is left as-is: the compact constructor already normalized it to a
// real true/false, and false (the gate off) is already the safe default — there is no
// further defaulting to do here.
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown, mc, coordinator, worktreeGroup, memberLoginShell, memberSkills,
idleSleepGuard, models, leadRollover);
idleSleepGuard, models, leadRollover, promptBoxGateEnabled);
}
/**
@@ -26,6 +26,9 @@ import java.util.regex.Pattern;
* <p>A pane that holds for {@link #HOLD_WARN_STREAK} consecutive checks gets one warning, so a box
* that never clears is visible instead of silent. The warning repeats only after the box has cleared
* again.
*
* <p>The check itself can be switched off (see the two-arg constructor): a disabled gate clears
* every target without reading its pane, logging a hold, or counting a streak.
*/
public final class PromptBox {
@@ -74,19 +77,35 @@ public final class PromptBox {
}
private final AgentControl agents;
private final boolean gateEnabled;
/** Consecutive holds per target, so a box that never clears can be warned about once. */
private final Map<String, Integer> holdStreaks = new ConcurrentHashMap<>();
/** Equivalent to {@link #PromptBox(AgentControl, boolean)} with the gate on. */
public PromptBox(AgentControl agents) {
this(agents, true);
}
/**
* @param gateEnabled {@code false} makes {@link #clearToSubmit} return {@code true} for every
* target without reading its pane, logging a hold, or counting a streak —
* every delivery proceeds exactly as if the box were always empty.
*/
public PromptBox(AgentControl agents, boolean gateEnabled) {
this.agents = agents;
this.gateEnabled = gateEnabled;
}
/**
* Whether {@code target}'s input box is empty, so a delivery would submit only its own text.
* {@code false} means hold and come back; it never means the delivery failed.
* {@code false} means hold and come back; it never means the delivery failed. Always
* {@code true} when the gate is disabled.
*/
public boolean clearToSubmit(String target) {
if (!gateEnabled) {
return true;
}
Reading reading = inspect(target);
if (reading.state() == State.EMPTY) {
holdStreaks.remove(target);
@@ -232,7 +232,7 @@ public final class LeadRollover {
/**
* What is known about one token, right now — the answer {@link #status} gives. Distinguishes
* five terminal outcomes an approved roll can finish with, one in-flight outcome for a roll
* terminal outcomes an approved roll can finish with, one in-flight outcome for a roll
* that has been approved but has not finished yet, and two answers for a token that names no
* active work at all: still pending confirmation, or nothing known about this token at all.
*/
@@ -256,7 +256,8 @@ public final class LeadRollover {
* that is, in fact, actively running. This is not sticky: the deferred continuation
* overwrites this same entry with a terminal state ({@link #ROLLED}, {@link
* #TURN_NEVER_SETTLED}, {@link #OLD_PANE_NEVER_DIED}, {@link #RELAUNCH_FAILED}, {@link
* #RELAUNCH_NEVER_READY}, {@link #RELAUNCH_NOT_RECOGNISED}, or {@link #FAILED}) once it
* #RELAUNCH_NEVER_READY}, {@link #RELAUNCH_NOT_RECOGNISED}, {@link #BOOTSTRAP_NEVER_SENT},
* or {@link #FAILED}) once it
* finishes — including by throwing, which {@link #runRollover}'s catch turns into {@link
* #FAILED} instead of leaving this entry stuck forever.
*/
@@ -305,6 +306,12 @@ public final class LeadRollover {
* an operator should check why the tab was not recognised.
*/
RELAUNCH_NOT_RECOGNISED,
/**
* A fresh lead was ready, but herdr kept reporting {@code agent_not_ready} while this class
* retried {@code bootstrapText} for {@code relaunchReadySeconds}. The fresh session did not
* receive its handover instruction.
*/
BOOTSTRAP_NEVER_SENT,
/**
* The deferred continuation threw a {@link RuntimeException} and the continuation thread
* died with it. Without this state, that throw would leave {@link #outcomes} holding {@link
@@ -731,7 +738,19 @@ public final class LeadRollover {
IdentityResult identityResult = waitUntilRecognisedAsLead(newAgent.terminalId(),
cfg.relaunchReadySeconds());
agents.send(newAgent.terminalId(), cfg.bootstrapTextFor(p.handoverPath()));
BootstrapResult bootstrapResult = sendBootstrapWithRetry(newAgent.terminalId(),
cfg.bootstrapTextFor(p.handoverPath()), cfg.relaunchReadySeconds());
if (!bootstrapResult.sent()) {
log.warn("lead-rollover: bootstrapText was never sent to fresh terminal {} for lead '{}' "
+ "after agent_not_ready persisted for {}ms (token={}, configured={}s)",
newAgent.terminalId(), leadName, bootstrapResult.elapsedMillis(), p.token(),
cfg.relaunchReadySeconds());
outcomes.put(p.token(), new RollStatus(RollState.BOOTSTRAP_NEVER_SENT,
"fresh terminal " + newAgent.terminalId() + " kept rejecting bootstrapText with "
+ "agent_not_ready for relaunchReadySeconds=" + cfg.relaunchReadySeconds()
+ "s (measured elapsed=" + bootstrapResult.elapsedMillis() + "ms)"));
return;
}
if (!identityResult.ready()) {
log.warn("lead-rollover: fresh terminal {} for lead '{}' is alive and bootstrapped, but "
+ "was never recognised as a live lead — an operator should check why "
@@ -755,6 +774,30 @@ public final class LeadRollover {
+ newAgent.terminalId()));
}
/**
* Sends {@code bootstrapText}, retrying only the transient herdr {@code agent_not_ready} refusal
* until {@code readySeconds} elapses. All other failures propagate to {@link #runRollover}.
*/
private BootstrapResult sendBootstrapWithRetry(String terminal, String bootstrapText, int readySeconds) {
long startMillis = nowMillis.getAsLong();
long deadline = startMillis + TimeUnit.SECONDS.toMillis(readySeconds);
while (nowMillis.getAsLong() < deadline) {
try {
agents.send(terminal, bootstrapText);
return new BootstrapResult(true, nowMillis.getAsLong() - startMillis);
} catch (HerdrException e) {
if (!"agent_not_ready".equals(e.code())) {
throw e;
}
pollSleeper.run();
}
}
return new BootstrapResult(false, nowMillis.getAsLong() - startMillis);
}
/** The result of {@link #sendBootstrapWithRetry}. */
private record BootstrapResult(boolean sent, long elapsedMillis) {}
/** Attempts {@link #captureAgentWithRetry} makes before letting the failure propagate. */
static final int CAPTURE_RETRIES = 3;
@@ -121,9 +121,9 @@ public final class FleetMcp {
/**
* Kept as a field (rather than only captured by the {@code contextExtractor} closure) so
* {@link #denyFor} can read {@link CallerResolver#knownLeadOrCollaborator()} and {@link
* CallerResolver#sendableObserverTarget()} — the classifiers a collaborator's and an
* observer's {@code SEND} are each checked against, built from the same maps {@link #identity}-
* based resolution reads.
* CallerResolver#observerSendTarget()} — the classifiers a collaborator's and an observer's
* {@code SEND} are each checked against, built from the same maps {@link #identity}-based
* resolution reads.
*/
private final CallerResolver callers;
private final Metrics metrics; // CB-502: null → auth failures not counted
@@ -323,12 +323,12 @@ public final class FleetMcp {
* injector, keyed by terminal id — never a second, separately-derived check
* @param architectSlot the same classifier {@link CallerResolver#boundToArchitectSlot} resolves
* a caller against — never a second, separately-derived check
* @param sendableToObserver the same predicate {@link CallerResolver#sendableObserverTarget}
* @param observerSendTarget the same predicate {@link CallerResolver#observerSendTarget}
* builds for the {@code SEND} gate — never a second, separately-derived
* check
*/
public record PaneSource(Supplier<Map<String, String>> tabLabels, Supplier<Map<String, String>> workspaceLabels,
Predicate<String> deliverable, Predicate<String> architectSlot, Predicate<String> sendableToObserver) {
Predicate<String> deliverable, Predicate<String> architectSlot, Predicate<String> observerSendTarget) {
/** Inert source — no labels, no deliverable targets, no architect slots, nothing sendable. */
public static PaneSource none() {
return new PaneSource(Map::of, Map::of, _ -> false, _ -> false, _ -> false);
@@ -516,7 +516,10 @@ public final class FleetMcp {
// Answering a worker's fleet_ask (CB-205): resolve its blocked question and
// block for the worker's reply as it resumes the same turn. This is the same
// delegation, so ownership is left untouched (CB-548) — never re-recorded.
return answer(messages, turnId, content, timeoutMs(a), callerOwner);
// ANSWER is primary-or-architect, and an architect is refused DRAIN too, so a
// timeout receipt must not name fleet_poll{target} to an architect caller.
boolean mayDrainPoll = Authz.permits(caller, Authz.Action.DRAIN, target);
return answer(messages, turnId, content, timeoutMs(a), callerOwner, mayDrainPoll);
}
// CB-548: delegator ownership (which lead's reply nudge this worker routes to,
// CB-532) is recorded only once the send is ACCEPTED — MessageService has won the
@@ -528,12 +531,18 @@ public final class FleetMcp {
// name would never match there anyway, but passing null for them keeps the intent
// explicit rather than relying on that lookup to filter it out.
String delegatorName = caller.isPrimary() ? caller.name() : null;
Runnable onAccepted = () ->
primaryRegistry.recordDelegation(target, callerTerminal, delegatorName);
// A nudge about this target must never tell its recipient to run a call Authz
// would refuse it — asked once, here, while the real Principal is still in
// scope, rather than re-derived from a bare terminal string later.
boolean mayDrainNudge = Authz.permits(caller, Authz.Action.DRAIN, target);
boolean mayAnswerNudge = Authz.permits(caller, Authz.Action.ANSWER, target);
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, callerTerminal,
delegatorName, mayDrainNudge, mayAnswerNudge);
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
return Boolean.FALSE.equals(a.get("wait"))
? sendAsync(messages, target, content, onAccepted, workers.profiles(), caller)
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), callerOwner);
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles(), callerOwner,
mayDrainNudge);
};
// fleet_reply's identity is the CONNECTION, never an argument. The authz check is
// terminal ownership, not a role test: the caller may reply only for its own pane,
@@ -574,10 +583,19 @@ public final class FleetMcp {
Map<String, Object> a = req.arguments();
String target = str(a, "target");
String coordId = str(a, "coordId");
String ticket = str(a, "ticket");
// The action depends on the ARGUMENTS, not on the tool name -- see pollAction.
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_poll", a), target);
Authz.Action action = toolAction("fleet_poll", a);
// TASK_READ reads a ticket, not a DRAIN target -- the gate must see the
// ticket id, and an observer's or a collaborator's grant is confined to a
// ticket it created.
String authzTarget = action == Authz.Action.TASK_READ ? ticket : target;
Predicate<String> ticketOwnedByCaller = action == Authz.Action.TASK_READ
? t -> messages.ownsTicket(t, principal(exchange).ownerKey())
: Authz.NO_OWNED_TICKET;
McpSchema.CallToolResult denied = deny(exchange, action, authzTarget, ticketOwnedByCaller);
if (denied != null) return denied;
return poll(messages, leadChannel, str(a, "ticket"), target, coordId,
return poll(messages, leadChannel, ticket, target, coordId,
principal(exchange).ownerKey());
};
// CB-307 Increment 3: per-msgId ack (not needed in v1 but supported by the inbox).
@@ -616,7 +634,7 @@ public final class FleetMcp {
PaneSource panes = new PaneSource(() -> identity.panes().tabLabelsByTabId(),
() -> identity.panes().workspaceLabelsByWorkspaceId(),
Fleetd.deliverableTo(presence, callers::leads, callers::collaborators),
callers::boundToArchitectSlot, callers.sendableObserverTarget());
callers::boundToArchitectSlot, callers.observerSendTarget());
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, leadContextGauge, leadConfigDirs, callers.leads(),
callerTerminal(exchange),
@@ -745,6 +763,16 @@ public final class FleetMcp {
return denyFor(principal(exchange), action, target);
}
/**
* As {@link #deny(McpSyncServerExchange, Authz.Action, String)}, with a real classifier for
* an observer's or a collaborator's {@code TASK_READ} — the only call site that can supply
* one is a ticket poll, which knows the ticket id and the caller's owner key.
*/
private McpSchema.CallToolResult deny(McpSyncServerExchange exchange, Authz.Action action,
String target, Predicate<String> ticketOwnedByCaller) {
return denyFor(principal(exchange), action, target, ticketOwnedByCaller);
}
/**
* The policy half of {@link #deny}: everything except pulling the caller out of the MCP
* exchange. Kept separate so the authorization decision — the actual control — is unit-testable
@@ -758,13 +786,22 @@ public final class FleetMcp {
* @return {@code null} when the call may proceed, or the error result to return when it may not
*/
McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target) {
return denyFor(caller, action, target, Authz.NO_OWNED_TICKET);
}
/**
* As {@link #denyFor(Principal, Authz.Action, String)}, with a real classifier for an
* observer's or a collaborator's {@code TASK_READ}.
*/
McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target,
Predicate<String> ticketOwnedByCaller) {
// The enforcement switch lives HERE rather than in the exchange-facing wrapper: any future
// tool that calls this directly must not be able to skip the gate by accident.
if (!authorizationEnforced) {
return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518)
}
if (Authz.permits(caller, action, target, callers.knownLeadOrCollaborator(),
callers.sendableObserverTarget())) {
callers.observerSendTarget(), ticketOwnedByCaller)) {
if (action != Authz.Action.READ && action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
}
@@ -826,13 +863,15 @@ public final class FleetMcp {
}
/**
* Who may see {@code fleet_list}'s {@code leads} array — exactly the roles that may
* {@link Authz.Action#SEND} to a lead: the primary, an architect, and a collaborator. A
* collaborator's own {@code fleet_whoami} carries no lead address, and {@code leads} is the
* only place this tool gives one, so a collaborator needs this array to use the send it
* already holds. A worker can never {@code SEND} at all, so it still sees neither this array
* nor {@code members}; a worker's own facts come from {@code fleet_whoami} instead. Split
* out for the same reason as {@link #coordinatorVisibleTo} and
* Who may see {@code fleet_list}'s {@code leads} array — the primary, an architect, and a
* collaborator. A collaborator's own {@code fleet_whoami} carries no lead address, and this is
* the only place this tool gives one, so a collaborator needs this array to use the
* {@link Authz.Action#SEND} it already holds. An observer holds that send to a lead too, but
* learns the address from its filtered {@code panes} rows instead: a {@code leads} row carries
* a lead's name, its configured context window and its config dir, which are the fleet's own
* shape rather than an address. A worker can never {@code SEND} at all, so it still sees
* neither this array nor {@code members}; a worker's own facts come from {@code fleet_whoami}
* instead. Split out for the same reason as {@link #coordinatorVisibleTo} and
* {@link #collaboratorsVisibleTo}: the decision must be unit-testable without fabricating an
* SDK {@code McpSyncServerExchange}, and the handler must call this named predicate rather
* than inlining the check.
@@ -855,9 +894,10 @@ public final class FleetMcp {
/**
* Who may see {@code fleet_list}'s {@code panes} array — every role that may
* {@link Authz.Action#SEND} to some other pane. A plain worker holds {@code READ} but never
* {@code SEND}, so it still does not see this array. An observer does hold {@code SEND}, to
* another observer pane only, so it sees the array too — but {@code listFleet} filters its rows
* to {@link CallerResolver#sendableObserverTarget} and reduces each one; see {@code paneRows}.
* {@code SEND}, so it still does not see this array. An observer does hold {@code SEND}, to a
* lead or another observer pane, so it sees the array too — and it is the only place this tool
* gives it a lead's address. {@code listFleet} filters its rows to
* {@link CallerResolver#observerSendTarget} and reduces each one; see {@code paneRows}.
*/
static boolean panesVisibleTo(Principal caller) {
return caller.isPrimary() || caller.isArchitect() || caller.isCollaborator() || caller.isObserver();
@@ -996,12 +1036,17 @@ public final class FleetMcp {
* {@code fleet_send}: delegate {@code content} to a worker session and block for its reply.
* The configured profiles are required so a profile name can never bypass target validation.
*
* (CB-548): {@code onAccepted} records delegator ownership the instant the send is accepted, so
* a BUSY interloper never claims a turn it did not win. {@code null} disables recording.
* <p>(CB-548): {@code onAccepted} records delegator ownership the instant the send is
* accepted, so a BUSY interloper never claims a turn it did not win. {@code null} disables
* recording.
*
* @param mayDrainPoll whether this caller's {@code Authz.Action.DRAIN} is granted — a
* TIMED_OUT or BUSY receipt must name {@code fleet_poll} only when this
* is {@code true}, since that route is refused otherwise
*/
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
Long timeoutMs, Runnable onAccepted, Set<String> profiles,
String callerOwner) {
String callerOwner, boolean mayDrainPoll) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
@@ -1011,7 +1056,8 @@ public final class FleetMcp {
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
try {
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerOwner), timeout);
return formatReply(messages.send(sessionId, content, timeout, onAccepted, callerOwner), timeout,
sessionId, mayDrainPoll);
} catch (HerdrException e) {
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
}
@@ -1022,14 +1068,20 @@ public final class FleetMcp {
* {@code fleet_ask} (CB-205). Resolves the worker's blocked question and blocks for its reply as
* it resumes the same turn — surfaced to the primary identically to a normal send.
* {@code callerOwner} must match the turn's recorded owner or this is refused.
*
* @param mayDrainPoll whether this caller's {@code Authz.Action.DRAIN} is granted — see
* {@link #send(MessageService, String, String, Long, Runnable, Set, String,
* boolean)}'s same parameter. {@code ANSWER} is primary-or-architect, and
* an architect is refused {@code DRAIN}, so this cannot default to true.
*/
static McpSchema.CallToolResult answer(MessageService messages, String turnId, String content, Long timeoutMs,
String callerOwner) {
String callerOwner, boolean mayDrainPoll) {
if (isBlank(turnId) || isBlank(content)) {
return error("turnId and content are required to answer a worker's question");
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
return formatReply(messages.answer(turnId, content, timeout, callerOwner), timeout);
String target = messages.sessionForTurn(turnId);
return formatReply(messages.answer(turnId, content, timeout, callerOwner), timeout, target, mayDrainPoll);
}
/**
@@ -1056,8 +1108,15 @@ public final class FleetMcp {
};
}
/** Render a {@link MessageService.Reply} as a tool result — shared by {@link #send} and {@link #answer}. */
private static McpSchema.CallToolResult formatReply(MessageService.Reply r, long timeout) {
/**
* Render a {@link MessageService.Reply} as a tool result — shared by {@link #send} and
* {@link #answer}. {@code target} is the worker session the send went to (may be {@code null}
* if it could not be resolved), named in a TIMED_OUT_* receipt so the caller knows which
* inbox to poll — but only when {@code mayDrainPoll} is {@code true}: naming a route the
* caller's own {@code Authz.Action.DRAIN} grant would refuse is worse than naming none.
*/
private static McpSchema.CallToolResult formatReply(MessageService.Reply r, long timeout, String target,
boolean mayDrainPoll) {
return switch (r.outcome()) {
case REPLIED -> text(r.text());
// The worker's turn finished but it never called fleet_reply — hand back the scraped
@@ -1080,8 +1139,16 @@ public final class FleetMcp {
+ "answered (turnId stale)");
case NOT_TURN_OWNER -> error("this turn belongs to a different delegation — only the caller "
+ "that opened it may answer it");
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker "
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> {
String base = "[no reply within " + timeout + "ms — worker "
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + ". "
+ MessageService.NO_TICKET_NO_RESEND + ". ";
yield text(mayDrainPoll
? base + "Poll fleet_poll{target=\""
+ (target != null ? target : "<the worker session you sent to>")
+ "\"} to drain its inbox for a late reply.]"
: base + "A late reply cannot be recovered on this channel — ask the lead.]");
}
// fleetd #571: delivery is unknown here — agent.prompt pastes and submits in one call,
// so the message may already be sitting in the pane. Do not invite a blind retry the way
// the case above does; a resend on this route can double-deliver the same brief.
@@ -2155,8 +2222,8 @@ public final class FleetMcp {
}
/**
* As above, plus fleetd #758: an observer sees the {@code panes} array too, but filtered to
* {@link CallerResolver#sendableObserverTarget} and each row reduced to the five fields an
* As above, plus: an observer sees the {@code panes} array too, but filtered to
* {@link CallerResolver#observerSendTarget} and each row reduced to the five fields an
* observer may learn — see {@code paneRows}/{@code paneRow}.
*
* @param callerIsObserver whether the {@code fleet_list} caller is an observer; every wrapper
@@ -2579,9 +2646,10 @@ public final class FleetMcp {
* already read, so a pane neither configured as a lead nor spawned as a member — a hand-opened
* tab — still gets a row here.
*
* <p>fleetd #758: for an observer caller ({@code observerView}), the rows are filtered to
* {@link PaneSource#sendableToObserver} before being built, and each row is reduced — see
* {@code paneRow}.
* <p>For an observer caller ({@code observerView}), the rows are filtered to
* {@link PaneSource#observerSendTarget} before being built, and each row is reduced — see
* {@code paneRow}. A lead's pane survives that filter, so it is where an observer reads a
* lead's {@code sessionId}.
*/
private static List<Map<String, Object>> paneRows(Map<String, Agent> live, List<MemberSession> roster,
Map<String, String> leads, Map<String, String> collaborators, PaneSource panes,
@@ -2592,7 +2660,7 @@ public final class FleetMcp {
.filter(s -> s.terminalId() != null)
.collect(Collectors.toMap(MemberSession::terminalId, Function.identity(), (_, b) -> b));
return live.values().stream()
.filter(a -> !observerView || panes.sendableToObserver().test(a.terminalId()))
.filter(a -> !observerView || panes.observerSendTarget().test(a.terminalId()))
.sorted(Comparator.comparing(Agent::terminalId))
.map(a -> paneRow(a, byTerminal.get(a.terminalId()), leads, collaborators, tabLabels,
workspaceLabels, panes, observerView))
@@ -2607,9 +2675,9 @@ public final class FleetMcp {
* @param workspaceLabels workspace id → its herdr display label (the space name); a workspace
* absent here, or carrying a {@code null} label itself, projects as a
* {@code null} "workspaceLabel"
* @param observerView fleetd #758: an observer's row carries only {@code sessionId}, {@code
* label}, {@code status}, {@code role}, {@code deliverable} — never {@code
* paneId} (the {@code fleet_stop} handle), {@code workspaceId},
* @param observerView an observer's row carries only {@code sessionId}, {@code label},
* {@code status}, {@code role}, {@code deliverable} — never {@code paneId}
* (the {@code fleet_stop} handle), {@code workspaceId},
* {@code workspaceLabel}, {@code tabId}, {@code agentType}, or {@code cwd}
* (a member's worktree path is the lead's business)
*/
@@ -2857,9 +2925,12 @@ public final class FleetMcp {
return tool(FleetTool.LIST.wireName(),
"List the whole fleet the bridge tracks, in two parts. 'members' is visible to "
+ "the primary and an architect only. 'leads' is visible to those two AND a "
+ "collaborator — exactly the roles that may fleet_send to a lead, so a "
+ "collaborator can learn a lead's sessionId before using the send it already "
+ "holds. A worker holds READ to call this tool at all, but gets neither "
+ "collaborator, so a collaborator can learn a lead's sessionId before using "
+ "the send it already holds. An observer may fleet_send to a lead too, but "
+ "reads that sessionId from its own 'panes' rows instead: a 'panes' row for "
+ "an observer is filtered to the panes it may send to — a lead's pane and "
+ "another observer's — and reduced to sessionId, label, status, role and "
+ "deliverable. A worker holds READ to call this tool at all, but gets neither "
+ "array, never an empty one; a worker reads its own session, "
+ "profile, state, worktree, branch and owner from fleet_whoami instead. "
+ "'leads' are your PEERS — other "
@@ -48,8 +48,13 @@ public final class PrimaryRegistry {
*/
private final ConcurrentHashMap<String, Delegation> leadByTarget = new ConcurrentHashMap<>();
/** A recorded delegator: the terminal learned from call traffic, and its name, if it has one. */
private record Delegation(String terminal, String name) {
/**
* A recorded delegator: the terminal learned from call traffic, its name if it has one, and
* whether a nudge may tell it to {@code fleet_poll(target=…)} or {@code fleet_send(turnId=…)}
* — the real authorization decision for this delegator's role, captured once by the caller at
* delegation time rather than re-derived from a bare terminal string later.
*/
private record Delegation(String terminal, String name, boolean mayDrainNudge, boolean mayAnswerNudge) {
}
/**
@@ -128,14 +133,28 @@ public final class PrimaryRegistry {
/**
* As {@link #recordDelegation(String, String)}, additionally recording the delegating lead's
* name when the caller carries one. See {@link #record(String, String)} for why the name
* matters.
* name when the caller carries one, and granting a nudge about this target everything a
* primary may be told to run. See {@link #record(String, String)} for why the name matters,
* and {@link #recordDelegation(String, String, String, boolean, boolean)} for why every other
* caller must supply its own, real grant instead of this default.
*/
public void recordDelegation(String target, String leadTerminal, String leadName) {
recordDelegation(target, leadTerminal, leadName, true, true);
}
/**
* As {@link #recordDelegation(String, String, String)}, with the delegating caller's own
* {@code DRAIN} and {@code ANSWER} grants — the real {@code Authz.permits} decision for its
* role, asked once by the caller at delegation time, since a terminal string alone cannot be
* resolved back to a role here. A nudge about this target must offer only what
* {@link #nudgeMayDrainFor(String)}/{@link #nudgeMayAnswerFor(String)} report back as true.
*/
public void recordDelegation(String target, String leadTerminal, String leadName,
boolean mayDrainNudge, boolean mayAnswerNudge) {
if (target == null || target.isBlank() || leadTerminal == null || leadTerminal.isBlank()) {
return;
}
leadByTarget.put(target, new Delegation(leadTerminal, blankToNull(leadName)));
leadByTarget.put(target, new Delegation(leadTerminal, blankToNull(leadName), mayDrainNudge, mayAnswerNudge));
}
/** Forget a worker's delegating lead — call on release, so a torn-down session leaks nothing. */
@@ -167,6 +186,32 @@ public final class PrimaryRegistry {
return currentPrimaryTerminal();
}
/**
* Whether a nudge about {@code target} may tell its recipient to {@code fleet_poll(target=…)}
* — present on exactly the same condition as {@link #nudgeTargetFor(String)}, carrying the
* grant recorded for a delegation, or {@code true} for the singleton primary fallback, which
* is always genuinely the primary.
*/
public Optional<Boolean> nudgeMayDrainFor(String target) {
Delegation delegation = target == null ? null : leadByTarget.get(target);
if (delegation != null) {
return Optional.of(delegation.mayDrainNudge());
}
return currentPrimaryTerminal().isPresent() ? Optional.of(true) : Optional.empty();
}
/**
* As {@link #nudgeMayDrainFor(String)}, for whether a nudge may tell its recipient to
* {@code fleet_send(turnId=…)}.
*/
public Optional<Boolean> nudgeMayAnswerFor(String target) {
Delegation delegation = target == null ? null : leadByTarget.get(target);
if (delegation != null) {
return Optional.of(delegation.mayAnswerNudge());
}
return currentPrimaryTerminal().isPresent() ? Optional.of(true) : Optional.empty();
}
/**
* The known primary terminal, or empty if not yet learned (and not pinned) — the raw value as
* it was recorded, with no attempt to resolve a named lead's current pane. Callers that need a
@@ -63,6 +63,12 @@ public final class LeadCoordLoop {
private volatile boolean running;
/** Equivalent to {@link #LeadCoordLoop(LeadChannel, AgentControl, Supplier, ScheduledExecutorService, long, boolean)} with the prompt-box gate on. */
public LeadCoordLoop(LeadChannel channel, AgentControl agents, Supplier<Map<String, String>> leads,
ScheduledExecutorService scheduler, long intervalMs) {
this(channel, agents, leads, scheduler, intervalMs, true);
}
/**
* @param channel this daemon's own lead mailbox
* @param agents herdr control, for the status gate and the pane injection
@@ -71,12 +77,13 @@ public final class LeadCoordLoop {
* the tab scan after startup becomes reachable without a restart
* @param scheduler the loop's own scheduler; the caller owns its shutdown
* @param intervalMs how long between ticks
* @param promptBoxGateEnabled passed straight to {@link PromptBox#PromptBox(AgentControl, boolean)}
*/
public LeadCoordLoop(LeadChannel channel, AgentControl agents, Supplier<Map<String, String>> leads,
ScheduledExecutorService scheduler, long intervalMs) {
ScheduledExecutorService scheduler, long intervalMs, boolean promptBoxGateEnabled) {
this.channel = channel;
this.agents = agents;
this.promptBox = new PromptBox(agents);
this.promptBox = new PromptBox(agents, promptBoxGateEnabled);
this.leads = leads;
this.scheduler = scheduler;
this.intervalMs = intervalMs;
@@ -137,9 +137,21 @@ public final class LeadHeartbeatLoop {
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
LeadContextSource contextSource, boolean contextHighNudge,
boolean requireOperatorConfirm) {
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock,
idleAfterNanos, backoffMs, quietNudgeCap, metrics, contextSource, contextHighNudge,
requireOperatorConfirm, true);
}
/** As above, plus the prompt-box gate's on/off switch — see {@link PromptBox#PromptBox(AgentControl, boolean)}. */
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
ScheduledExecutorService scheduler, LongSupplier clock,
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
LeadContextSource contextSource, boolean contextHighNudge,
boolean requireOperatorConfirm, boolean promptBoxGateEnabled) {
this.primaryRegistry = primaryRegistry;
this.agents = agents;
this.promptBox = new PromptBox(agents);
this.promptBox = new PromptBox(agents, promptBoxGateEnabled);
this.inbox = inbox;
this.roster = roster;
this.pushLoop = pushLoop;
@@ -134,6 +134,14 @@ public final class MessageService {
NOT_TURN_OWNER
}
/**
* The fact shared by the MCP and REST receipts for {@link Outcome#TIMED_OUT_WORKING},
* {@link Outcome#TIMED_OUT_QUEUED} and {@link Outcome#BUSY}: a blocking send is never tracked
* by a ticket, so there is no id to poll by, and resending it risks a duplicate delivery.
*/
public static final String NO_TICKET_NO_RESEND =
"this send created no ticket, and a resend can duplicate the delivery";
/**
* @param outcome how the send ended (or paused)
* @param text the worker's answer when {@link #completed()} (a structured {@code fleet_reply}
@@ -353,6 +361,16 @@ public final class MessageService {
* target ({@link #send} opening a fresh waiter) or a teardown ({@link #abandon}).
*/
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
/**
* The captured waiter a timed-out send left open for a late turn-completion scrape
* ({@link #strandLateResolution}), keyed by target so a late {@code fleet_reply} can find and
* claim the exact one its own turn belongs to ({@link #reply}). The waiter's own {@code
* complete()} call is the discriminator: whichever caller resolves it first — this reply, or
* the eventual completion scrape — decides what the other one sees, and only the winner's kind
* ever reaches the inbox. An entry removes itself once its waiter resolves.
*/
private final ConcurrentHashMap<String, CompletableFuture<Rendezvous.Resolution>> lateWaiters =
new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
/**
* Minted once per {@code MessageService} instance and folded into every ticket id (see
@@ -612,6 +630,14 @@ public final class MessageService {
+ "answers, queuing to the inbox instead of guessing", session, candidates.size(),
tickets);
}
// Claim the matching timed-out send's captured waiter, if one is still open, before
// publishing this reply. The waiter is this exact turn's own CompletableFuture, so
// completing it here means a completion-fallback scrape that resolves the same waiter
// afterward finds it already answered and does not publish a second entry for this turn.
CompletableFuture<Rendezvous.Resolution> lateWaiter = lateWaiters.get(session);
if (lateWaiter != null) {
rendezvous.resolveLateReply(lateWaiter, content);
}
inbox.publish(session, UUID.randomUUID().toString(), content);
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
// worker whose replies keep missing their waiter, not only the queue depth this leaves behind.
@@ -808,6 +834,7 @@ public final class MessageService {
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
strandedReplies.remove(target);
queuedDeliveries.remove(target);
lateWaiters.remove(target);
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
boolean asyncFailed = false;
@@ -1050,6 +1077,7 @@ public final class MessageService {
queuedDeliveries.put(target, Boolean.TRUE);
outcome = Outcome.TIMED_OUT_QUEUED;
}
strandLateResolution(target, reply);
return recorded(new Reply(outcome, null));
} catch (ExecutionException e) {
Throwable cause = e.getCause();
@@ -1067,6 +1095,42 @@ public final class MessageService {
}
}
/**
* Once a blocking send has given up on {@code waiter} and already reported a TIMED_OUT_*
* outcome, the CB-106 completion fallback can still resolve it later — it completes the exact
* captured future ({@link Rendezvous#resolveCompletion}), not a session lookup, so
* {@link Rendezvous#close} does not stop it. With nobody left awaiting {@code waiter}, route
* that resolution to {@code target}'s inbox instead, exactly where a late explicit
* {@code fleet_reply} already lands (see {@link #reply}), so {@code fleet_poll{target}} can
* recover it. A no-op if {@code waiter} never resolves, or already has by the time this runs.
*
* <p>Only {@link Rendezvous.Kind#COMPLETION} is routed here. An explicit {@code fleet_reply}
* ({@link Rendezvous.Kind#REPLY}) still reaches {@code waiter} through {@link #reply}'s own
* session lookup, not through this captured reference, so leaving it out here avoids a second,
* racing publish of the same reply.
*/
private void strandLateResolution(String target, CompletableFuture<Rendezvous.Resolution> waiter) {
lateWaiters.put(target, waiter);
waiter.whenComplete((resolution, error) -> {
lateWaiters.remove(target, waiter);
if (resolution != null && resolution.kind() == Rendezvous.Kind.COMPLETION) {
try {
// inbox.publish reaches a broker and can throw. An exception thrown inside a
// whenComplete callback is swallowed into the discarded dependent stage, not
// rethrown to any caller — log it here, or a broker failure is invisible and the
// late reply is simply lost.
inbox.publish(target, UUID.randomUUID().toString(), resolution.text());
strandedReplies.put(target, Boolean.TRUE);
if (pushLoop != null) {
pushLoop.onReplyQueued(target);
}
} catch (RuntimeException e) {
log.warn("strandLateResolution: publishing a late turn-completion for {} failed", target, e);
}
}
});
}
/**
* A worker's mid-turn question (CB-205 reverse rendezvous): surface {@code question} to the
* primary by resolving its open blocking {@code fleet_send}, then block this (worker) call until
@@ -1167,6 +1231,14 @@ public final class MessageService {
}
}
/**
* The worker session an open {@code fleet_ask} turn belongs to, or {@code null} if
* {@code turnId} is unknown or has already lapsed.
*/
public String sessionForTurn(String turnId) {
return rendezvous.askSession(turnId);
}
/**
* The primary's answer to a worker's {@code fleet_ask} (CB-205): resolve the worker's blocked
* question identified by {@code turnId}, then — like a fresh {@link #send} — block for the worker's
@@ -1492,6 +1564,16 @@ public final class MessageService {
|| Objects.equals(callerOwner, task.creatorOwner);
}
/**
* Whether {@code callerOwner} owns {@code ticket}, with no side effect — unlike
* {@link #poll(String, String)}, this never fires the ticket's collection hook. A ticket this
* daemon has never heard of is owned by nobody.
*/
public boolean ownsTicket(String ticket, String callerOwner) {
Task task = ticket == null ? null : tasks.get(ticket);
return task != null && ownsTicket(task, callerOwner);
}
/**
* Test seam only — carries no production behaviour, and nothing in this class calls it;
* {@link #pruneTerminalTickets} still reads {@link Task#completedNanos} directly.
@@ -304,6 +304,20 @@ public final class Rendezvous {
return waiter != null && waiter.complete(new Resolution(Kind.BACKEND_EXHAUSTED, reason));
}
/**
* Resolve a specific captured {@code waiter} as a reply — for a {@code fleet_reply} that
* arrives after the send which opened {@code waiter} already gave up on it, so the ordinary
* session-keyed {@link #resolve} finds no live waiter to match. Like
* {@link #resolveCompletion(CompletableFuture, String)} it targets the exact captured waiter
* (CB-116): a no-op if that waiter was already resolved — first resolution wins, and whichever
* one wins is what any later resolution attempt on the same waiter will see.
*
* @return true if this call resolved the waiter, false if it was null or already resolved
*/
public boolean resolveLateReply(CompletableFuture<Resolution> waiter, String content) {
return waiter != null && waiter.complete(new Resolution(Kind.REPLY, content));
}
private boolean complete(String session, Resolution resolution) {
ForwardWaiter waiter = waiters.get(session);
return waiter != null && waiter.future().complete(resolution);
@@ -130,9 +130,16 @@ public final class ReplyPushLoop {
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
ScheduledExecutorService scheduler,
int maxReminders, long backoffMs, Metrics metrics) {
this(primaryRegistry, agents, inbox, scheduler, maxReminders, backoffMs, metrics, true);
}
/** As above, with the prompt-box gate's on/off switch — see {@link PromptBox#PromptBox(AgentControl, boolean)}. */
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
ScheduledExecutorService scheduler,
int maxReminders, long backoffMs, Metrics metrics, boolean promptBoxGateEnabled) {
this.primaryRegistry = primaryRegistry;
this.agents = agents;
this.promptBox = new PromptBox(agents);
this.promptBox = new PromptBox(agents, promptBoxGateEnabled);
this.inbox = inbox;
this.scheduler = scheduler;
this.maxReminders = maxReminders;
@@ -178,13 +185,20 @@ public final class ReplyPushLoop {
pendingReplies.remove(target, owning);
continue;
}
// Never nudge a caller to run a DRAIN its own role could never perform -- the durable
// inbox stays the backstop for it, same as when no nudge target is known at all.
if (!owning.mayDrainNudge()) continue;
result.add(target);
}
return result;
}
/** A pending reply target: which lead to nudge, and how many nudges have named it so far. */
private record ReplyEntry(String lead, int nudgeCount) {
/**
* A pending reply target: which nudge target to notify, how many nudges have named it so
* far, and whether that nudge target's own role may actually run the
* {@code fleet_poll(target=…)} a reply nudge would tell it to run.
*/
private record ReplyEntry(String lead, int nudgeCount, boolean mayDrainNudge) {
}
/**
@@ -207,11 +221,12 @@ public final class ReplyPushLoop {
/**
* An open question awaiting the lead's answer: which ticket it belongs to, which worker asked,
* which lead to nudge, the question text, and how many nudges have named it so far (CB-598 —
* tracked per question, not per lead per source).
* which lead to nudge, the question text, how many nudges have named it so far (CB-598 —
* tracked per question, not per lead per source), and whether that nudge target's own role may
* actually run the {@code fleet_send(turnId=…)} a question nudge would tell it to run.
*/
private record PendingQuestion(String turnId, String ticket, String target, String lead,
String question, int nudgeCount) {
String question, int nudgeCount, boolean mayAnswerNudge) {
}
private record IncidentLead(String incidentId, String lead) {
@@ -229,7 +244,11 @@ public final class ReplyPushLoop {
/** Questions still open for {@code lead}, snapshotted fresh for one tick. */
private List<PendingQuestion> pendingQuestionsFor(String lead) {
return pendingQuestions.values().stream().filter(q -> lead.equals(q.lead())).toList();
// Never nudge a caller to run an ANSWER its own role could never perform -- the 55-second
// ask window lapses on its own, same as when no nudge target is known at all.
return pendingQuestions.values().stream()
.filter(q -> lead.equals(q.lead()) && q.mayAnswerNudge())
.toList();
}
/** Question turnIds still open for {@code lead} — a plain snapshot for race comparison. */
@@ -385,7 +404,7 @@ public final class ReplyPushLoop {
boolean hasIncidentWork = !pendingIncidentKeysFor(lead).isEmpty();
boolean hasUnmappedTargetWork = !pendingUnmappedTargetKeysFor(lead).isEmpty();
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork && !hasIncidentWork && !hasUnmappedTargetWork) {
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
log.debug("push: nothing pending for nudge target {}, stopping reminder", lead);
return Action.STOP;
}
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
@@ -394,7 +413,7 @@ public final class ReplyPushLoop {
boolean incidentEligible = hasIncidentWork && incidentReminderCount < maxReminders;
boolean unmappedTargetEligible = hasUnmappedTargetWork && unmappedTargetReminderCount < maxReminders;
if (!replyEligible && !ticketEligible && !questionEligible && !incidentEligible && !unmappedTargetEligible) {
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
log.debug("push: reminder cap ({}) reached for nudge target {} on every source with pending work, stopping",
maxReminders, lead);
countNudge("exhausted");
return Action.STOP;
@@ -403,15 +422,15 @@ public final class ReplyPushLoop {
try {
status = agents.status(lead);
} catch (RuntimeException e) {
log.debug("push: status check failed for lead {}, will retry", lead, e);
log.debug("push: status check failed for nudge target {}, will retry", lead, e);
return Action.WAIT_BUSY;
}
if (!status.injectable()) {
log.debug("push: lead {} is {} (not injectable), waiting", lead, status);
log.debug("push: nudge target {} is {} (not injectable), waiting", lead, status);
return Action.WAIT_BUSY;
}
if (!promptBox.clearToSubmit(lead)) {
log.debug("push: lead {} has unsubmitted text in its prompt box, waiting", lead);
log.debug("push: nudge target {} has unsubmitted text in its prompt box, waiting", lead);
return Action.WAIT_BUSY;
}
return Action.INJECT;
@@ -466,14 +485,14 @@ public final class ReplyPushLoop {
if (lead.isEmpty() || isLive(lead.get())) {
return lead;
}
log.debug("push: lead {} delegated to for {} is no longer live, forgetting the stale binding "
log.debug("push: nudge target {} delegated to for {} is no longer live, forgetting the stale binding "
+ "and falling back", lead.get(), target);
primaryRegistry.forgetDelegation(target);
Optional<String> fallback = primaryRegistry.nudgeTargetFor(target);
if (fallback.isEmpty() || isLive(fallback.get())) {
return fallback;
}
log.debug("push: fallback lead {} for {} is also not live, skipping this tick", fallback.get(), target);
log.debug("push: fallback nudge target {} for {} is also not live, skipping this tick", fallback.get(), target);
return Optional.empty();
}
@@ -496,9 +515,9 @@ public final class ReplyPushLoop {
} catch (RuntimeException e) {
boolean gone = e instanceof HerdrException he && "agent_not_found".equals(he.code());
if (gone) {
log.debug("push: lead {} no longer exists ({})", lead, e.toString());
log.debug("push: nudge target {} no longer exists ({})", lead, e.toString());
} else {
log.debug("push: liveness check for lead {} was inconclusive ({}); treating as live "
log.debug("push: liveness check for nudge target {} was inconclusive ({}); treating as live "
+ "rather than risk destroying a live binding", lead, e.toString());
}
return !gone;
@@ -517,11 +536,12 @@ public final class ReplyPushLoop {
public void onReplyQueued(String target) {
var lead = resolveLiveLead(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
log.debug("push: no nudge target is known to be waiting on {}, skipping reminder", target);
return;
}
boolean mayDrainNudge = primaryRegistry.nudgeMayDrainFor(target).orElse(false);
pendingReplies.compute(target, (t, existing) ->
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount()));
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount(), mayDrainNudge));
startOrCoalesce(lead.get());
}
@@ -545,7 +565,7 @@ public final class ReplyPushLoop {
public void onTicketTerminal(String ticket, String target, boolean failed) {
var lead = resolveLiveLead(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
log.debug("push: no nudge target is known to be waiting on ticket {} (target {}), skipping nudge",
ticket, target);
return;
}
@@ -580,11 +600,13 @@ public final class ReplyPushLoop {
public void onQuestionOpened(String ticket, String target, String turnId, String question) {
var lead = resolveLiveLead(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on {}'s question (turnId {}), skipping nudge",
log.debug("push: no nudge target is known to be waiting on {}'s question (turnId {}), skipping nudge",
target, turnId);
return;
}
pendingQuestions.put(turnId, new PendingQuestion(turnId, ticket, target, lead.get(), question, 0));
boolean mayAnswerNudge = primaryRegistry.nudgeMayAnswerFor(target).orElse(false);
pendingQuestions.put(turnId,
new PendingQuestion(turnId, ticket, target, lead.get(), question, 0, mayAnswerNudge));
startOrCoalesce(lead.get());
}
@@ -608,7 +630,7 @@ public final class ReplyPushLoop {
for (String target : targets) {
var lead = resolveLiveLead(target);
if (lead.isEmpty()) {
log.warn("push: backend incident {} has no known lead for target {}", incidentId, target);
log.warn("push: backend incident {} has no known nudge target for worker {}", incidentId, target);
continue;
}
targetsByLead.computeIfAbsent(lead.get(), _ -> new ArrayList<>()).add(target);
@@ -644,10 +666,10 @@ public final class ReplyPushLoop {
/** Start a reminder schedule for {@code lead}, or join the one already running. */
private void startOrCoalesce(String lead) {
if (activeLeads.putIfAbsent(lead, Boolean.TRUE) != null) {
log.debug("push: reminder loop already active for lead {}, work coalesced in", lead);
log.debug("push: reminder loop already active for nudge target {}, work coalesced in", lead);
return;
}
log.debug("push: starting reminder loop for lead {}", lead);
log.debug("push: starting reminder loop for nudge target {}", lead);
scheduleNext(lead);
}
@@ -752,11 +774,11 @@ public final class ReplyPushLoop {
|| pendingIncidentKeysFor(lead).stream().anyMatch(i -> !incidentsBefore.contains(i))
|| pendingUnmappedTargetKeysFor(lead).stream().anyMatch(i -> !unmappedTargetsBefore.contains(i));
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
log.debug("push: new work for nudge target {} raced the reminder loop's stop — restarting", lead);
scheduleNext(lead);
return;
}
log.debug("push: reminder loop ended for lead {}", lead);
log.debug("push: reminder loop ended for nudge target {}", lead);
}
/** Send one combined nudge covering everything currently pending for {@code lead}. */
@@ -773,13 +795,13 @@ public final class ReplyPushLoop {
List<PendingUnmappedTarget> unmappedTargets = pendingUnmappedTargetsFor(lead);
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty() && incidents.isEmpty()
&& unmappedTargets.isEmpty()) {
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
log.debug("push: pending work for nudge target {} drained before the nudge could be sent", lead);
return;
}
String nudge = formatNudge(replyTargets, tickets, questions, incidents, unmappedTargets);
try {
agents.send(lead, nudge);
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; "
log.debug("push: nudge sent to nudge target {} (reply {}/{}, ticket {}/{}, question {}/{}; "
+ "{} reply target(s), {} ticket(s), {} question(s))",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
questionReminderCount + 1, maxReminders,
@@ -796,7 +818,7 @@ public final class ReplyPushLoop {
}
}
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
log.warn("push: failed to nudge target {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
questionReminderCount + 1, maxReminders, e.toString());
}
@@ -814,7 +836,8 @@ public final class ReplyPushLoop {
List<PendingQuestion> questions, List<PendingIncident> incidents,
List<PendingUnmappedTarget> unmappedTargets) {
for (String target : replyTargets) {
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
pendingReplies.computeIfPresent(target,
(t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1, e.mayDrainNudge()));
}
for (PendingTicket ticket : tickets) {
pendingTickets.computeIfPresent(ticket.ticket(),
@@ -823,7 +846,7 @@ public final class ReplyPushLoop {
for (PendingQuestion question : questions) {
pendingQuestions.computeIfPresent(question.turnId(), (id, e) ->
new PendingQuestion(e.turnId(), e.ticket(), e.target(), e.lead(), e.question(),
e.nudgeCount() + 1));
e.nudgeCount() + 1, e.mayAnswerNudge()));
}
for (PendingIncident incident : incidents) {
pendingIncidents.computeIfPresent(incident.key(), (id, e) -> new PendingIncident(e.key(),
@@ -285,13 +285,26 @@ public final class FleetApp {
/**
* As {@link #permitsFor(Principal, Authz.Action, String, Predicate)}, also threading the
* classifier an observer's {@code SEND} is checked against; pass {@link #auth}'s own
* {@code sendableObserverTarget()} to exercise the real production gate, as {@link #allow}
* does.
* {@code observerSendTarget()} to exercise the real production gate, as {@link #allow} does.
*/
static boolean permitsFor(Principal caller, Authz.Action action, String target,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> knownObserverTarget) {
return Authz.permits(caller, action, target, knownLeadOrCollaborator, knownObserverTarget);
Predicate<String> observerSendTarget) {
return Authz.permits(caller, action, target, knownLeadOrCollaborator, observerSendTarget);
}
/**
* As {@link #permitsFor(Principal, Authz.Action, String, Predicate, Predicate)}, also
* threading the classifier an observer's or a collaborator's {@code TASK_READ} is checked
* against; pass a real {@code MessageService#ownsTicket(String, String)}-backed predicate to
* exercise the real production gate for a ticket poll, as {@link #allow} does.
*/
static boolean permitsFor(Principal caller, Authz.Action action, String target,
Predicate<String> knownLeadOrCollaborator,
Predicate<String> observerSendTarget,
Predicate<String> ticketOwnedByCaller) {
return Authz.permits(caller, action, target, knownLeadOrCollaborator, observerSendTarget,
ticketOwnedByCaller);
}
/**
@@ -303,11 +316,22 @@ public final class FleetApp {
* yours" (a worker reaching for another worker's session, or for orchestration).
*/
private boolean allow(Context ctx, Authz.Action action, String target) {
return allow(ctx, action, target, Authz.NO_OWNED_TICKET);
}
/**
* As {@link #allow(Context, Authz.Action, String)}, with a real classifier for an observer's
* or a collaborator's {@code TASK_READ} — the only call site that can supply one is a ticket
* poll, which knows the ticket id and the caller's owner key.
*/
private boolean allow(Context ctx, Authz.Action action, String target,
Predicate<String> ticketOwnedByCaller) {
if (auth == null) {
return true; // legacy: authorization not enforced
}
Principal caller = ctx.attribute(CALLER);
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.sendableObserverTarget())) {
if (permitsFor(caller, action, target, auth.knownLeadOrCollaborator(), auth.observerSendTarget(),
ticketOwnedByCaller)) {
if (action != Authz.Action.READ && action != Authz.Action.METRICS
&& action != Authz.Action.TASK_READ) {
AuditLog.allowed(caller, action, target); // reads would drown the trail
@@ -718,9 +742,15 @@ public final class FleetApp {
content = FleetMcp.attributeIfObserver(caller, content);
timeout = Math.clamp(timeout, 1, MAX_MESSAGE_TIMEOUT_MS);
// Whether THIS caller's own Authz.Action.DRAIN is granted — a TIMED_OUT_*/BUSY receipt
// must not name GET /sessions/{id}/replies when that same caller would be refused it.
// Mirrors allow()'s own legacy bypass: with no CallerResolver configured, that route is
// not gated at all, so every caller may reach it.
boolean mayDrainPoll = auth == null || Authz.permits(caller, Authz.Action.DRAIN, id);
// Answering a worker's fleet_ask (CB-205): always blocks, and derives the worker from turnId.
if (turnId != null && !turnId.isBlank()) {
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerOwner), timeout);
writeReply(ctx, id, messages.answer(turnId, content, timeout, callerOwner), timeout, mayDrainPoll);
return;
}
@@ -732,7 +762,7 @@ public final class FleetApp {
}
try {
writeReply(ctx, id, messages.send(id, content, timeout, callerOwner), timeout);
writeReply(ctx, id, messages.send(id, content, timeout, callerOwner), timeout, mayDrainPoll);
} catch (HerdrException e) {
herdrError(ctx, e);
}
@@ -742,8 +772,13 @@ public final class FleetApp {
* Render a {@link MessageService.Reply} onto the response — shared by a normal send and a
* fleet_ask answer. A structured/scraped completion is 200; a worker's mid-turn question a 202
* (with its {@code turnId}); a stale answer a 409; every other non-terminal outcome a typed 202.
*
* @param mayDrainPoll whether this caller's {@code Authz.Action.DRAIN} is granted — a
* TIMED_OUT or BUSY {@code detail} names {@code GET /sessions/{id}/replies}
* only when this is {@code true}, since that route shares the same grant.
*/
private void writeReply(Context ctx, String id, MessageService.Reply reply, long timeout) {
private void writeReply(Context ctx, String id, MessageService.Reply reply, long timeout,
boolean mayDrainPoll) {
switch (reply.outcome()) {
case QUESTION -> ctx.status(202).json(Map.of(
"sessionId", id, "status", "question",
@@ -792,7 +827,10 @@ public final class FleetApp {
|| reply.outcome() == MessageService.Outcome.BACKEND_EXHAUSTED)
&& reply.text() != null
? reply.text()
: "no reply within " + timeout + "ms; poll status or retry"));
: "no reply within " + timeout + "ms; " + MessageService.NO_TICKET_NO_RESEND + "; "
+ (mayDrainPoll
? "GET /sessions/" + id + "/replies drains its inbox for a late reply"
: "a late reply cannot be recovered on this channel — ask the lead")));
}
}
@@ -929,11 +967,15 @@ public final class FleetApp {
/** Poll an async (wait:false) delegation by ticket. 404 for an unknown/expired ticket. */
private void taskStatus(Context ctx) {
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), null)) {
String ticket = ctx.pathParam("ticket");
Principal caller = ctx.attribute(CALLER);
String callerOwner = caller == null ? null : caller.ownerKey();
// The gate must see the ticket id, not a literal null -- an observer's or a
// collaborator's grant is confined to a ticket its own caller created.
if (!allow(ctx, routeAction("GET /tasks/{ticket}"), ticket, t -> messages.ownsTicket(t, callerOwner))) {
return;
}
Principal caller = ctx.attribute(CALLER);
MessageService.TaskView v = messages.poll(ctx.pathParam("ticket"), caller == null ? null : caller.ownerKey());
MessageService.TaskView v = messages.poll(ticket, callerOwner);
if (v == null) {
ctx.status(404).json(Map.of("error", "unknown_ticket", "detail", "no such task (or it has expired)"));
return;
@@ -259,18 +259,39 @@ class AuthzTest {
}
/**
* Every action denied to a collaborator, asserted denied even when the classifier would
* accept any target — proving none of these is actually gated on the classifier at all.
* Every action denied to a collaborator, asserted denied even when the
* {@code knownLeadOrCollaborator} classifier would accept any target — proving none of these
* is actually gated on that classifier at all. {@code TASK_READ} is excluded here and given
* its own matrix below, since — unlike every action in this loop — its grant is conditional
* on the separate ticket-ownership classifier, not on this one.
*/
@Test
void aCollaboratorIsDeniedLifecycleCoordinationAndTicketPolling() {
void aCollaboratorIsDeniedLifecycleAndCoordination() {
for (Authz.Action a : new Authz.Action[]{SPAWN, STOP, DRAIN, HANDOVER, ANSWER, COORD_SEND,
COORD_READ, TASK_READ}) {
COORD_READ}) {
assertFalse(Authz.permits(COLLABORATOR, a, "term_lead", target -> true),
"a collaborator must not " + a + " even when the classifier accepts every target");
}
}
/**
* {@code TASK_READ} is conditional for a collaborator too (fleetd #804), on the same
* ticket-ownership classifier an observer's grant uses: it may poll a ticket its own
* {@code fleet_send(wait:false)} created, and nothing else.
*/
@Test
void aCollaboratorMayTaskReadOnlyATicketTheClassifierSaysItOwns() {
assertTrue(Authz.permits(COLLABORATOR, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
Authz.NO_OBSERVER_SEND_TARGET, target -> true),
"the classifier accepting the ticket as the caller's own must grant TASK_READ");
assertFalse(Authz.permits(COLLABORATOR, TASK_READ, "someone_elses_ticket",
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, Authz.NO_OBSERVER_SEND_TARGET, target -> false),
"the classifier refusing the ticket must deny TASK_READ");
assertFalse(Authz.permits(COLLABORATOR, TASK_READ, "own-ticket"),
"the three-argument convenience form fails closed, so TASK_READ is refused without "
+ "a real classifier");
}
@Test
void aCollaboratorIsNotCountedAsPrimaryWorkerOrArchitect() {
assertFalse(COLLABORATOR.isPrimary());
@@ -310,17 +331,17 @@ class AuthzTest {
}
/**
* Every action beyond READ/METRICS/REPLY/ASK/INBOX/SEND, asserted denied for an observer —
* including {@code TASK_READ}, which is the entire point of this role: an unconfigured pane
* must not be able to poll a ticket or read another session's status. The exempt set is the
* three only-as-itself actions plus the two open reads. {@code SEND} is excluded here and
* given its own matrix below, since — unlike every action in this loop — its grant is
* conditional on the target, not fixed.
* Every action beyond READ/METRICS/REPLY/ASK/INBOX/SEND, asserted denied for an observer.
* The exempt set is the three only-as-itself actions plus the two open reads. {@code SEND}
* and {@code TASK_READ} are excluded here and given their own matrices below, since — unlike
* every action in this loop — each one's grant is conditional on more than the caller's role
* alone.
*/
@Test
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskInboxAndSend() {
void anObserverIsDeniedEverythingBeyondReadMetricsReplyAskInboxSendAndTaskRead() {
for (Authz.Action a : Authz.Action.values()) {
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == INBOX || a == SEND) {
if (a == READ || a == METRICS || a == REPLY || a == ASK || a == INBOX || a == SEND
|| a == TASK_READ) {
continue;
}
assertFalse(Authz.permits(OBSERVER, a, "term_observer", target -> true),
@@ -328,6 +349,49 @@ class AuthzTest {
}
}
/**
* {@code TASK_READ} is conditional for an observer, not fixed: it may poll a ticket its own
* {@code fleet_send(wait:false)} created, confined by {@code ticketOwnedByCaller}, and nothing
* else — a session's status is not ticket-scoped, so no call site can ever supply a predicate
* that grants it, and the 3-/4-/5-argument convenience forms stay fail-closed for it too.
*/
@Test
void anObserverMayTaskReadOnlyATicketTheClassifierSaysItOwns() {
assertTrue(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
Authz.NO_OBSERVER_SEND_TARGET, target -> true),
"the classifier accepting the ticket as the caller's own must grant TASK_READ");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "someone_elses_ticket",
Authz.NO_KNOWN_LEAD_OR_COLLABORATOR, Authz.NO_OBSERVER_SEND_TARGET, target -> false),
"the classifier refusing the ticket must deny TASK_READ");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket"),
"the three-argument convenience form fails closed, so TASK_READ is refused without "
+ "a real classifier");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR),
"the four-argument form still fails closed for TASK_READ");
assertFalse(Authz.permits(OBSERVER, TASK_READ, "own-ticket", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
Authz.NO_OBSERVER_SEND_TARGET),
"the five-argument form still fails closed for TASK_READ, since it supplies no "
+ "ticketOwnedByCaller classifier either");
}
/**
* Control for the test above: every other action's result for an observer does not move when
* the ticket-ownership classifier does. Only {@code TASK_READ} is wired to it.
*/
@Test
void theTicketOwnershipClassifierMovesOnlyTaskReadForAnObserver() {
for (Authz.Action a : Authz.Action.values()) {
if (a == TASK_READ) {
continue;
}
assertEquals(
Authz.permits(OBSERVER, a, "term_observer"),
Authz.permits(OBSERVER, a, "term_observer", Authz.NO_KNOWN_LEAD_OR_COLLABORATOR,
Authz.NO_OBSERVER_SEND_TARGET, target -> true),
a + " must not depend on the ticket-ownership classifier at all");
}
}
@Test
void anObserverIsNotCountedAsAnyOtherRole() {
assertFalse(OBSERVER.isPrimary());
@@ -912,42 +912,65 @@ class CallerResolverTest {
assertFalse(r.knownLeadOrCollaborator().test("term_a"));
}
// ── fleetd #743: sendableObserverTarget() reads the same maps and functions resolve() does ────
// ── observerSendTarget() reads the same maps and functions resolve() does, in the same order ──
/**
* A terminal this resolver recognises as none of the privileged roles is exactly the one
* {@code resolve} would itself hand back {@link Role#OBSERVER} for.
*/
@Test
void sendableObserverTargetIsTrueForATerminalKnownAsNoOtherRole() {
void observerSendTargetIsTrueForATerminalKnownAsNoOtherRole() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_lead", "opus-5.0"), new MemberRegistry(null),
t -> "term_worker".equals(t) ? MemberRole.DEV : null,
() -> Map.of("term_collab", "ops"));
assertTrue(r.sendableObserverTarget().test("term_other"));
assertTrue(r.observerSendTarget().test("term_other"));
}
/**
* A configured lead's own terminal is reachable: an observer may open a conversation with a
* lead, and this is the classifier that grant is checked against.
*/
@Test
void sendableObserverTargetIsFalseForALeadTerminal() {
void observerSendTargetIsTrueForALeadTerminal() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_lead", "opus-5.0"), new MemberRegistry(null), t -> null, Map::of);
assertFalse(r.sendableObserverTarget().test("term_lead"),
"a lead's own terminal must never be a sendable observer target");
assertTrue(r.observerSendTarget().test("term_lead"),
"a lead's own terminal must be reachable from an observer pane");
}
/**
* A pane named as a lead AND bound to an architect slot resolves as the lead, because
* {@code resolve} reads the lead map first — so the classifier must accept it, or a pane's
* resolved role and its reachability would disagree.
*/
@Test
void observerSendTargetIsTrueForALeadTerminalThatIsAlsoABoundArchitectSlot() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_a", "opus-5.0"),
boundMembers("architect:lead-designer", MemberRole.ARCHITECT), t -> null, Map::of);
assertEquals(Role.PRIMARY, r.resolve("127.0.0.1", 42, null).role(),
"premise: the lead map is read before the architect registry");
assertTrue(r.observerSendTarget().test("term_a"));
}
@Test
void sendableObserverTargetIsFalseForACollaboratorTerminal() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
void observerSendTargetIsFalseForACollaboratorTerminal() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_lead", "opus-5.0"),
new MemberRegistry(null), t -> null, () -> Map.of("term_collab", "ops"));
assertFalse(r.sendableObserverTarget().test("term_collab"),
"a collaborator's own terminal must never be a sendable observer target");
assertFalse(r.observerSendTarget().test("term_collab"),
"a collaborator's own terminal must never be reachable from an observer pane");
// CONTROL: the same wiring, a target recognised as no configured role at all.
assertTrue(r.observerSendTarget().test("term_other"));
}
@Test
void sendableObserverTargetIsFalseForALiveSpawnedMembersTerminal() {
void observerSendTargetIsFalseForALiveSpawnedMembersTerminal() {
// Covers both a worker and an architect: spawnedMemberRole.apply(target) is non-null for
// either, and resolve() never falls through to OBSERVER once it is.
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
@@ -958,26 +981,47 @@ class CallerResolverTest {
default -> null;
}, Map::of);
assertFalse(r.sendableObserverTarget().test("term_worker"));
assertFalse(r.sendableObserverTarget().test("term_architect"));
assertFalse(r.observerSendTarget().test("term_worker"));
assertFalse(r.observerSendTarget().test("term_architect"));
// CONTROL: the same wiring, a target the member lookup above answers null for.
assertTrue(r.observerSendTarget().test("term_other"));
}
/**
* A lead map entry does not rescue a terminal a live spawned member occupies: the member
* lookup runs first, exactly as in {@code resolve}.
*/
@Test
void observerSendTargetIsFalseForASpawnedMemberOnATerminalTheLeadMapAlsoNames() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null,
() -> Map.of("term_worker", "opus-5.0"), new MemberRegistry(null),
t -> "term_worker".equals(t) ? MemberRole.DEV : null, Map::of);
assertFalse(r.observerSendTarget().test("term_worker"));
// CONTROL: the same wiring, the same lead map, a terminal no member occupies.
assertTrue(r.observerSendTarget().test("term_other"));
}
@Test
void sendableObserverTargetIsFalseForABoundArchitectSlotWithNoLiveMember() {
void observerSendTargetIsFalseForABoundArchitectSlotWithNoLiveMember() {
// The edge case resolve() itself carries: a terminal bound to a configured architect slot
// but with no live spawned-member session yet.
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
boundMembers("architect:lead-designer", MemberRole.ARCHITECT), t -> null, Map::of);
assertFalse(r.sendableObserverTarget().test("term_a"));
assertFalse(r.observerSendTarget().test("term_a"));
// CONTROL: the same wiring, a terminal the bind above never touched.
assertTrue(r.observerSendTarget().test("term_other"));
}
@Test
void sendableObserverTargetIsFalseForANullTarget() {
void observerSendTargetIsFalseForANullTarget() {
CallerResolver r = CallerResolver.withLeadsAndMembers(workerIdentity(), false, null, Map::of,
new MemberRegistry(null), t -> null, Map::of);
assertFalse(r.sendableObserverTarget().test(null));
assertFalse(r.observerSendTarget().test(null));
// CONTROL: the same wiring, a non-null target.
assertTrue(r.observerSendTarget().test("term_other"));
}
}
@@ -114,6 +114,7 @@ class ConfigRefTopLevelReportingCoverageTest {
// fleetd #480: leadRollover joins placement/memberCredentials/memberLoginShell as
// hot-excluded — never compared by any changed*Keys method, so it stays null like them.
v.put("leadRollover", null);
v.put("promptBoxGateEnabled", true);
assertNamesMatchComponents(v);
return v;
}
@@ -159,6 +160,7 @@ class ConfigRefTopLevelReportingCoverageTest {
v.put("models", new FleetConfig.Models(
List.of(new FleetConfig.Models.ModelEntry("model-b"))));
v.put("leadRollover", null);
v.put("promptBoxGateEnabled", false);
assertNamesMatchComponents(v);
return v;
}
@@ -104,6 +104,7 @@ class FleetConfigWithDefaultsPreservesEveryComponentTest {
// here proves it, rather than leaving it null and proving nothing.
v.put("leadRollover", new FleetConfig.LeadRollover(
"/handover/guard.md", true, 3600, 20, 45, "read the handover file"));
v.put("promptBoxGateEnabled", true);
assertNamesMatchComponents(v);
return v;
}
@@ -177,4 +177,24 @@ class PromptBoxTest {
herdr.detectionText(EMPTY);
assertTrue(box.clearToSubmit("term_a"));
}
// --- the gate's on/off switch ---------------------------------------------
@Test
void aDisabledGateClearsADraftedBoxWithoutReadingThePane() {
FakeHerdr herdr = new FakeHerdr().detectionText(DRAFTED);
assertTrue(new PromptBox(new AgentControl(herdr), false).clearToSubmit("term_a"),
"a disabled gate must clear every target, draft or not");
assertFalse(herdr.called("agent.read"), "a disabled gate must never read the pane");
}
@Test
void anEnabledGateBehavesExactlyLikeTheSingleArgConstructor() {
FakeHerdr herdr = new FakeHerdr().detectionText(DRAFTED);
assertFalse(new PromptBox(new AgentControl(herdr), true).clearToSubmit("term_a"),
"gateEnabled=true must hold a draft exactly like the default constructor");
assertTrue(herdr.called("agent.read"), "an enabled gate still reads the pane");
}
}
@@ -2,10 +2,12 @@ package dev.ltms.fleet.lead;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.spi.ILoggingEvent;
import com.fasterxml.jackson.databind.JsonNode;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.testing.CapturedLog;
import org.junit.jupiter.api.DisplayName;
@@ -20,6 +22,7 @@ import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Function;
import java.util.function.LongSupplier;
@@ -1323,6 +1326,66 @@ class LeadRolloverTest {
+ "so an operator reading status() has something to act on: " + status.detail());
}
@Test
@DisplayName("[BOOTSTRAP 1] one agent_not_ready bootstrap refusal is retried and sends the "
+ "handover instruction exactly once")
void transientAgentNotReadyRetriesBootstrapAndRolls() throws IOException {
FakeHerdr fake = herdrReadyForAFullRoll();
AtomicInteger refusedSends = new AtomicInteger();
HerdrClient transientRefusal = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) {
if ("agent.prompt".equals(method) && refusedSends.getAndIncrement() == 0) {
throw new HerdrException("herdr error [agent_not_ready]: agent.prompt failed",
"agent_not_ready", null);
}
return fake.call(method, params);
}
@Override
public void close() {
}
};
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
AgentControl agents = new AgentControl(transientRefusal);
WorkspaceControl spaces = new WorkspaceControl(transientRefusal);
LeadRollover rollover = new LeadRollover(agents, spaces,
new LeadLauncher(agents, spaces, fleetConfigWithRelaunchableLead()),
() -> cfg(handover.toString()), _ -> null, _ -> LEAD_NAME,
() -> Map.of(NEW_TERMINAL, LEAD_NAME), fixedClock(clock), () -> { }, Runnable::run);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
assertTrue(rollover.confirm(LEAD, pending.token(), true).accepted());
assertEquals(2, refusedSends.get(), "one agent_not_ready refusal must be followed by one retry");
assertEquals(1, promptCallCount(fake), "only the successful retry reaches herdr delivery");
assertEquals(LeadRollover.RollState.ROLLED, rollover.status(pending.token()).state());
}
@Test
@DisplayName("[BOOTSTRAP 2] persistent agent_not_ready records BOOTSTRAP_NEVER_SENT and "
+ "releases the single-flight claim")
void persistentAgentNotReadyRecordsBootstrapNeverSentAndReleasesClaim() throws IOException {
FakeHerdr fake = herdrReadyForAFullRoll();
fake.agentSendFailsWith("agent_not_ready");
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRolloverForAFullRoll(fake, cfg(handover.toString()),
() -> clock.addAndGet(1_000), Map.of(NEW_TERMINAL, LEAD_NAME));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
assertTrue(rollover.confirm(LEAD, pending.token(), true).accepted());
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.BOOTSTRAP_NEVER_SENT, status.state());
assertNotEquals(LeadRollover.RollState.ROLLED, status.state());
assertNotEquals(LeadRollover.RollState.FAILED, status.state());
LeadRollover.PendingRollover retry = rollover.open(LEAD, "retry after bootstrap timeout");
assertTrue(rollover.confirm(LEAD, retry.token(), true).accepted(),
"the terminal outcome must release the single-flight claim");
}
// ---- fleetd #726 unit 3: confirm() single-flights one roll at a time per lead ---------------
@Test
@@ -56,6 +56,7 @@ class FleetMcpAuthzTest {
private final AgentControl agents = new AgentControl(herdr);
private Metrics metrics;
private FleetMcp mcp;
private MessageService messages;
@AfterEach
void close() {
@@ -82,7 +83,7 @@ class FleetMcpAuthzTest {
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> "tok");
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(),
messages = new MessageService(agents, new Injector(agents), new Rendezvous(),
new InMemoryReplyInbox());
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
metrics = FleetMetrics.create(sessions, new InMemoryReplyInbox());
@@ -270,18 +271,19 @@ class FleetMcpAuthzTest {
"a spawned member's own terminal must stay unreachable, even once the classifier is real");
}
// --- fleetd #743: the observer SEND matrix, over MCP's denyFor -------------------------------
// --- the observer SEND matrix, over MCP's denyFor -------------------------------------------
private static final Principal OBSERVER = Principal.observer("term_observer", 700);
/**
* Wires one real {@link CallerResolver} that recognises a lead, a collaborator, and a live
* spawned worker, leaving "term_other_observer" classified as none of them — so the same
* wiring both denies an observer's {@code SEND} to every privileged role and grants it to
* another unclassified pane, proving the refusals are the rule and not a missing fixture.
* wiring denies an observer's {@code SEND} to a collaborator and to a member while granting
* it to a lead and to another unclassified pane, proving the refusals are the rule and not a
* missing fixture.
*/
@Test
void anObserverMaySendOnlyToAnotherObserverNeverToALeadWorkerOrArchitect() {
void anObserverMaySendToALeadOrAnotherObserverButNeverToACollaboratorOrAMember() {
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
() -> Map.of("term_lead_known", "lead-x"), new MemberRegistry(null),
@@ -289,22 +291,23 @@ class FleetMcpAuthzTest {
() -> Map.of("term_collab_known", "ops2"));
FleetMcp m = mcp(true, callers);
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_lead_known"),
"an observer must never reach a lead's terminal");
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_lead_known"),
"an observer must reach a lead's terminal, so a peer session can open a "
+ "conversation with a lead");
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_collab_known"),
"an observer must never reach a collaborator's terminal");
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_a"),
"an observer must never reach a live spawned member's terminal");
// CONTROL: the same wiring, the same denyFor call, a target recognised as none of the
// three privileged roles above -- this is what proves the three refusals above are the
// rule working, not a classifier that refuses every target regardless of what it is.
// configured roles above -- this is what proves the two refusals above are the rule
// working, not a classifier that refuses every target regardless of what it is.
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_other_observer"),
"an observer must reach another pane that resolves as an observer itself");
}
/**
* As {@link #anObserverMaySendOnlyToAnotherObserverNeverToALeadWorkerOrArchitect}, for a
* As {@link #anObserverMaySendToALeadOrAnotherObserverButNeverToACollaboratorOrAMember}, for a
* terminal bound to a configured architect slot but hosting no live spawned-member session --
* the case {@link CallerResolver#resolve} itself treats separately from a live worker/architect.
*/
@@ -324,6 +327,93 @@ class FleetMcpAuthzTest {
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_other_observer"));
}
/**
* An observer's reach is widened for {@code SEND} alone. It still holds no {@code TASK_READ},
* so it cannot poll a ticket or read a lead's session status, and no {@code SPAWN}/
* {@code STOP}/{@code COORD_SEND}, so it cannot drive the fleet it can now message.
*/
@Test
void anObserverReachingALeadStillHoldsNoTicketReadAndNoLifecycle() {
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
() -> Map.of("term_lead_known", "lead-x"), new MemberRegistry(null), t -> null, Map::of);
FleetMcp m = mcp(true, callers);
assertNull(m.denyFor(OBSERVER, Authz.Action.SEND, "term_lead_known"),
"premise: this wiring grants the observer's send to that lead");
assertNotNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, "term_lead_known"));
assertNotNull(m.denyFor(OBSERVER, Authz.Action.SPAWN, null));
assertNotNull(m.denyFor(OBSERVER, Authz.Action.STOP, null));
assertNotNull(m.denyFor(OBSERVER, Authz.Action.COORD_SEND, null));
}
/**
* An observer may read a ticket its own {@code fleet_send(wait:false)} created, never one
* another caller created — exercised against the real {@link MessageService#ownsTicket}, not
* a stand-in classifier.
*/
@Test
void anObserverMayTaskReadATicketItCreatedButNotOneAnotherCallerCreated() {
FleetMcp m = mcp(true);
String ownTicket = messages.sendAsync("term_a", "do it", null, OBSERVER);
String othersTicket = messages.sendAsync("term_a", "do it too", null, WORKER_A);
assertNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, ownTicket,
t -> messages.ownsTicket(t, OBSERVER.ownerKey())),
"an observer must read a ticket its own send created");
assertNotNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, othersTicket,
t -> messages.ownsTicket(t, OBSERVER.ownerKey())),
"an observer must not read a ticket a different caller created");
}
/**
* As {@link #anObserverMayTaskReadATicketItCreatedButNotOneAnotherCallerCreated} (fleetd
* #804): a collaborator's {@code fleet_send(wait:false)} creates a ticket too, and the only
* caller that can ever read it back is the collaborator that created it.
*/
@Test
void aCollaboratorMayTaskReadATicketItCreatedButNotOneAnotherCallerCreated() {
FleetMcp m = mcp(true);
String ownTicket = messages.sendAsync("term_a", "do it", null, COLLABORATOR);
String othersTicket = messages.sendAsync("term_a", "do it too", null, WORKER_A);
assertNull(m.denyFor(COLLABORATOR, Authz.Action.TASK_READ, ownTicket,
t -> messages.ownsTicket(t, COLLABORATOR.ownerKey())),
"a collaborator must read a ticket its own send created");
assertNotNull(m.denyFor(COLLABORATOR, Authz.Action.TASK_READ, othersTicket,
t -> messages.ownsTicket(t, COLLABORATOR.ownerKey())),
"a collaborator must not read a ticket a different caller created");
}
/**
* fleetd #803: {@code fleet_poll{}} with no {@code ticket} argument resolves to
* {@code TASK_READ} with a {@code null} ticket. Before the fix, the real {@code
* ownsTicket(String, String)} classifier reached {@code tasks.get(null)} on a
* {@code ConcurrentHashMap} and threw a {@code NullPointerException} instead of refusing —
* the gate never got to write its {@code AuditLog} refusal, the caller saw an internal error.
* Asserts the refusal itself, not merely that nothing throws: a test that only catches the
* absence of a throw would also pass if this path quietly started granting the call.
*/
@Test
void anObserverPollingWithNoTicketIsRefusedRatherThanThrowing() {
FleetMcp m = mcp(true);
McpSchema.CallToolResult denied = m.denyFor(OBSERVER, Authz.Action.TASK_READ, null,
t -> messages.ownsTicket(t, OBSERVER.ownerKey()));
assertNotNull(denied, "a missing ticket must be refused, not treated as owned");
}
/**
* {@code fleet_status} stays refused for an observer even once {@code TASK_READ} opens: it is
* not ticket-scoped, so no call site supplies a real {@code ticketOwnedByCaller} classifier,
* and the default 3-argument {@code denyFor} fails closed.
*/
@Test
void anObserverMayNotReadSessionStatusEvenAfterTaskReadOpensForTickets() {
FleetMcp m = mcp(true);
assertNotNull(m.denyFor(OBSERVER, Authz.Action.TASK_READ, "term_a"),
"a session id is not a ticket the caller could ever own");
}
@Test
void theLegacyConstructorLeavesTheGateOpen() {
// The 22 pre-existing FleetMcpTest cases rely on no authorization being enforced.
@@ -477,9 +567,10 @@ class FleetMcpAuthzTest {
/**
* {@link FleetMcp#leadsVisibleTo} is the whole policy decision for {@code fleet_list}'s
* {@code leads} array: visible to exactly the roles that may {@code SEND} to a lead -- the
* primary, an architect, and a collaborator -- never a worker, which holds {@code READ} but
* can never {@code SEND} at all, and never an anonymous caller.
* {@code leads} array: visible to the primary, an architect, and a collaborator -- never a
* worker, which holds {@code READ} but can never {@code SEND} at all, never an observer, which
* reads a lead's address from its filtered {@code panes} rows instead, and never an anonymous
* caller.
*/
@Test
void primaryArchitectAndCollaboratorMaySeeTheLeadsArray() {
@@ -489,6 +580,9 @@ class FleetMcpAuthzTest {
"a collaborator may SEND to a lead, so it must see the leads array to learn where");
assertFalse(FleetMcp.leadsVisibleTo(WORKER_A),
"a worker holds READ but can never SEND, so it must not see the leads array");
assertFalse(FleetMcp.leadsVisibleTo(Principal.observer("term_obs", 700)),
"an observer may SEND to a lead but learns the address from its panes rows, which "
+ "carry no lead name, context window or config dir");
assertFalse(FleetMcp.leadsVisibleTo(ANON), "authenticated as nothing must not see it either");
}
@@ -513,8 +607,8 @@ class FleetMcpAuthzTest {
/**
* {@link FleetMcp#panesVisibleTo} is the whole policy decision for {@code fleet_list}'s
* {@code panes} array: visible to every role that may {@code SEND} to some other pane -- the
* primary, an architect, a collaborator, and an observer (to another observer pane only, with
* its row filtered and reduced -- see {@code listFleet}) -- never a worker, never an
* primary, an architect, a collaborator, and an observer (to a lead or another observer pane,
* with its rows filtered and reduced -- see {@code listFleet}) -- never a worker, never an
* anonymous caller.
*/
@Test
@@ -525,8 +619,8 @@ class FleetMcpAuthzTest {
assertFalse(FleetMcp.panesVisibleTo(WORKER_A),
"a worker holds READ but can never SEND, so it must not see the panes array");
assertTrue(FleetMcp.panesVisibleTo(Principal.observer("term_obs", 700)),
"an observer holds SEND to another observer pane, so it must see the (filtered, "
+ "reduced) panes array");
"an observer holds SEND to a lead and to another observer pane, so it must see the "
+ "(filtered, reduced) panes array");
assertFalse(FleetMcp.panesVisibleTo(ANON), "authenticated as nothing must not see it either");
}
@@ -0,0 +1,178 @@
package dev.ltms.fleet.mcp;
import dev.ltms.fleet.Fleetd;
import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.inject.TurnListener;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.session.FakeWorktrees;
import dev.ltms.fleet.session.SessionManager;
import io.modelcontextprotocol.client.McpClient;
import io.modelcontextprotocol.client.McpSyncClient;
import io.modelcontextprotocol.client.transport.HttpClientStreamableHttpTransport;
import io.modelcontextprotocol.spec.McpClientTransport;
import io.modelcontextprotocol.spec.McpSchema;
import org.eclipse.jetty.server.Server;
import org.eclipse.jetty.server.ServerConnector;
import org.eclipse.jetty.servlet.ServletContextHandler;
import org.eclipse.jetty.servlet.ServletHolder;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.function.Predicate;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* An observer's {@code fleet_send} to a lead, driven end to end: a real MCP client over a real
* HTTP transport, resolved by the real {@link CallerResolver} to
* {@link dev.ltms.fleet.auth.Role#OBSERVER}, through the real {@link MessageService} and the real
* {@link Injector} to the point of its real herdr call.
*
* <p>{@link FleetMcpAuthzTest} proves {@code denyFor} grants this case. This is the path that also
* proves the grant is not dead at the injector's readiness gate: that gate is the production
* {@link Fleetd#deliverableTo} predicate, and the lead's terminal carries no
* {@link MemberPresence} entry, so the delivery can only pass by the lead being a configured lead.
*/
class FleetMcpObserverSendToLeadDeliveryTest {
private static final String LEAD = "term_lead_pane";
private static final String COLLABORATOR = "term_collab_pane";
private final FakeHerdr herdr = new FakeHerdr();
private final AgentControl agents = new AgentControl(herdr);
private final Rendezvous rendezvous = new Rendezvous();
private final MemberPresence presence = new MemberPresence();
private Injector injector;
private MessageService messages;
private FleetMcp mcp;
private Server server;
private String baseUrl;
@BeforeEach
void startServer() throws Exception {
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(agents, new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> "tok");
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
// The fake's pane list carries a second pane, "term_shell", whose shell pid is 9001 and
// which hosts no agent -- so the real resolver lands the caller on the observer floor.
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 9001L);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
() -> Map.of(LEAD, "fleet01-lead"), new MemberRegistry(null),
_ -> null, () -> Map.of(COLLABORATOR, "ops"));
Predicate<String> deliverable =
Fleetd.deliverableTo(presence, callers::leads, callers::collaborators);
injector = new Injector(agents, TurnListener.NOOP, deliverable);
messages = new MessageService(agents, injector, rendezvous, new InMemoryReplyInbox());
mcp = new FleetMcp(messages, workers, sessions, identity, presence,
new PrimaryRegistry(null), callers, FleetMcp.AuthorizationMode.ENFORCED,
null, FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
FleetMcp.QuarantineSource.none(), null, FleetMcp.OutageSource.none(),
FleetMcp.LeadSeatSource.none(), List.of(), null);
ServletContextHandler handler = new ServletContextHandler();
handler.setContextPath("/");
handler.addServlet(new ServletHolder(mcp.servlet()), "/mcp");
server = new Server(0);
server.setHandler(handler);
server.start();
baseUrl = "http://127.0.0.1:"
+ ((ServerConnector) server.getConnectors()[0]).getLocalPort();
}
@AfterEach
void tearDown() throws Exception {
if (server != null) {
server.stop();
}
if (mcp != null) {
mcp.close();
}
}
@Test
void anObserversSendToALeadIsAttributedAndReachesTheRealInjector() throws Exception {
assertFalse(presence.isPresent(LEAD),
"premise: the lead's terminal is deliverable only as a configured lead, never "
+ "through a presence entry");
McpSchema.CallToolResult result = sendFleetSend(LEAD, "can we split the review?");
assertFalse(result.isError(), "an observer sending to a lead must be accepted: "
+ textOf(result));
long waiterDeadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(LEAD) && System.currentTimeMillis() < waiterDeadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(LEAD), "the async send must have opened its rendezvous waiter");
injector.onStatus(LEAD, AgentStatus.IDLE); // drives the real delivery attempt to herdr
long deliveryDeadline = System.currentTimeMillis() + 3000;
while (!herdr.called("agent.prompt") && System.currentTimeMillis() < deliveryDeadline) {
Thread.sleep(5);
}
assertTrue(herdr.called("agent.prompt"), "the delivery attempt must have reached herdr");
@SuppressWarnings("unchecked")
Map<String, Object> params = (Map<String, Object>) herdr.lastCall("agent.prompt").params();
assertEquals("[fleet_send from observer term_shell]\ncan we split the review?",
params.get("text"),
"a lead must see the sender's own daemon-resolved terminal, never a raw echo of "
+ "the content and never a client-supplied name");
}
/**
* Control for the test above, over the same live server and the same wiring: the grant is
* specific to a lead target, so a collaborator's terminal is still refused at the handler and
* nothing is ever queued for it.
*/
@Test
void thatSameObserverIsStillRefusedACollaboratorsTerminal() {
McpSchema.CallToolResult result = sendFleetSend(COLLABORATOR, "can we split the review?");
assertTrue(result.isError(), "an observer must not reach a collaborator's terminal");
assertFalse(rendezvous.isWaiting(COLLABORATOR),
"a refused send must never open a waiter for its target");
}
private McpSchema.CallToolResult sendFleetSend(String target, String content) {
McpClientTransport transport = HttpClientStreamableHttpTransport.builder(baseUrl)
.endpoint("/mcp")
.build();
try (McpSyncClient client = McpClient.sync(transport).build()) {
client.initialize();
return client.callTool(McpSchema.CallToolRequest.builder("fleet_send")
.arguments(Map.of("sessionId", target, "content", content, "wait", false))
.build());
}
}
private static String textOf(McpSchema.CallToolResult r) {
return ((McpSchema.TextContent) r.content().getFirst()).text();
}
}
@@ -78,7 +78,7 @@ class FleetMcpTest {
private void assertSendRoundTrips(String target, Set<String> profiles) throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles, null));
() -> FleetMcp.send(messages, target, "hi", 4000L, null, profiles, null, true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -104,7 +104,7 @@ class FleetMcpTest {
void sendThenReplyRoundTrips() throws Exception {
// fleet_send blocks; fleet_reply resolves it with the worker's structured answer.
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of(), null));
() -> FleetMcp.send(messages, "term_a", "review this", 4000L, null, Set.of(), null, true));
// Wait until the send has opened its waiter so the reply resolves it (CB-307: reply now
// queues in the inbox if no waiter is open, which would break the round-trip).
@@ -182,7 +182,7 @@ class FleetMcpTest {
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null, true));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
@@ -280,7 +280,7 @@ class FleetMcpTest {
assertEquals("late reply", messages.drainReplies("term_a").getFirst().content());
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L, null));
() -> FleetMcp.answer(messages, firstTurnId, "config.yaml", 5000L, null, true));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -298,7 +298,7 @@ class FleetMcpTest {
@Test
void aDifferentCallersMcpAnswerIsRefusedForABlockingSendButTheRealOwnerSucceeds() throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "do X", 5000L, null, Set.of(), "term_owner"));
() -> FleetMcp.send(messages, T, "do X", 5000L, null, Set.of(), "term_owner", true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
@@ -313,12 +313,12 @@ class FleetMcpTest {
String afterTurnId = questionText.substring(questionText.indexOf("turnId=\"") + "turnId=\"".length());
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker");
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L, "term_attacker", true);
assertTrue(hijacked.isError(), "a caller that did not open this turn must get an error, not an answer");
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner"));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, "term_owner", true));
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
@@ -357,14 +357,14 @@ class FleetMcpTest {
String turnId = asking.turnId();
McpSchema.CallToolResult hijacked = FleetMcp.answer(messages, turnId, "evil.yaml", 500L,
"worker:term_attacker");
"worker:term_attacker", true);
assertTrue(hijacked.isError(), "a caller that did not create this delegation must get an error");
assertFalse(ask.isDone(), "a refused answer must not resolve the worker's blocked fleet_ask");
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase(),
"a refused answer must not advance the async ticket's phase");
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, owner.ownerKey()));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, owner.ownerKey(), true));
assertEquals("config.yaml", textOf(ask.get(5, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
@@ -383,11 +383,71 @@ class FleetMcpTest {
@Test
void sendTimesOutWithAWorkingNote() {
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null);
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null, true);
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
}
/**
* fleetd #801: a blocking send creates no ticket, ever, so the TIMED_OUT_QUEUED receipt must
* name the one route that recovers a late answer — draining the target's inbox — and must not
* invite a resend, which can duplicate the delivery.
*/
@Test
void sendTimesOutQueuedNamesFleetPollTargetAndNotTheBareWordRetry() {
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null, true);
String text = textOf(res);
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(text.contains("fleet_poll{target=\"term_a\"}"), "got: " + text);
assertFalse(text.contains("retry"), "got: " + text);
}
/** As above, for the TIMED_OUT_WORKING arm (delivered, but the worker never replies). */
@Test
void sendTimesOutWorkingNamesFleetPollTargetAndNotTheBareWordRetry() throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null, true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker starts, never replies
McpSchema.CallToolResult res = send.get(5, TimeUnit.SECONDS);
String text = textOf(res);
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(text.contains("fleet_poll{target=\"" + T + "\"}"), "got: " + text);
assertFalse(text.contains("retry"), "got: " + text);
}
/** As above, for the BUSY arm (another send already held the session for the whole window). */
@Test
void sendTimesOutBusyNamesFleetPollTargetAndNotTheBareWordRetry() throws Exception {
CompletableFuture<McpSchema.CallToolResult> first = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "first", 5000L, null, Set.of(), null, true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(T), "the first send should hold the session lock");
McpSchema.CallToolResult busy = FleetMcp.send(messages, T, "second", 100L, null, Set.of(), null, true);
String text = textOf(busy);
assertNotEquals(Boolean.TRUE, busy.isError(), "a timeout is informational, not a tool error");
assertTrue(text.contains("fleet_poll{target=\"" + T + "\"}"), "got: " + text);
assertFalse(text.contains("retry"), "got: " + text);
// Release the first send so the test thread is not left pinned.
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
FleetMcp.reply(messages, T, Role.WORKER, "first done");
first.get(5, TimeUnit.SECONDS);
}
/**
* fleetd #571 (ticket CORRECTION 5): {@code formatReply}'s {@code TIMED_OUT_UNCONFIRMED} arm is
* the one message whose whole job is to stop a caller retrying a delivery that may already have
@@ -399,7 +459,7 @@ class FleetMcpTest {
void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception {
herdr.agentSendFailsWith("send_failed");
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null));
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of(), null, true));
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -417,15 +477,59 @@ class FleetMcpTest {
+ "a resend here can double-deliver the same brief: got " + text);
}
/**
* fleetd #801 defect 1: {@code SEND} is granted to callers {@code DRAIN} refuses (an architect,
* a collaborator, an observer), so a TIMED_OUT/BUSY receipt must not name {@code fleet_poll}
* when the caller's own grant would be refused there — it must say the reply is unrecoverable
* instead.
*/
@Test
void sendTimesOutWithoutDrainGrantNamesNoRouteAndAdmitsItCannotBeRecovered() {
McpSchema.CallToolResult res = FleetMcp.send(messages, "term_a", "hi", 120L, null, Set.of(), null, false);
String text = textOf(res);
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertFalse(text.contains("fleet_poll{target"), "got: " + text);
assertTrue(text.contains("cannot be recovered"), "got: " + text);
}
/** As above, for {@code answer()} — {@code ANSWER} is primary-or-architect, and an architect is refused DRAIN too. */
@Test
void answerTimesOutWithoutDrainGrantNamesNoRoute() throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null, true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"), "the send must be waiting for the ask to surface to");
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
() -> FleetMcp.ask(messages, "term_a", "which config?", 5000L));
McpSchema.CallToolResult q = send.get(6, TimeUnit.SECONDS);
String qt = textOf(q);
String afterMarker = qt.substring(qt.indexOf("turnId=\"") + "turnId=\"".length());
String turnId = afterMarker.substring(0, afterMarker.indexOf('"'));
// The resumed worker never replies, so the answer's own wait times out.
McpSchema.CallToolResult answer = FleetMcp.answer(messages, turnId, "config.yaml", 150L, null, false);
String text = textOf(answer);
assertNotEquals(Boolean.TRUE, answer.isError(), "a timeout is informational, not a tool error");
assertFalse(text.contains("fleet_poll{target"), "got: " + text);
assertTrue(text.contains("cannot be recovered"), "got: " + text);
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
}
@Test
void sendRejectsMissingArgs() {
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of(), null).isError());
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of(), null).isError());
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of(), null, true).isError());
assertTrue(FleetMcp.send(messages, "term_a", " ", null, null, Set.of(), null, true).isError());
}
@Test
void sendRejectsAConfiguredProfileNameBeforeAcceptingIt() {
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"), null);
McpSchema.CallToolResult blocking = FleetMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"), null, true);
McpSchema.CallToolResult async = FleetMcp.sendAsync(messages, "sol", "hi", null, Set.of("sol"));
assertTrue(blocking.isError());
@@ -444,7 +548,7 @@ class FleetMcpTest {
assertSendRoundTrips("term_live_member", profiles);
// A herdr-owned pane outside the bridge roster cannot be classified at accept time.
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles, null);
McpSchema.CallToolResult result = FleetMcp.send(messages, "external-pane", "hi", 10L, null, profiles, null, true);
assertFalse(result.isError(), "an unclassified target must not be rejected at acceptance time");
}
@@ -536,7 +640,7 @@ class FleetMcpTest {
void askThenAnswerRoundTrips() throws Exception {
// The primary delegates and blocks; wait until its waiter is open before the worker asks.
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null));
() -> FleetMcp.send(messages, "term_a", "do X", 5000L, null, Set.of(), null, true));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -558,7 +662,7 @@ class FleetMcpTest {
// The primary answers via fleet_send(turnId); this blocks again for the worker's reply.
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null));
() -> FleetMcp.answer(messages, turnId, "config.yaml", 5000L, null, true));
// The worker's ask returns the answer — it resumes the same turn.
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
@@ -584,7 +688,7 @@ class FleetMcpTest {
@Test
void answerToAStaleTurnIsAnError() {
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L, null);
McpSchema.CallToolResult res = FleetMcp.answer(messages, "term_a#999", "too late", 500L, null, true);
assertTrue(res.isError());
assertTrue(textOf(res).contains("no longer open"), textOf(res));
}
@@ -1239,7 +1343,7 @@ class FleetMcpTest {
/**
* fleetd #756: a pane bound to a configured architect slot with no live member session must
* report {@code role: "architect"}, read from {@link CallerResolver#boundToArchitectSlot} —
* the same classifier {@link CallerResolver#sendableObserverTarget} refuses as a {@code SEND}
* the same classifier {@link CallerResolver#observerSendTarget} refuses as a {@code SEND}
* target — rather than falling through to {@code "observer"}.
*/
@Test
@@ -1276,14 +1380,13 @@ class FleetMcpTest {
}
/**
* fleetd #758: an observer's {@code fleet_list} now carries a {@code panes} key, filtered to
* {@link CallerResolver#sendableObserverTarget} (so a lead's pane, a spawned member's pane, a
* collaborator's pane, and an unoccupied architect-slot pane are all absent) and every
* surviving row reduced to exactly {@code sessionId}, {@code label}, {@code status},
* An observer's {@code fleet_list} carries a {@code panes} key, filtered to
* {@link CallerResolver#observerSendTarget} (so a lead's pane survives, while a spawned
* member's pane, a collaborator's pane and an unoccupied architect-slot pane are all absent)
* and every surviving row reduced to exactly {@code sessionId}, {@code label}, {@code status},
* {@code role}, {@code deliverable} — never {@code paneId}, {@code workspaceId},
* {@code workspaceLabel} (fleetd #771 — a space name is host shape, a stronger disclosure than
* a pane id, so it stays out of the reduced row too), {@code tabId}, {@code agentType}, or
* {@code cwd}.
* {@code workspaceLabel} (a space name is host shape, a stronger disclosure than a pane id, so
* it stays out of the reduced row too), {@code tabId}, {@code agentType}, or {@code cwd}.
*/
@Test
void listFiltersAndReducesThePanesArrayForAnObserver() {
@@ -1306,7 +1409,7 @@ class FleetMcpTest {
FleetMcp.PaneSource panes = new FleetMcp.PaneSource(
() -> new PaneLocator(h).tabLabelsByTabId(),
() -> new PaneLocator(h).workspaceLabelsByWorkspaceId(), _ -> true,
callers::boundToArchitectSlot, callers.sendableObserverTarget());
callers::boundToArchitectSlot, callers.observerSendTarget());
String out = textOf(listFleetAsObserver(h, callers.leads(),
Map.of("term_collab_pane", "ops"), panes));
@@ -1314,7 +1417,10 @@ class FleetMcpTest {
assertTrue(out.contains("\"panes\":["), out);
assertTrue(out.contains("\"sessionId\":\"term_sendable\""),
"an ordinary unclassified pane must still be sendable and visible: " + out);
assertFalse(out.contains("term_lead_pane"), "a lead's pane must not be enumerated: " + out);
assertTrue(out.contains("\"sessionId\":\"term_lead_pane\""),
"a lead's pane is where an observer reads the sessionId its send needs: " + out);
assertTrue(out.contains("\"role\":\"lead\""),
"the lead's row must name the role, so an observer can tell it from a peer pane: " + out);
assertFalse(out.contains("term_member_pane"), "a spawned member's pane must not be enumerated: " + out);
assertFalse(out.contains("term_collab_pane"), "a collaborator's pane must not be enumerated: " + out);
assertFalse(out.contains("term_architect_pane"),
@@ -266,4 +266,61 @@ class PrimaryRegistryTest {
assertEquals("term_opus_after_roll", reg.nudgeTargetFor("term_never_seen").orElseThrow());
}
// ── fleetd #778: a nudge must offer only what the delegator's own role may run ──────────────
@Test
void theShortDelegationOverloadsGrantEverythingAPrimaryMayRun() {
var reg = new PrimaryRegistry(null);
reg.recordDelegation("term_worker", "term_lead");
assertEquals(Boolean.TRUE, reg.nudgeMayDrainFor("term_worker").orElseThrow());
assertEquals(Boolean.TRUE, reg.nudgeMayAnswerFor("term_worker").orElseThrow());
}
@Test
void theFullOverloadCarriesTheCallersOwnGrants() {
var reg = new PrimaryRegistry(null);
reg.recordDelegation("term_worker", "term_observer", "observer-name", false, false);
assertEquals(Boolean.FALSE, reg.nudgeMayDrainFor("term_worker").orElseThrow());
assertEquals(Boolean.FALSE, reg.nudgeMayAnswerFor("term_worker").orElseThrow());
}
@Test
void theTwoGrantsAreIndependent() {
var reg = new PrimaryRegistry(null);
reg.recordDelegation("term_worker", "term_architect", "architect-name", false, true);
assertEquals(Boolean.FALSE, reg.nudgeMayDrainFor("term_worker").orElseThrow(),
"an architect may not DRAIN");
assertEquals(Boolean.TRUE, reg.nudgeMayAnswerFor("term_worker").orElseThrow(),
"an architect may ANSWER");
}
@Test
void theSingletonPrimaryFallbackIsAlwaysGrantedEverything() {
var reg = new PrimaryRegistry("term_pinned");
assertEquals(Boolean.TRUE, reg.nudgeMayDrainFor("term_never_seen").orElseThrow());
assertEquals(Boolean.TRUE, reg.nudgeMayAnswerFor("term_never_seen").orElseThrow());
}
@Test
void withNoDelegationAndNoPrimaryTheGrantsAreUnknown() {
var reg = new PrimaryRegistry(null);
assertTrue(reg.nudgeMayDrainFor("term_never_seen").isEmpty());
assertTrue(reg.nudgeMayAnswerFor("term_never_seen").isEmpty());
}
@Test
void forgettingADelegationForgetsItsGrantsToo() {
var reg = new PrimaryRegistry(null);
reg.recordDelegation("term_worker", "term_observer", "observer-name", false, false);
reg.forgetDelegation("term_worker");
assertTrue(reg.nudgeMayDrainFor("term_worker").isEmpty());
assertTrue(reg.nudgeMayAnswerFor("term_worker").isEmpty());
}
}
@@ -223,4 +223,19 @@ class LeadCoordLoopTest {
assertEquals(List.of(), channel.acked());
}
@Test
void aDisabledGateDeliversAHeldMessageDespiteUnsubmittedTextInThePromptBox() {
var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked"));
var herdr = new FakeHerdr().detectionText(FakeHerdr.DRAFTED_PROMPT_CARET).agentStatus("idle");
var loop = new LeadCoordLoop(channel, new AgentControl(herdr), () -> Map.of(LEAD_TERM, SELF),
NO_SCHEDULER, 3_000L, false);
loop.tick();
assertEquals(1, prompts(herdr).size(), "a disabled gate delivers even though the box holds a draft");
assertEquals(List.of("m1"), channel.acked());
assertTrue(herdr.calls.stream().noneMatch(c -> c.method().equals("agent.read")),
"a disabled gate never reads the pane");
}
}
@@ -25,6 +25,7 @@ import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.*;
@@ -521,12 +522,20 @@ class LeadHeartbeatLoopTest {
private static LeadHeartbeatLoop tickableLoop(FailableHerdrClient herdr, AtomicLong now,
List<MemberSession>[] rosterBox, InMemoryReplyInbox inbox,
int quietNudgeCap, ScheduledExecutorService scheduler) {
return tickableLoop(herdr, now, rosterBox, inbox, quietNudgeCap, scheduler, true);
}
/** As above, with the prompt-box gate's on/off switch. */
private static LeadHeartbeatLoop tickableLoop(FailableHerdrClient herdr, AtomicLong now,
List<MemberSession>[] rosterBox, InMemoryReplyInbox inbox,
int quietNudgeCap, ScheduledExecutorService scheduler,
boolean promptBoxGateEnabled) {
AgentControl agents = new AgentControl(herdr);
PrimaryRegistry registry = new PrimaryRegistry(LEAD);
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, agents, inbox, scheduler, 5, 100_000);
return new LeadHeartbeatLoop(registry, agents, inbox, () -> rosterBox[0], pushLoop, scheduler,
now::get, IDLE_AFTER_NANOS, 100_000L, quietNudgeCap, null,
new LeadHeartbeatLoop.LeadContextSource(t -> highReading()), true);
new LeadHeartbeatLoop.LeadContextSource(t -> highReading()), true, true, promptBoxGateEnabled);
}
@Test
@@ -695,6 +704,7 @@ class LeadHeartbeatLoopTest {
private final String lead;
private final List<String> sentTexts = new ArrayList<>();
private final List<String> promptTargets = new ArrayList<>();
private final AtomicInteger reads = new AtomicInteger();
private boolean throwOnNextSend = false;
/** What {@code agent.read} reports — the loop reads the lead's input box before it nudges. */
private String paneTail = FakeHerdr.IDLE_PROMPT_CARET;
@@ -720,6 +730,10 @@ class LeadHeartbeatLoopTest {
return List.copyOf(promptTargets);
}
int readCount() {
return reads.get();
}
@Override
@SuppressWarnings("unchecked")
public JsonNode call(String method, Object params) {
@@ -730,6 +744,7 @@ class LeadHeartbeatLoopTest {
.put("agent_status", "idle"));
}
if ("agent.read".equals(method)) {
reads.incrementAndGet();
return MAPPER.createObjectNode()
.set("read", MAPPER.createObjectNode().put("text", paneTail));
}
@@ -794,4 +809,23 @@ class LeadHeartbeatLoopTest {
"a pane whose box cannot be found may be holding a draft");
}
@Test
void aDisabledGateSendsTheNudgeDespiteUnsubmittedTextInThePromptBox() {
var herdr = new FailableHerdrClient(LEAD);
herdr.paneTail(FakeHerdr.DRAFTED_PROMPT_CARET);
var now = new AtomicLong(NOW);
@SuppressWarnings("unchecked")
List<MemberSession>[] rosterBox = new List[]{List.of()};
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
LeadHeartbeatLoop loop = tickableLoop(herdr, now, rosterBox, inbox, 0, scheduler, false);
loop.tick(); // opens the idle window
now.addAndGet(TimeUnit.SECONDS.toNanos(400));
loop.tick(); // a disabled gate must not hold on the operator's draft
assertEquals(1, herdr.sentTexts().size(),
"a disabled gate must send the nudge even though the box holds a draft");
assertEquals(0, herdr.readCount(), "a disabled gate must never read the lead's pane");
}
}
@@ -107,6 +107,219 @@ class MessageServiceTest {
"an unreplied but finished turn resolves via the completion fallback");
assertEquals("BUILD GREEN: 391 files", reply.text(), "the scraped transcript tail is returned");
assertTrue(reply.completed(), "a scraped completion still counts as completed");
assertTrue(messages.drainReplies(T).isEmpty(),
"a completion resolved by a live caller must not also land in the inbox");
}
/** Poll {@code messages.drainReplies(target)} until it is non-empty or the deadline passes. */
private java.util.List<ReplyInbox.InboxMessage> awaitDrained(String target) throws InterruptedException {
var drained = messages.drainReplies(target);
long deadline = System.currentTimeMillis() + 2000;
while (drained.isEmpty() && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
drained = messages.drainReplies(target);
}
return drained;
}
/**
* Polls {@code target}'s inbox for {@code millis} and returns {@code true} only if nothing new
* ever arrived. Used to assert an absence against an asynchronous completion race: the window
* has to be long enough for a virtual thread already scheduled to actually run.
*/
private boolean noNewReplyWithin(String target, long millis) throws InterruptedException {
long deadline = System.currentTimeMillis() + millis;
while (System.currentTimeMillis() < deadline) {
if (!messages.drainReplies(target).isEmpty()) {
return false;
}
//noinspection BusyWait
Thread.sleep(5);
}
return true;
}
/**
* fleetd #801: {@code send}'s TimeoutException branch never completes {@code reply} itself, and
* {@link Rendezvous#close} only deregisters it from future session lookups — it does not cancel
* the future. The CB-106 completion fallback resolves the exact captured future (CB-116), so it
* can still succeed after the caller gave up and returned TIMED_OUT_WORKING. With nobody left
* awaiting that future, the late scrape must land in the target's inbox instead of being
* silently discarded.
*/
@Test
void aCompletionThatArrivesAfterTheCallerGaveUpLandsInTheInbox() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync("do the task", 150);
awaitWaiting();
herdr.readText("$ prompt"); // pre-turn baseline
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker starts — no reply yet
MessageService.Reply timedOut = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, timedOut.outcome(),
"the caller gives up while the worker is still delivered and working");
// The worker's turn finishes AFTER the caller already gave up and nobody is listening.
herdr.readText("late answer, nobody was listening");
injector.onStatus(T, AgentStatus.IDLE); // working -> idle: the completion fallback fires
var drained = awaitDrained(T);
assertEquals(1, drained.size(), "the late completion must land in the inbox, not vanish");
assertEquals("late answer, nobody was listening", drained.get(0).content());
}
/**
* fleetd #808: once a send has timed out, an explicit {@code fleet_reply} for that same turn
* must land in the inbox exactly once — not twice, with the second copy a completion scrape of
* the same answer the reply already delivered.
*/
@Test
void anExplicitReplyAfterATimeoutLandsExactlyOnceNotTwice() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync("do the task", 150);
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE); // deliver
injector.onStatus(T, AgentStatus.WORKING); // worker starts — no reply yet
MessageService.Reply timedOut = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, timedOut.outcome(),
"the caller gives up while the worker is still delivered and working");
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "structured answer"),
"no live waiter is open by now, so the explicit reply is held for a later drain");
var afterReply = awaitDrained(T);
assertEquals(1, afterReply.size(), "the explicit reply must land in the inbox exactly once");
assertEquals("structured answer", afterReply.get(0).content());
// The worker's turn then finishes for real; the completion fallback races the same waiter
// the reply already claimed, and must not publish a second entry for it.
herdr.readText("late scrape that must not also publish");
injector.onStatus(T, AgentStatus.IDLE);
assertTrue(noNewReplyWithin(T, 500),
"a completion scrape for a turn that already got an explicit fleet_reply must not "
+ "publish a second entry");
}
/**
* fleetd #808: the suppression above is per-turn, not per-target. A second, independent turn on
* the same member that never calls {@code fleet_reply} must still have its completion scrape
* land, even though an earlier turn on that same target already got an explicit reply.
*/
@Test
void theSuppressionIsPerTurnNotPerTarget() throws Exception {
// First turn: times out, then gets an explicit fleet_reply.
CompletableFuture<MessageService.Reply> first = sendAsync("first task", 150);
awaitWaiting();
herdr.readText("$ prompt 1");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, first.get(5, TimeUnit.SECONDS).outcome());
assertEquals(MessageService.ReplyOutcome.QUEUED, messages.reply(T, "first answer"));
var firstDrain = awaitDrained(T);
assertEquals(1, firstDrain.size());
assertEquals("first answer", firstDrain.get(0).content());
// Finishing turn 1 for real also frees the worker for turn 2, and exercises that turn 1's
// own late scrape stays suppressed right up to that point.
herdr.readText("first turn's late scrape — must stay suppressed");
injector.onStatus(T, AgentStatus.IDLE);
assertTrue(noNewReplyWithin(T, 300), "turn 1's own scrape must still be suppressed");
// Second, independent turn on the same target: times out, and no fleet_reply is ever called.
CompletableFuture<MessageService.Reply> second = sendAsync("second task", 150);
awaitWaiting();
herdr.readText("$ prompt 2");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, second.get(5, TimeUnit.SECONDS).outcome());
herdr.readText("second turn's late scrape — must still land");
injector.onStatus(T, AgentStatus.IDLE);
var secondDrain = awaitDrained(T);
assertEquals(1, secondDrain.size(),
"a later, independent turn's completion must still land even though an earlier "
+ "turn on the same target already got an explicit reply");
assertEquals("second turn's late scrape — must still land", secondDrain.get(0).content());
}
/**
* fleetd #801 defect 2: the publish inside {@code strandLateResolution}'s {@code whenComplete}
* callback runs with no caller left to surface a throw to — {@code AmqpReplyInbox.publish} can
* throw {@code IllegalStateException}, and an exception thrown inside a {@code whenComplete}
* action is captured into the discarded dependent stage, never rethrown. It must be caught and
* logged there, or a broker failure on this path is invisible and the late reply is just lost.
*/
@Test
void aLateCompletionThatFailsToPublishIsLoggedNotLost() throws Exception {
FakeHerdr localHerdr = new FakeHerdr();
Rendezvous localRendezvous = new Rendezvous();
java.util.concurrent.atomic.AtomicLong localClock = new java.util.concurrent.atomic.AtomicLong();
CompletionResolver localCompletion = new CompletionResolver(new AgentControl(localHerdr), localRendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(),
() -> localClock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1));
Injector localInjector = new Injector(new AgentControl(localHerdr), localCompletion);
ReplyInbox throwingInbox = new ReplyInbox() {
private final InMemoryReplyInbox delegate = new InMemoryReplyInbox();
@Override public void own(String target) { delegate.own(target); }
@Override public void release(String target) { delegate.release(target); }
@Override public void publish(String target, String msgId, String content) {
throw new IllegalStateException("PROBE-801-strandLateResolution");
}
@Override public java.util.List<InboxMessage> peek(String target) { return delegate.peek(target); }
@Override public boolean ack(String target, String msgId) { return delegate.ack(target, msgId); }
};
MessageService localMessages = new MessageService(new AgentControl(localHerdr), localInjector,
localRendezvous, throwingInbox);
CompletableFuture<MessageService.Reply> send = CompletableFuture.supplyAsync(
() -> localMessages.send(T, "do the task", 150, (String) null));
long deadline = System.currentTimeMillis() + 2000;
while (!localRendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(localRendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
localHerdr.readText("$ prompt");
localInjector.onStatus(T, AgentStatus.IDLE); // deliver
localInjector.onStatus(T, AgentStatus.WORKING); // worker starts — no reply yet
MessageService.Reply timedOut = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, timedOut.outcome(),
"the caller gives up while the worker is still delivered and working");
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
try {
// The worker's turn finishes AFTER the caller gave up — the completion fallback fires,
// and its publish into throwingInbox fails.
localHerdr.readText("late answer, publish fails");
localInjector.onStatus(T, AgentStatus.IDLE);
long logDeadline = System.currentTimeMillis() + 2000;
while (appender.list.isEmpty() && System.currentTimeMillis() < logDeadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertFalse(appender.list.isEmpty(), "a failed late-completion publish must be logged");
assertEquals(Level.WARN, appender.list.get(0).getLevel());
} finally {
detachMessageServiceLog(appender);
}
assertFalse(localMessages.hasStrandedReply(T),
"strandedReplies must not be set when the publish it depends on failed");
}
/**
@@ -1034,6 +1247,39 @@ class MessageServiceTest {
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
}
/**
* {@link MessageService#ownsTicket(String, String)} is the side-effect-free ownership check
* the authorization gate consults ahead of {@link MessageService#poll(String, String)}, which
* also fires the ticket's collection hook. Read-only: polling the same ticket afterwards still
* sees it, which {@link #theOneArgPollOverloadBypassesOwnershipEntirely} nearby does not
* guarantee for every overload.
*/
@Test
void publicOwnsTicketIsASideEffectFreeOwnershipCheck() {
Principal lead = Principal.leader("opus", "term_lead", 1);
String ticket = messages.sendAsync(T, "long task", null, lead);
assertTrue(messages.ownsTicket(ticket, lead.ownerKey()));
assertFalse(messages.ownsTicket(ticket, Principal.anonymous().ownerKey()));
assertFalse(messages.ownsTicket("no-such-ticket", lead.ownerKey()),
"a ticket this daemon never heard of is owned by nobody");
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket, lead.ownerKey()).phase(),
"the ownership check must not have consumed or altered the ticket");
}
/**
* fleetd #803: the gap between "unknown ticket" (above) and "no ticket at all". A null ticket
* used to reach {@code ConcurrentHashMap.get(null)} and throw a {@code NullPointerException}
* instead of returning false.
*/
@Test
void publicOwnsTicketWithANullTicketReturnsFalseRatherThanThrowing() {
Principal lead = Principal.leader("opus", "term_lead", 1);
assertFalse(messages.ownsTicket(null, lead.ownerKey()),
"no ticket at all is owned by nobody, the same as an unknown one");
}
@Test
void architectOwnershipUsesTerminalRatherThanSlot() {
Principal oldArchitect = Principal.architect("opus", "term_OLD", 1);
@@ -339,6 +339,85 @@ class ReplyPushLoopTest {
assertEquals(0, rec.promptTargets().size(), "nobody live was found, so nothing was ever sent");
}
// --- fleetd #778: never nudge a delegator to run a call its own role is refused -------------
@Test
void onReplyQueuedNeverNudgesADelegatorThatMayNotRunDrain() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
loop(1, 50).onReplyQueued(WORKER);
Thread.sleep(200);
assertEquals(0, rec.sendCount(),
"a delegator whose role cannot run fleet_poll(target=...) must never be nudged to run it");
}
@Test
void onReplyQueuedStillNudgesADelegatorThatMayRunDrain() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
registry.recordDelegation(WORKER, OTHER_PRIMARY, "lead", true, true);
loop(1, 50).onReplyQueued(WORKER);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("fleet_poll(target=" + WORKER + ")"),
"a delegator that may run DRAIN must still be nudged about it: " + nudge);
}
@Test
void onQuestionOpenedNeverNudgesADelegatorThatMayNotRunAnswer() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
Thread.sleep(200);
assertEquals(0, rec.sendCount(),
"a delegator whose role cannot run fleet_send(turnId=...) must never be nudged to run it");
}
@Test
void onQuestionOpenedStillNudgesADelegatorThatMayRunAnswer() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
registry.recordDelegation(WORKER, OTHER_PRIMARY, "architect", false, true);
loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("fleet_send(turnId="),
"an architect may run ANSWER, so the nudge must still name it: " + nudge);
}
/**
* A delegator barred from DRAIN can still have a ticket genuinely pending for it — tickets are
* out of this unit's scope (fleetd #778 Part 1 already makes a ticket nudge correct for an
* observer). The reply half must still be suppressed even though the ticket half fires.
*/
@Test
void aForbiddenReplyNudgeIsSuppressedWhileAnEligibleTicketStillFires() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
registry.recordDelegation(WORKER, OTHER_PRIMARY, "observer", false, false);
var loop = loop(1, 300); // wide backoff so both calls land before the first tick fires
loop.onReplyQueued(WORKER);
loop.onTicketTerminal("task-1", WORKER, false);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS));
String nudge = rec.sentParams().getFirst().getValue().toString();
assertFalse(nudge.contains("fleet_poll(target=" + WORKER + ")"),
"the reply half is forbidden for this delegator and must never be named: " + nudge);
}
// --- nudge format --------------------------------------------------------------------------
@Test
@@ -1150,6 +1229,21 @@ class ReplyPushLoopTest {
"the nudge lands on the first tick after the box empties");
}
// --- the gate's on/off switch ----------------------------------------------------------------
@Test
void aDisabledGateNudgesALeadWithUnsubmittedTextInItsPromptBox() {
var herdr = new PaneTextHerdrClient(FakeHerdr.DRAFTED_PROMPT_CARET);
agents = new AgentControl(herdr);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(5, 100_000, false);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
"a disabled gate delivers even though the box holds a draft");
assertEquals(0, herdr.readCount(), "a disabled gate never reads the pane");
}
// --- helpers -------------------------------------------------------------------------------
private ReplyPushLoop loop() {
@@ -1164,6 +1258,11 @@ class ReplyPushLoopTest {
return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs, metrics);
}
private ReplyPushLoop loop(int maxReminders, long backoffMs, boolean promptBoxGateEnabled) {
return new ReplyPushLoop(registry, agents, inbox, scheduler, maxReminders, backoffMs, null,
promptBoxGateEnabled);
}
private static AgentControl agentWithStatus(String status) {
return new AgentControl(new FakeHerdrClient(status));
}
@@ -1483,6 +1582,7 @@ class ReplyPushLoopTest {
private static final class PaneTextHerdrClient implements HerdrClient {
private final List<Map.Entry<String, Object>> prompts =
Collections.synchronizedList(new ArrayList<>());
private final AtomicInteger reads = new AtomicInteger();
private volatile String paneText;
private volatile boolean failReads = false;
volatile CountDownLatch sendLatch = new CountDownLatch(1);
@@ -1504,6 +1604,10 @@ class ReplyPushLoopTest {
return prompts.size();
}
int readCount() {
return reads.get();
}
@Override
public JsonNode call(String method, Object params) {
if ("agent.get".equals(method)) {
@@ -1513,6 +1617,7 @@ class ReplyPushLoopTest {
.put("agent_status", "idle"));
}
if ("agent.read".equals(method)) {
reads.incrementAndGet();
if (failReads) {
throw new HerdrException("herdr socket read timed out");
}
@@ -285,6 +285,67 @@ class FleetAppAuthTest {
}
}
/**
* An observer may read a ticket its own {@code fleet_send(wait:false)} created, over the REST
* route and not only the unit-level classifier — and still not a ticket a different caller
* created, even another observer pane's.
*/
@Test
void restPollLetsAnObserverReadItsOwnTicketButNotAnothersOverRest() throws Exception {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
Injector injector = new Injector(agents);
MessageService messages = new MessageService(agents, injector, new Rendezvous());
Javalin observerApp = startOnSharedServiceAsObserver(messages, herdr, 9001L); // -> term_shell
Javalin otherObserverApp = startOnSharedServiceAsObserver(messages, herdr, FakeHerdr.WORKER_PID); // -> term_a
try {
Principal observer = Principal.observer("term_shell", 9001);
String ticket = messages.sendAsync("term_a", "long task", null, observer);
HttpResponse<String> own = send(observerApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(200, own.statusCode());
assertFalse(own.body().contains("forbidden"),
"the observer that created the ticket must read it: " + own.body());
// Unlike a worker's or the primary's mismatched ticket (refused by MessageService's own
// ownership check, inside a 200 response), a different observer is refused at the
// authorization gate itself, since its grant is conditional on the classifier -- so
// this is a 403, never reaching poll().
HttpResponse<String> other = send(otherObserverApp.port(), "GET", "/tasks/" + ticket, null, null);
assertEquals(403, other.statusCode(), other.body());
assertFalse(other.body().contains("\"reply\""),
"a refusal must never carry reply text: " + other.body());
} finally {
observerApp.stop();
otherObserverApp.stop();
}
}
/**
* As {@link #startOnSharedService(MessageService, FakeHerdr, long)}, but resolves every
* connecting pane to the unconfigured-pane {@link Role#OBSERVER} floor instead of a spawned
* worker — no lead, collaborator, or member map claims it.
*/
private Javalin startOnSharedServiceAsObserver(MessageService messages, FakeHerdr herdr, long pid) {
FleetConfig.Profile wcfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
AgentControl agents = new AgentControl(herdr);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(
agents, new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
Map.of(wcfg.profile(), wcfg), wcfg.profile(),
k -> "FLEETD_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null,
Map::of, new MemberRegistry(null), t -> null, Map::of);
Metrics appMetrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
return new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
callers, appMetrics).build().start("127.0.0.1", 0);
}
/**
* {@code POST /sessions/{id}/message} with {@code wait:false} must record the creating
* caller's own terminal on the ticket it returns, so that caller can still poll its own
@@ -694,6 +755,27 @@ class FleetAppAuthTest {
assertEquals(403, toSpawnedMembersTerminal.statusCode(), toSpawnedMembersTerminal.body());
}
/**
* The REST route must give the same answer as MCP for an observer: pid 9001 resolves to
* "term_shell", a herdr pane recognised as no configured role, so the real
* {@link CallerResolver#observerSendTarget()} classifier reaches the known lead and refuses
* the known collaborator -- over the route, not just the unit-level classifier, so a grant
* covering only MCP cannot leave this one behind.
*/
@Test
void anObserverMaySendToAKnownLeadButNotToAKnownCollaboratorOverRest() throws Exception {
int port = startWithRealClassifier(9001L,
Map.of("term_lead_known", "lead-x"), Map.of("term_collab_known", "ops2"));
HttpResponse<String> toLead = send(port, "POST", "/sessions/term_lead_known/message",
"{\"content\":\"hi\",\"wait\":false}", null);
assertEquals(202, toLead.statusCode(), toLead.body());
HttpResponse<String> toCollaborator = send(port, "POST", "/sessions/term_collab_known/message",
"{\"content\":\"hi\",\"wait\":false}", null);
assertEquals(403, toCollaborator.statusCode(), toCollaborator.body());
}
/**
* As {@link #start}, but with explicit lead/collaborator maps and no spawned-member roster, so
* a test can wire the real {@link CallerResolver#knownLeadOrCollaborator()} classifier instead
@@ -724,6 +806,59 @@ class FleetAppAuthTest {
return app.port();
}
/**
* As {@link #startWithRealClassifier}, but with a {@link MemberRegistry} bound to an architect
* slot at {@code terminal} (the defense-in-depth path, no live roster) — lets a test drive a
* real architect {@link Principal} over HTTP.
*/
private int startWithArchitect(long pid, String terminal) {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile wcfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
AgentControl agents = new AgentControl(herdr);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(
agents, new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
Map.of(wcfg.profile(), wcfg), wcfg.profile(),
k -> "FLEETD_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
Injector injector = new Injector(agents);
MessageService messages = new MessageService(agents, injector, new Rendezvous());
MemberRegistry members = new MemberRegistry(new FleetConfig.Fleet(
Map.of(), Map.of("lead-designer", new FleetConfig.Slot("sonnet")), Map.of(), Map.of(), null));
assertTrue(members.bind("architect:lead-designer", terminal));
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> pid);
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, false, null, Map::of, members);
metrics = FleetMetrics.create(sessions, new dev.ltms.fleet.msg.InMemoryReplyInbox());
app = new FleetApp(herdr, workers, sessions, messages, sessions.asPresence(), null,
callers, metrics).build().start("127.0.0.1", 0);
return app.port();
}
/**
* fleetd #801 defect 1: {@code SEND} is granted to an architect but {@code DRAIN} is not, so a
* TIMED_OUT/BUSY receipt to an architect must not name {@code GET /sessions/{id}/replies} —
* that route refuses the same caller — and must say the reply is unrecoverable instead.
*/
@Test
void anArchitectsTimedOutSendDoesNotNameTheRepliesRouteItCannotDrain() throws Exception {
int port = startWithArchitect(FakeHerdr.WORKER_PID, "term_a"); // FakeHerdr resolves this pid to term_a
HttpResponse<String> res = send(port, "POST", "/sessions/term_b/message",
"{\"content\":\"hi\",\"timeoutMs\":150}", null);
assertEquals(202, res.statusCode(), "an architect is granted SEND: " + res.body());
JsonNode body = new ObjectMapper().readTree(res.body());
String detail = body.get("detail").asText();
assertFalse(detail.contains("/replies"), "got: " + detail);
assertTrue(detail.contains("cannot be recovered"), "got: " + detail);
assertEquals(403, send(port, "GET", "/sessions/term_b/replies", null, null).statusCode(),
"the architect this receipt was sent to is indeed refused DRAIN");
}
// --- loopback-trust: the caller is the primary -------------------------------------------
@Test
@@ -663,8 +663,10 @@ class FleetAppTest {
HttpResponse<String> res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":150}");
assertEquals(202, res.statusCode());
assertEquals("queued", mapper.readTree(res.body()).get("status").asText());
JsonNode body = mapper.readTree(res.body());
assertEquals("queued", body.get("status").asText());
assertFalse(herdr.called("agent.prompt"), "no injection while the worker is mid-turn");
assertDetailNamesTheInboxRouteNotARetry(body);
}
@Test
@@ -674,8 +676,21 @@ class FleetAppTest {
HttpResponse<String> res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":250}");
assertEquals(202, res.statusCode());
assertEquals("working", mapper.readTree(res.body()).get("status").asText());
JsonNode body = mapper.readTree(res.body());
assertEquals("working", body.get("status").asText());
assertTrue(herdr.called("agent.prompt"), "message was injected");
assertDetailNamesTheInboxRouteNotARetry(body);
}
/**
* fleetd #801: the TIMED_OUT_WORKING / TIMED_OUT_QUEUED / BUSY receipt is never tracked by a
* ticket, so it must name the one route that recovers a late reply — draining the session's
* inbox — and must not invite a bare resend, which can duplicate the delivery.
*/
private void assertDetailNamesTheInboxRouteNotARetry(JsonNode body) {
String detail = body.get("detail").asText();
assertTrue(detail.contains("/sessions/term_a/replies"), "got: " + detail);
assertFalse(detail.contains("retry"), "got: " + detail);
}
/**
+2 -2
View File
@@ -1,7 +1,7 @@
{
"name": "fleet",
"description": "Make a project fleet-ready: mount the fleetd MCP gateway and set up standard Claude Code settings so this session can orchestrate a fleet of delegated workers. Lead-side only — member skills and agents travel in the worktree. Ships no credentials.",
"version": "0.2.0",
"description": "Set up standard Claude Code settings so this session can orchestrate a fleet of delegated workers, and run the fleet mod for cross-session messaging. Lead-side only — member skills and agents travel in the worktree, and mounting the fleetd MCP gateway is now the instance's or the project's job, not this plugin's. Ships no credentials.",
"version": "0.3.0",
"author": {
"name": "LTMS"
},
-8
View File
@@ -1,8 +0,0 @@
{
"mcpServers": {
"fleet": {
"type": "http",
"url": "${FLEETD_MCP_URL}"
}
}
}
+35 -19
View File
@@ -1,7 +1,7 @@
# fleet (Claude Code plugin)
Makes a project **fleet-ready**: mounts the `fleetd` MCP gateway and applies standard Claude Code
settings, so the session can orchestrate a fleet of delegated workers.
Makes a project **fleet-ready**: applies standard Claude Code settings and runs the fleet mod, so
the session can orchestrate a fleet of delegated workers.
**This plugin ships no credentials.** Every secret is referenced by environment-variable *name*;
the values stay with the user. Nothing the plugin writes is unsafe to commit.
@@ -22,9 +22,10 @@ Member-facing assets travel in the worktree, not in this plugin. See fleetd #362
## What it is not
The plugin is the **client-side setup**, not the bridge. `fleetd` is a separate daemon and `herdr`
is a separate PTY multiplexer, each with its own lifecycle and install. The plugin mounts an
already-running daemon and tells you what is missing when one isn't there — it deliberately does
not try to install system services on your behalf.
is a separate PTY multiplexer, each with its own lifecycle and install, and the plugin does not try
to install either on your behalf. It also does not mount the daemon for you — mounting is the
instance's or the project's own `.mcp.json`, and `/fleet:setup` is the one thing in this plugin that
still helps with that (it writes the project-level entry).
## Install
@@ -33,11 +34,20 @@ not try to install system services on your behalf.
/plugin install fleet@fleetd
```
Export the gateway URL — the plugin mounts `${FLEETD_MCP_URL}`, not a hardcoded address, so one
plugin serves hosts that run the daemon on different ports:
Mount the daemon yourself — this plugin carries no mount of its own. Either add the entry below to
your Claude Code instance's own `.claude.json`, so every project you open there gets it, or run
`/fleet:setup` in the project you want to onboard, which writes the same entry into that project's
`.mcp.json`:
```shell
export FLEETD_MCP_URL=http://127.0.0.1:8765/mcp
```json
{
"mcpServers": {
"fleet": {
"type": "http",
"url": "http://127.0.0.1:8765/mcp"
}
}
}
```
Then, in the project you want to onboard:
@@ -50,18 +60,24 @@ Then, in the project you want to onboard:
| Component | Effect |
|---|---|
| `.mcp.json` | mounts `fleet` at `${FLEETD_MCP_URL}` for any session with the plugin enabled |
| `skills/setup` | `/fleet:setup` — preflight, project settings, credential guidance, and verification |
| `skills/setup` | `/fleet:setup` — preflight, project settings, credential guidance, and verification. Also the only thing in this plugin that helps mount `fleet`: it writes the project `.mcp.json` entry shown above. |
| `hooks/register.js` (the fleet mod) | Cross-account session messaging while the plugin is enabled: `/fleet-peers`, `/fleet-mail`, `/fleet-whoami`, and a background poll that delivers mail fleetd queued for this pane. A spawned worker or architect skips that poll, because it already gets its brief pasted into its pane. |
The server is named **`fleet`** on purpose: that is `PeerLauncher.MCP_MOUNT_NAME` in the daemon and
the name a spawned member's own mount carries. Version 0.1.0 named it `fleetd`, which produced two
mounts of one daemon for anyone who also had a project-level `.mcp.json`. Upgrading from 0.1.0 is a
**breaking change** — a project that pre-allowed `mcp__fleetd__fleet_whoami` in
`.claude/settings.json` must be updated to `mcp__fleet__*`.
Whichever file mounts the daemon, name the server **`fleet`**. That is `PeerLauncher.MCP_MOUNT_NAME`
in the daemon, the name a spawned member's own mount carries, and the name the `mcp__fleet__*`
role heuristic in `CLAUDE.md` keys on.
Because the plugin carries its own `.mcp.json`, an installed plugin needs no project-level MCP
file at all. The setup skill writes one only when you want the mount to work *without* the plugin —
for teammates who haven't installed it, or for CI.
## Upgrading from 0.2.0 — breaking
The plugin no longer mounts the daemon. It used to carry its own `.mcp.json`, pointed at
`${FLEETD_MCP_URL}`, and that file is gone. The mod still reads `FLEETD_MCP_URL`, but only as an
optional override of the address it calls (default `http://127.0.0.1:8765/mcp`). Mount `fleet`
yourself: add the entry under **Install** above to your instance's `.claude.json` or to the
project's own `.mcp.json`, by hand or with `/fleet:setup`.
Version 0.1.0 named the mounted server `fleetd`, which produced two mounts of one daemon for
anyone who also had a project-level `.mcp.json`. A project that pre-allowed
`mcp__fleetd__fleet_whoami` in `.claude/settings.json` must be updated to `mcp__fleet__*`.
## Verifying a setup
+30 -10
View File
@@ -67,23 +67,28 @@ async function livePeers($, now) {
return rows
}
const FLEETD_MCP = 'http://127.0.0.1:8765/mcp'
const FLEETD_MCP_DEFAULT = 'http://127.0.0.1:8765/mcp'
const MCP_HEADERS = { 'Content-Type': 'application/json', Accept: 'application/json, text/event-stream' }
// The status fleetd answers, with "Session not found", for an Mcp-Session-Id it no longer holds.
const MCP_SESSION_GONE = 404
// The MCP session every call below shares. fleetd keeps a server-side session per initialize and
// drops it only on a DELETE, so one initialize per call would leave a session behind every time.
// The MCP session every call below shares, and the URL it was opened against. fleetd keeps a
// server-side session per initialize and drops it only on a DELETE, so one initialize per call
// would leave a session behind every time.
let mcpSessionId = null
let fleetdUrl = FLEETD_MCP_DEFAULT
/**
* Open an MCP session on the local fleetd and hold it for later calls.
*
* Cleared first, so a failure here leaves no dead id behind for the next call to reuse.
* Cleared first, so a failure here leaves no dead id behind for the next call to reuse. Reads
* FLEETD_MCP_URL fresh on every open, so a session opened after the daemon moves uses the new
* address.
*/
async function openFleetSession($) {
mcpSessionId = null
const init = await $.http.fetch(FLEETD_MCP, {
fleetdUrl = (await $.env.get('FLEETD_MCP_URL')) || FLEETD_MCP_DEFAULT
const init = await $.http.fetch(fleetdUrl, {
method: 'POST',
headers: MCP_HEADERS,
body: JSON.stringify({
@@ -93,7 +98,7 @@ async function openFleetSession($) {
})
if (!init.ok) throw new Error('fleetd initialize failed with status ' + init.status)
const opened = init.headers['mcp-session-id']
await $.http.fetch(FLEETD_MCP, {
await $.http.fetch(fleetdUrl, {
method: 'POST',
headers: { ...MCP_HEADERS, 'Mcp-Session-Id': opened },
body: JSON.stringify({ jsonrpc: '2.0', method: 'notifications/initialized' }),
@@ -103,7 +108,7 @@ async function openFleetSession($) {
/** Send one tools/call on the session this mod holds, and return the raw HTTP answer. */
function sendFleetToolCall($, tool, args) {
return $.http.fetch(FLEETD_MCP, {
return $.http.fetch(fleetdUrl, {
method: 'POST',
headers: { ...MCP_HEADERS, 'Mcp-Session-Id': mcpSessionId },
body: JSON.stringify({ jsonrpc: '2.0', id: 2, method: 'tools/call', params: { name: tool, arguments: args } }),
@@ -170,17 +175,32 @@ export function register(on) {
}
})
// A spawned worker or architect already gets its brief pasted into its pane, so this
// poll would be a second, redundant delivery path for it. Every other role collects
// its own mail through this poll.
let role = null
// Collect what the fleet queued for this pane and hand each message to
// Claude. Every call also renews fleetd's record that this pane collects its
// own mail, so an empty answer still has to be asked for.
$.clock.every(FLEETD_INBOX_POLL_MS, async () => {
if (role === null) {
try {
role = JSON.parse(await fleetTool($, 'fleet_whoami', {})).role
} catch {
// A daemon that is down, or a pane fleetd cannot place, is the ordinary
// case on a host with no fleet running. Retry on the next tick.
return
}
}
if (role === 'worker' || role === 'architect') return
let collected
try {
collected = JSON.parse(await fleetTool($, 'fleet_inbox', {}))
} catch {
// A daemon that is down, or a pane fleetd cannot place, is the ordinary
// case on a host with no fleet running. The timer survives a throw, so
// this only keeps every tick from writing an error to the debug log.
// The timer survives a throw, so this only keeps every tick from
// writing an error to the debug log.
return
}
for (const text of collected.messages || []) {
+6 -12
View File
@@ -36,16 +36,11 @@ a time.
command -v herdr && herdr --version 2>&1 | head -1 || echo "MISSING: herdr"
command -v ccs && ccs version 2>&1 | head -1 || echo "MISSING: ccs (needed for worker profiles)"
command -v codex && codex --version 2>&1 | head -1 || echo "absent: codex (optional)"
curl -s -m 5 "${FLEETD_MCP_URL%/mcp}/healthz" 2>/dev/null \
|| curl -s -m 5 http://127.0.0.1:8765/healthz \
|| echo "MISSING: fleetd daemon is not reachable"
[ -n "$FLEETD_MCP_URL" ] && echo "FLEETD_MCP_URL is set" || echo "MISSING: FLEETD_MCP_URL"
curl -s -m 5 http://127.0.0.1:8765/healthz || echo "MISSING: fleetd daemon is not reachable"
```
**`FLEETD_MCP_URL` is required.** The plugin's own `.mcp.json` mounts `${FLEETD_MCP_URL}` rather
than a hardcoded address, so one plugin can serve hosts that run the daemon on different ports. If
it is unset the mount does not resolve. The usual value is `http://127.0.0.1:8765/mcp`; tell the
user to export it, do not write it into a file for them.
The usual address is `http://127.0.0.1:8765`. If this daemon runs elsewhere, use that address
instead wherever this skill writes `http://127.0.0.1:8765/mcp` below.
A healthy daemon answers with its status **and the herdr protocol it negotiated**:
@@ -95,10 +90,9 @@ If `.mcp.json` already exists, add only the `fleet` key and leave every other se
If a `fleet` entry is already there with a different URL, **ask** rather than assuming yours is
right — a non-default port usually means a deliberate second daemon.
> **If this plugin is installed, you can skip this step entirely.** The plugin ships its own
> `.mcp.json`, so `fleetd` is already mounted for any session with the plugin enabled. Write the
> project-level file only when the user wants the mount to work *without* the plugin — for
> teammates who have not installed it, or for CI.
This write is the only way this plugin helps mount `fleet` — the plugin carries no mount of its
own. A session that wants the mount without running this skill can instead add the same entry to
its own Claude Code instance's `.claude.json`.
**Before writing it, settle whether `.mcp.json` is committed here:**
+151 -10
View File
@@ -32,6 +32,7 @@ function stubEngine(on: any): Map<string, any> {
on('session.cwd', () => ({ value: '/Users/x/claude-bridge' }))
on('ui.log', () => ({ value: undefined }))
on('prompt.submit', () => ({ value: undefined }))
on('env.get', () => ({ value: undefined }))
// Claude Code's own delivery, which the mod's receive hook must reach.
on('session.receive', (_$: any, e: any) => e)
return store
@@ -181,6 +182,7 @@ function stubSessionStart(on: any, submitted: string[]): Map<string, any> {
on('session.cwd', () => ({ value: '/Users/x/claude-bridge' }))
on('command.register', () => ({ value: undefined }))
on('ui.log', () => ({ value: undefined }))
on('env.get', () => ({ value: undefined }))
// The engine skips a prompt.submit hook that answers anything but { text } or { drop }, and
// the mod's callback then throws, so this must hand the text straight back.
on('prompt.submit', (_$: any, e: any) => {
@@ -194,10 +196,12 @@ function stubSessionStart(on: any, submitted: string[]): Map<string, any> {
/**
* Answer fleetd's MCP requests, with the tool result read fresh on every call.
*
* `fetches` collects every method sent, and `sessions` the Mcp-Session-Id of each tools/call.
* Each initialize hands out the next id, so a reused session and a reopened one differ.
* `sessionGone` makes a tools/call answer the way fleetd answers for a session id it no longer
* holds.
* `fetches` collects every method sent, with a tools/call entry naming its tool (such as
* `tools/call:fleet_inbox`), and `sessions` the Mcp-Session-Id of each tools/call. Each
* initialize hands out the next id, so a reused session and a reopened one differ.
* `sessionGone` makes a tools/call answer the way fleetd answers for a session id it no
* longer holds. `whoamiText` answers a `fleet_whoami` call apart from `toolText`, which
* answers every other tool.
*/
function stubFleetdDynamic(
on: any,
@@ -205,24 +209,28 @@ function stubFleetdDynamic(
fetches: string[],
sessions: string[] = [],
sessionGone: () => boolean = () => false,
whoamiText: () => string = () => '{"role":"primary"}',
) {
let opened = 0
on('http.fetch', (_$: any, e: any) => {
const body = JSON.parse(e.init.body)
fetches.push(body.method)
if (body.method === 'initialize') {
fetches.push(body.method)
opened += 1
return { value: { ok: true, status: 200, headers: { 'mcp-session-id': 'sid-' + opened }, text: '{}' } }
}
if (body.method === 'tools/call') {
fetches.push(body.method + ':' + body.params.name)
sessions.push(e.init.headers['Mcp-Session-Id'])
if (sessionGone()) {
const gone = '{"jsonRpcError":{"code":-32603,"message":"Session not found"}}'
return { value: { ok: false, status: 404, headers: {}, text: gone } }
}
const result = { jsonrpc: '2.0', id: 2, result: { content: [{ type: 'text', text: toolText() }] } }
const text = body.params.name === 'fleet_whoami' ? whoamiText() : toolText()
const result = { jsonrpc: '2.0', id: 2, result: { content: [{ type: 'text', text }] } }
return { value: { ok: true, status: 200, headers: {}, text: 'data: ' + JSON.stringify(result) + '\n' } }
}
fetches.push(body.method)
return { value: { ok: true, status: 202, headers: {}, text: '' } }
})
}
@@ -311,11 +319,11 @@ test('a second inbox poll reuses the first MCP session', async ($, on) => {
await clock.advance(3_000)
await clock.advance(3_000)
// fleetd holds a session per initialize and drops it only on a DELETE, so one initialize per
// poll would leave one behind every three seconds.
expect(sessions.length).toBe(2) // control: both polls really reached fleetd
// The first poll also reads the role once, so it makes two tools/call (whoami, then
// inbox); the second poll already knows the role and makes only one (inbox).
expect(sessions.length).toBe(3) // control: all three calls really reached fleetd
expect(initializes(fetches)).toBe(1)
expect(sessions[1]).toBe(sessions[0])
expect(sessions.every((s) => s === sessions[0])).toBe(true)
})
test('a session fleetd no longer holds is opened again and the call retried', async ($, on) => {
@@ -347,3 +355,136 @@ test('a session fleetd no longer holds is opened again and the call retried', as
expect(sessions[sessions.length - 1]).toBe('sid-2')
expect(submitted.length).toBe(2)
})
/** How many of `fetches` were a tools/call for the named tool. */
function toolCalls(fetches: string[], tool: string): number {
return fetches.filter((m) => m === 'tools/call:' + tool).length
}
test('a worker role stops the poll from calling fleet_inbox', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, [], () => false, () => '{"role":"worker"}')
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_whoami')).toBe(1)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(0)
expect(submitted.length).toBe(0)
})
test('an architect role stops the poll from calling fleet_inbox', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, [], () => false, () => '{"role":"architect"}')
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_whoami')).toBe(1)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(0)
expect(submitted.length).toBe(0)
})
test('a primary role keeps the poll calling fleet_inbox', async ($, on) => {
// The positive control for the two tests above: the same stubs, a role neither
// gates, so a bug that silenced every role would still pass them.
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, [], () => false, () => '{"role":"primary"}')
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(1)
expect(submitted.length).toBe(1)
})
test('an observer role keeps the poll calling fleet_inbox', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, [], () => false, () => '{"role":"observer"}')
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(1)
expect(submitted.length).toBe(1)
})
test('a whoami that fails on the first tick is retried and the poll resumes', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 1, messages: ['do the task'] }
let whoamiFails = true
on('http.fetch', (_$: any, e: any) => {
const body = JSON.parse(e.init.body)
if (body.method === 'tools/call' && body.params.name === 'fleet_whoami' && whoamiFails) {
fetches.push(body.method + ':' + body.params.name)
return { value: { ok: false, status: 503, headers: {}, text: '' } }
}
if (body.method === 'initialize') {
fetches.push(body.method)
return { value: { ok: true, status: 200, headers: { 'mcp-session-id': 'sid-1' }, text: '{}' } }
}
if (body.method === 'tools/call') {
fetches.push(body.method + ':' + body.params.name)
const text = body.params.name === 'fleet_whoami' ? '{"role":"observer"}' : JSON.stringify(inbox)
const result = { jsonrpc: '2.0', id: 2, result: { content: [{ type: 'text', text }] } }
return { value: { ok: true, status: 200, headers: {}, text: 'data: ' + JSON.stringify(result) + '\n' } }
}
fetches.push(body.method)
return { value: { ok: true, status: 202, headers: {}, text: '' } }
})
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(0) // control: a failed whoami calls no fleet_inbox
expect(submitted.length).toBe(0)
whoamiFails = false
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(1)
expect(submitted.length).toBe(1)
})
test('whoami is called once, not on every tick, once the role is known', async ($, on) => {
const clock = mock.clock(on)
const submitted: string[] = []
stubSessionStart(on, submitted)
const fetches: string[] = []
const inbox = { sessionId: 'term_self', count: 0, messages: [] as string[] }
stubFleetdDynamic(on, () => JSON.stringify(inbox), fetches, [], () => false, () => '{"role":"primary"}')
await $.session.start({ surface: 'terminal', isInteractive: true, cwd: '/Users/x/claude-bridge' })
await clock.advance(3_000)
await clock.advance(3_000)
await clock.advance(3_000)
expect(toolCalls(fetches, 'fleet_whoami')).toBe(1)
expect(toolCalls(fetches, 'fleet_inbox')).toBe(3)
})