Compare commits

...

12 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 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
23 changed files with 791 additions and 74 deletions
+40 -23
View File
@@ -33,9 +33,11 @@ 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`/`INBOX` on its own pane, and 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 never a
ticket. It finds such a target in `fleet_list`'s `panes` array, which for an observer is filtered to
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.
@@ -71,20 +73,29 @@ and the sender silently receives nothing. Fail toward the recoverable error.
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 call outside your role is refused, not queued.
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. A direct `fleet_send` to a lead's
pane does **not** wait for the box yet (fleetd #793); a lead that runs the fleet mod collects
that mail instead of having it pasted, so it is not exposed. 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
@@ -184,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.
@@ -291,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.
+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 {
@@ -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)
@@ -1937,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) {
@@ -2734,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);
@@ -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
@@ -538,7 +541,8 @@ public final class FleetMcp {
// 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,
@@ -1032,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");
}
@@ -1047,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());
}
@@ -1058,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);
}
/**
@@ -1092,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
@@ -1116,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.
@@ -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
@@ -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;
@@ -742,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;
}
@@ -756,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);
}
@@ -766,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",
@@ -816,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")));
}
}
@@ -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");
}
}
@@ -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));
}
@@ -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");
}
/**
@@ -1229,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() {
@@ -1243,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));
}
@@ -1562,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);
@@ -1583,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)) {
@@ -1592,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");
}
@@ -806,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);
}
/**