Compare commits

...

13 Commits

Author SHA1 Message Date
Dai Ha 275ac0d251 fleetd #562: surface loop health
CI / shell-tests (pull_request) Successful in 12s
CI / contract (pull_request) Successful in 1m31s
CI / build (pull_request) Successful in 2m11s
2026-09-12 19:26:37 +07:00
Dai Ha 204da67d66 Merge #576: fleetd #575 — one finally covers answer()'s STALE_TURN exit
CI / shell-tests (push) Successful in 16s
CI / build (push) Successful in 1m42s
CI / contract (push) Successful in 1m57s
Third instance of the #572 shape on this file: one invariant kept at N sites,
asserted at fewer than N. Here answer()'s inner try opened AFTER the Task
lookup/registration and the STALE_TURN early return, so that return was covered
only by a hand-rolled copy of the finally's cleanup pair. The fix widens the try
upward and deletes the copy, so every exit runs the one finally exactly once.
rendezvous.open(workerSession) correctly stays outside it — nothing to clean if
it never opened.

This fixed no live leak, and the code comment says so: neither
rendezvous.answerAsk nor clearAsyncQuestion(turnId, false) can throw, so nothing
ever left through the old gap uncovered. It is a structure fix, the ticket's own
fallback case.

Verified here, not taken from the worker's report:
- baseline on the branch: 1766 tests, 0 failures, from Maven and from an
  independent sum over target/surefire-reports/*.txt.
- my own mutation, located fresh: move the inner `try {` back down below the
  STALE_TURN return — the exact pre-fix structure, minus the hand-rolled pair.
  Result: 1766 run, exactly 1 failure, and it is the new test —
  MessageServiceTest.answerLosingTheRaceToAnAlreadyAnsweredAskStillReturnsStaleTurnAndCleansUpOnce:545
  "the forward waiter this answer() call opened must be closed after a STALE_TURN
  return". One failure, not a crowd: the new test is the only thing holding this
  path.
- the worker's own mutation removed the whole finally and took 8 other tests with
  it. That proves the finally runs; it does not prove the STALE_TURN path reaches
  it. Mine does.

The new test hook answerAskLapseRaceHookForTest follows the file's existing
askTimeoutRaceHookForTest convention.

Closes #575.
2026-09-12 19:05:00 +07:00
Dai Ha b091c51eee fleetd #575: widen answer()'s try so one finally covers its STALE_TURN exit
CI / shell-tests (pull_request) Successful in 5s
CI / contract (pull_request) Successful in 58s
CI / build (pull_request) Successful in 2m30s
The waiter cleanup pair (asyncTasksByWaiter.remove + rendezvous.close) was
duplicated: two sites sit in a finally, the third was hand-rolled inline
before answer()'s early STALE_TURN return, structurally outside any finally.

Both rendezvous.answerAsk and clearAsyncQuestion(turnId, false) are total
(cannot throw), so the gap never leaked in practice. But the duplicate was
untested: mutating it away left all 1765 tests green, while the two
finally-protected sites are each killed by 8-36 tests. Same shape as #572.

Fix: widen the try to wrap the Task registration and the STALE_TURN check,
so the single finally covers every exit and the hand-rolled copy is gone.
Added a race hook + regression test that deterministically reproduces the
'ask lapsed between the lookup and the unblock' case and proves the fix
still returns STALE_TURN and cleans up exactly once.
2026-09-12 18:56:02 +07:00
ltms 84d631b030 Merge #572: answer()'s session-lock release is pinned on all four exits
CI / shell-tests (push) Successful in 9s
CI / contract (push) Successful in 1m27s
CI / build (push) Successful in 1m44s
fleetd #572. MessageService releases its per-session lock in a finally at two sites. Removing the
send() one failed 22 tests and errored 1. Removing the answer() one left the whole suite green.
Skip that unlock and the thread holds the session lock forever, so every later send or answer to
that session blocks permanently — no exception, no log line.

Tests only. Against current main the diff is one file, MessageServiceTest.java, +178 lines.

Four new tests, one per exit of answer(): normal REPLIED, TIMED_OUT_WORKING, the ExecutionException
rethrow, and the InterruptedException rethrow. Each proves REACQUISITION rather than the return
value — a bounded follow-up send on the SAME session must not come back BUSY, and send reports BUSY
only when tryLock itself timed out, so a non-BUSY probe is specifically evidence the lock was free.

Lead verification, re-running rather than accepting the worker's numbers, on the branch merged with
main at 3f8c38f:
  baseline   -> exit 0, 1765 tests from Maven and an independent sum over 131 reports (1761 + 4)
  mutation L -> BUILD FAILURE, exactly 4 failures, all four new tests by name

THE TRAP THE WORKER CAUGHT ITSELF, which is the most valuable thing in this PR. Their first draft of
the TIMED_OUT_WORKING test called answer() inline on the test thread and then probed with send() on
that same thread. It FALSE-PASSED under the mutation: lock is a ReentrantLock, so the same thread
reenters for free whether or not unlock() ran. A same-thread probe proves reentrancy, not release.
They found it only because they actually ran the mutation instead of trusting a green baseline, moved
answer() onto a background thread, and re-proved the kill. The pitfall is now documented in the test's
own javadoc, with the measurement that exposed it.

NOTE FOR ANYONE RE-RUNNING THE TICKET'S RECIPE: the ticket records `sed '1218s|...'`, measured before
#551 merged. #551 edited MessageService.java, and the line is now 1231. Locate it fresh — a
line-anchored sed against a stale number mutates the wrong line and reports a meaningless green. The
anchor count is the check that catches this: `grep -Fxc '            lock.unlock();'` must read 2
pristine and 1 after.

Shape reported by the worker, not investigated and not fixed: the same two-line cleanup
`asyncTasksByWaiter.remove(reply); rendezvous.close(...)` appears at three sites — send()'s finally,
answer()'s inner finally, and answer()'s inline STALE_TURN early return, which is NOT inside a
finally. Same maintained-at-N-sites shape as this finding. Filed separately.
2026-09-12 13:30:22 +02:00
ltms 3f8c38fc54 Merge #567: pin LeadMailbox.inspect's probe-channel close
CI / shell-tests (push) Successful in 8s
CI / contract (push) Successful in 1m1s
CI / build (push) Successful in 1m49s
fleetd #567. LeadMailbox.inspect opens a probe channel and closes it in a finally. The production
code was already correct; nothing asserted it, so a future refactor could drop the close and leak an
AMQP channel per inspect() call with the suite green.

Test only. LeadMailbox.java is untouched — sha256 a2cd99be77345b7e... before and after.

The test asserts the CONSEQUENCE rather than the return value: it caps the connection at three
channels (LeadMailbox uses two, consume and publish), runs a successful inspect, then requires a
replacement channel. If the probe is left open, the broker has no channel number left and
createChannel() returns null. A test that only checked inspect()'s MailboxState would pass under the
mutation, which is the whole reason this hole existed.

Lead verification, re-running rather than accepting the worker's numbers, on the branch merged with
main at ed2fd66:
  mvn -o clean install -Pcontract  ->  exit 0, 1792 tests from Maven and from an independent sum
                                       over 137 surefire reports.
  1792 = 1761 (main) + 30 (contract-only) + 1 (new), which also confirms the default-profile count
  is untouched: the new test is in an @Tag("contract") class.

My own mutation, located fresh rather than assuming the reported line number: delete
LeadMailbox.java:338 `probe.close();`, anchor by grep -Fxc 1 -> 0. Result under -Pcontract: exactly
1 failure, inspectClosesItsSuccessfulProbeChannel:226, "the replacement channel was null". Restored
to a2cd99be77345b7e..., git status --short empty.

THE PROFILE TRAP THIS TICKET EXISTS BECAUSE OF. An earlier sweep worker instrumented this line, saw
zero hits under the DEFAULT profile, and concluded "no test executes this line". That was false:
fleetd/pom.xml:264 sets excludedGroups=contract, so the covering tests were excluded from the run,
not absent. A surviving mutation has THREE causes — never executes, executes with nothing asserted,
or the covering tests were excluded from the profile — and only naming the profile tells them apart.
The true finding was covered-but-unasserted.

Known fragility, recorded rather than fixed: the test's channel cap of three assumes LeadMailbox
holds exactly two channels. If it ever holds more, this test fails loudly, which is fine. If it ever
holds fewer, a leak would no longer exhaust the cap and the test would go vacuous silently. Worth
re-checking if LeadMailbox's channel usage changes.

Not covered, and stated rather than faked: the defensive catch (RuntimeException) around the passive
declare, which needs a connection dying between createChannel() and the declare landing. The
method's own javadoc already admits that branch is unproven; the worker did not invent a test for it.
2026-09-12 13:26:30 +02:00
ltms ed2fd6646a Merge #561: the completion/session listener fan-out survives either half throwing, and both sites are pinned
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 50s
CI / build (push) Successful in 2m7s
fleetd #561. Fleetd composed two TurnListener halves as `completion.X(); sessions.X();`, so a throw
from the first half skipped the second. The anonymous class is now a package-private factory,
Fleetd.turnListener(completion, sessions), built on two helpers that always attempt both halves and
rethrow whatever escaped — a second failure attached with addSuppressed rather than dropped, so it
still reaches StatusPoller's catch (Throwable).

The wiring at Fleetd.java:498 calls that factory, so the seam under test is the real caller.

Lead verification, on a merged tree, re-running the checks rather than accepting the worker's:
exit 0, 1761 tests from Maven and from an independent sum over 131 surefire reports.

The interesting part is what the first round MISSED. Two helpers maintain one invariant — "the
second half always runs" — and the first round's five tests asserted it at only one site. Measured:

  bothMustRun                      reverted to the pre-fix bug -> 1 named failure   (pinned)
  bothMustRunKeepingSecondResult   the SAME bug                -> 1755/1755 GREEN   (unpinned)
  failure.addSuppressed(t) deleted                             -> 1755/1755 GREEN   (unpinned)

Both survivors are now killed by new tests, re-verified by the lead after the fix:
sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronouslyForPostAction and
bothFailuresEscapeWhenBothHalvesThrowDistinctExceptions, each failing alone under its own mutation.

The rule this cost us, worth carrying: COUNT ASSERTIONS PER SITE, NOT PER INVARIANT. The total being
non-zero is what hides a zero at one site, and extracting a shared helper makes it worse rather than
better — it does not reduce the number of sites, only how many are visible. Credit to the fleet01
lead, who predicted this shape before an instance was found.

onDelivered stays deliberately unguarded. Its comment now gives the real reason —
CompletionResolver.captureBaseline already catches RuntimeException around its scrape and fails open,
so that half does not realistically throw — instead of the previous reason, which was true but about
registration rather than about this pair. A correct conclusion resting on a wrong premise reads
exactly like a verified one.
2026-09-12 13:20:56 +02:00
Dai Ha 5441a2b321 fleetd #567: assert inspect closes probe channel
CI / shell-tests (pull_request) Successful in 9s
CI / contract (pull_request) Successful in 54s
CI / build (pull_request) Successful in 2m23s
2026-09-12 18:19:43 +07:00
ltms 384867dfa3 Merge #551: ATTEMPTED is its own cancellation answer, and the javadoc stops claiming a timed-out send never arrived
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 1m11s
CI / build (push) Successful in 1m44s
fleetd #551. The injector polls a queued entry off the queue and marks it ATTEMPTED BEFORE the
irreversible AgentControl#send call, not after. So a Throwable escaping that call can never leave
the entry QUEUED at the head (the #546 re-send hazard) and can never be recorded as a confident
NOT_DELIVERED for text that may already be in the pane.

Cancellation.ATTEMPTED is added as a third answer. NOT_DELIVERED stays reserved for confirmed
absence: the readiness grace expiring, drop(), or a herdr *_not_found error, which the rest of this
codebase already reads as definitely-absent rather than inconclusive.

Also fixes three javadoc/comment sites that claimed a timed-out send definitely did not arrive:
TIMED_OUT_QUEUED, hasQueuedDelivery, the queuedDeliveries field, and the comment in send()'s timeout
branch. Two claims were wrong, not merely stale: "on every route it will not arrive later" is false
on the ATTEMPTED route, and "may already hold a partial paste" understates it — agent.prompt pastes
AND SUBMITS in one call, so the target may hold a complete, running turn.

Verified by the lead on a merged tree: mvn -o clean install from fleetd/, exit 0, 1754 tests from
Maven and from an independent sum over 130 surefire reports. Mutation: folding ATTEMPTED back into
NOT_DELIVERED in cancellationOf gives 3 red, each naming the property
(anErrorFromSendRemovesTheMessageAndMarksItAttempted:740,
aHerdrExceptionFromSendStillSurfacesButNowReportsAttempted:781,
aHerdrExceptionAfterThePasteIsNeverRecordedAsConfidentlyNotDelivered:806).

The final round is comment-only, proven mechanically rather than by reading: stripping every comment
from MessageService.java before and after and collapsing whitespace gives byte-identical code.
2026-09-12 13:19:20 +02:00
Dai Ha f40c19ecf0 fleetd #551 shape sweep: fix stale ATTEMPTED-route javadoc/comments in MessageService
CI / shell-tests (pull_request) Successful in 6s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 2m26s
hasQueuedDelivery's javadoc, the queuedDeliveries field javadoc, and a code
comment in send()'s timeout path all still made two claims that ATTEMPTED
(fleetd #551) falsifies: a blanket "the message will not arrive later" across
every route, and "the terminal may already hold a partial paste" — but
agent.prompt pastes AND submits in one call, so the target may hold a
complete, already-submitted turn. Each now names the three Injector.Cancellation
routes (CANCELLED, NOT_DELIVERED, ATTEMPTED) and says plainly that only the
first two establish the message will not arrive later.

Comment/javadoc only. No behaviour change: Injector.java is untouched
(sha256 0c689b6cf36275c0da45497a74b2bd4f5a3d80c4dbda77d46670c66004e53b69) and
no test was added. mvn -o clean install: BUILD SUCCESS, Tests run: 1754,
Failures: 0, Errors: 0 (unchanged from before this commit).
2026-09-12 18:13:06 +07:00
Dai Ha 034e17bb32 fleetd #561 follow-up: pin the session half of bothMustRunKeepingSecondResult
CI / shell-tests (pull_request) Successful in 7s
CI / contract (pull_request) Successful in 1m14s
CI / build (pull_request) Successful in 1m41s
Two helpers maintain one invariant (the second callback half always runs,
even when the first throws): bothMustRun and bothMustRunKeepingSecondResult.
Only bothMustRun's "session half still runs" direction was asserted
(sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronously, via
onTurnComplete). bothMustRunKeepingSecondResult — the helper
onTurnCompleteWithPostAction uses — could be reverted to the pre-#561
broken shape and the suite stayed green.

Adds two tests to FleetdTurnListenerCompositionTest:
- sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronouslyForPostAction:
  mirrors the existing onTurnComplete case for onTurnCompleteWithPostAction/
  bothMustRunKeepingSecondResult.
- bothFailuresEscapeWhenBothHalvesThrowDistinctExceptions: proves a second,
  distinct failure from the session half is preserved via addSuppressed
  rather than silently dropped when both halves of bothMustRun throw.

Also rewords the onDelivered comment in Fleetd.turnListener: it previously
said this pair is safe because registration survives a throw via #556's
Injector wiring, which is true but is not why THIS pair is unguarded.
CompletionResolver.captureBaseline already catches RuntimeException around
its scrape read and fails open, so completion.onDelivered does not
realistically throw. Comment text only, no logic change.
2026-09-12 18:10:31 +07:00
Dai Ha 68b428c484 fleetd #551 rework (comment 17058): fix TIMED_OUT_QUEUED javadoc for ATTEMPTED
CI / shell-tests (pull_request) Successful in 13s
CI / contract (pull_request) Successful in 1m16s
CI / build (pull_request) Successful in 2m21s
Javadoc-only change. #551 added Injector.Cancellation.ATTEMPTED, which
MessageService.send's timeout path folds into Outcome.TIMED_OUT_QUEUED
alongside CANCELLED and NOT_DELIVERED (the collapse itself is unchanged
behaviour and is being tracked as a separate follow-up ticket).

TIMED_OUT_QUEUED's javadoc — landed by #513 to state the routes that
reach it — named only two routes and said "on every route it will not
arrive later", with a "may already hold a partial paste" caveat. Both
claims are now stale: ATTEMPTED is a third route, and because
agent.prompt pastes AND submits in one call, that route may mean the
target holds a complete, already-submitted turn and is working on it
right now.

Names all three routes, says which one is uncertain, and drops the
now-false blanket claim. No behaviour change.
2026-09-12 17:51:26 +07:00
Dai Ha e20ccab1eb fleetd #561: harden the completion/session TurnListener fan-out
CI / shell-tests (pull_request) Successful in 8s
CI / contract (pull_request) Successful in 1m4s
CI / build (pull_request) Successful in 1m56s
Fleetd's turnListener composition had four callbacks (onTurnComplete,
onTurnCompleteWithPostAction, and both onTurnFailed overloads) built from two
bare, unguarded statements each. onDelivered's registration was already fixed
structurally by #556; these four had the identical fragility and were still
untested: nothing enforced that the completion resolver's half ran before the
session half beyond call order in the source, so a future reorder (or a
throwing session listener sequenced first) could silently skip the
completion resolver's effect and strand a caller for its full timeout.

Extracted the composition to a package-private static factory,
Fleetd.turnListener(completion, sessions), and hardened it with
bothMustRun/bothMustRunKeepingSecondResult: both callback halves are always
attempted regardless of whether the other throws, and whatever escapes is
rethrown afterward (never swallowed) so it still reaches StatusPoller's
catch (Throwable) and logs at ERROR.

FleetdTurnListenerCompositionTest builds this real composition from a real
CompletionResolver and a throwing fake sessions half, and asserts the
completion resolver's effect (the waiter resolving) survives the session
half throwing, for all four callbacks, plus a mirror case showing the
session half still runs when the completion half throws first.

onTurnCompleteWithPostAction keeps completion-before-session as a functional
requirement (resolveBeforePostAction must run before the context-reset
housekeeping can erase the pane), not just fault tolerance, so it is not
reorder-symmetric like the other three — documented in Fleetd.turnListener's
javadoc.
2026-09-12 17:47:53 +07:00
Dai Ha d83821bbce fleetd #551: record the delivery attempt before the irreversible send
CI / shell-tests (pull_request) Successful in 5s
CI / contract (pull_request) Successful in 57s
CI / build (pull_request) Successful in 2m23s
The Injector wrote the delivery outcome AFTER calling AgentControl.send(),
so a HerdrException thrown from the response half of that call (herdr
already replied, or may have) was recorded as a confident NOT_DELIVERED
for text that may already be sitting in the worker's pane.

Poll the queue entry and mark it Pending.State.ATTEMPTED before send() is
called, not after. On success it is upgraded to DELIVERED; on an ordinary
failure it stays ATTEMPTED (honest uncertainty), except a herdr
*_not_found error, which the rest of this codebase already treats as a
confirmed absence and which now still writes NOT_DELIVERED.

cancellationOf gets a matching third answer (Cancellation.ATTEMPTED)
instead of folding the new state into NOT_DELIVERED, so a caller that
cancels an already-attempted delivery is told the truth too.
2026-09-12 17:43:21 +07:00
12 changed files with 972 additions and 140 deletions
+1 -1
View File
@@ -144,7 +144,7 @@ the merge — and merging on a reviewer's word is delegating it by proxy.
| Confirm your own role | `fleet_whoami` |
| See backends available | `fleet_profiles` |
| Start a member | `fleet_spawn{role?, profile?, cwd?, worktree?, ticket?, sessionName?, resumeSessionId?}` → `sessionId` + `paneId` |
| See the fleet | `fleet_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) · one peer's state: `fleet_status{sessionId}` |
| See the fleet | `fleet_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) + `loopHealth` (`RUNNING`, `STALLED`, or `STOPPED` for `statusPoller` and `sessionReaper`) · one peer's state: `fleet_status{sessionId}` |
| Delegate (blocking) | `fleet_send{sessionId, content}` |
| Delegate (long task) | `fleet_send{sessionId, content, wait:false}` → ticket → `fleet_poll{ticket}` |
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
+161 -37
View File
@@ -21,6 +21,7 @@ import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.inject.LiveExhaustedPatterns;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.inject.TurnListener;
import dev.ltms.fleet.inject.MemberPresence;
@@ -79,6 +80,7 @@ import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.BooleanSupplier;
import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.function.Predicate;
@@ -491,42 +493,10 @@ public final class Fleetd {
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
TurnListener turnListener = new TurnListener() {
@Override
public void onTurnComplete(String target) {
completion.onTurnComplete(target);
sessions.onTurnComplete(target);
}
@Override
public boolean hasPostTurnAction(String target) {
return sessions.hasPostTurnAction(target);
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
completion.resolveBeforePostAction(target);
return sessions.onTurnCompleteWithPostAction(target);
}
@Override
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
completion.onDelivered(target, token);
sessions.onDelivered(target, token);
}
@Override
public void onTurnFailed(String target) {
completion.onTurnFailed(target);
sessions.onTurnFailed(target);
}
@Override
public void onTurnFailed(String target, String reason) {
completion.onTurnFailed(target, reason);
sessions.onTurnFailed(target);
}
};
// fleetd #561: extracted to a static factory (see turnListener below) — the anonymous class
// this replaced had two bare, unguarded statements per callback, and nothing enforced that
// the completion half went first beyond call order in the source.
TurnListener turnListener = turnListener(completion, sessions);
Predicate<String> deliverable = deliverableTo(presence, leads);
// fleetd #556: registration is wired directly to `completion`, not folded into the
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
@@ -696,10 +666,13 @@ public final class Fleetd {
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy);
FleetMcp.LoopHealthSource loopHealth = new FleetMcp.LoopHealthSource(poller::health,
() -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health());
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics,
capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)),
healthCoverageSource(config),
loopHealth,
quarantineSource,
leadMailbox,
outageSource,
@@ -791,7 +764,7 @@ public final class Fleetd {
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
callers, metrics, deliverable,
() -> MemberCredentialPolicyView.of(config.get().memberCredentials()),
quarantineSource, outageSource).build();
quarantineSource, outageSource, loopHealth).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("fleetd listening on {}:{}, herdr socket {}",
cfg.bind().host(), cfg.bind().port(), socket);
@@ -1096,6 +1069,157 @@ public final class Fleetd {
.orElse(null);
}
/**
* fleetd #561: compose the production {@link TurnListener} from its two halves — the
* completion resolver (which resolves a blocked {@code fleet_send}'s waiter) and the session
* manager (which drives the member's lifecycle state) — hardened so a throw from either
* half's callback can never suppress the other half's callback for the same event.
*
* <p>Before this, each method below was two bare, unguarded statements: whichever ran first
* throwing meant the second one never ran at all, and nothing beyond call order in the
* source enforced "the completion half goes first". #556 fixed the identical shape for {@code
* onDelivered}'s registration by moving it off the fan-out entirely (see {@code
* TurnRegistrar}); this fixes the four remaining callbacks — {@code onTurnComplete}, {@code
* onTurnCompleteWithPostAction}, and both {@code onTurnFailed} overloads — by hardening the
* fan-out itself instead, since none of them can be pulled out of the listener the way
* registration was.
*
* <p>The invariant this composition guarantees: <b>a throwing session-listener must not
* prevent the completion resolver from being told the turn ended.</b> The completion half is
* always attempted first, and — as a bonus the resolver does not depend on — its own throw
* does not stop the session half from running either. Whatever escapes (from one half or
* both) is rethrown once both have been attempted, with a second failure recorded via {@link
* Throwable#addSuppressed} on the first, so it still reaches {@link StatusPoller}'s {@code
* catch (Throwable)} and logs at ERROR. Nothing here swallows a failure to make the two halves
* look safe.
*
* <p>{@code onTurnCompleteWithPostAction} is the one place order is not just fault-tolerance
* but a functional requirement: {@code completion.resolveBeforePostAction} must resolve the
* scrape before {@code sessions.onTurnCompleteWithPostAction}'s adapter housekeeping can erase
* the pane's rendered output (see {@link CompletionResolver#resolveBeforePostAction}). That is
* why this composition does not treat the pair symmetrically the way {@link #bothMustRun}
* does for the other three callbacks: the mirror case (completion half throws, session half's
* return value still observed) is not preserved here — once the completion half's failure
* escapes, the session half's return value is discarded, matching how the {@link Injector}
* already treats any throw from this callback as "the action did not start" (see the {@code
* started} default at its call site, {@code Injector.java} ~line 616).
*
* <p>Package-private so {@code FleetdTurnListenerCompositionTest} can build this listener
* directly from a real {@link CompletionResolver} and a fake {@link TurnListener} standing in
* for {@code sessions}, without booting the rest of {@code main}.
*/
static TurnListener turnListener(CompletionResolver completion, TurnListener sessions) {
return new TurnListener() {
@Override
public void onTurnComplete(String target) {
bothMustRun(() -> completion.onTurnComplete(target), () -> sessions.onTurnComplete(target));
}
@Override
public boolean hasPostTurnAction(String target) {
return sessions.hasPostTurnAction(target);
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
return bothMustRunKeepingSecondResult(() -> completion.resolveBeforePostAction(target),
() -> sessions.onTurnCompleteWithPostAction(target));
}
@Override
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
// fleetd #561: left unguarded on purpose, not because registration survives a
// throw elsewhere. completion.onDelivered runs CompletionResolver.captureBaseline,
// which already wraps its scrape read in its own catch (RuntimeException) and
// fails open (baseline = null) — so this call does not realistically throw, and
// there is nothing here for bothMustRun to protect.
completion.onDelivered(target, token);
sessions.onDelivered(target, token);
}
@Override
public void onTurnFailed(String target) {
bothMustRun(() -> completion.onTurnFailed(target), () -> sessions.onTurnFailed(target));
}
@Override
public void onTurnFailed(String target, String reason) {
bothMustRun(() -> completion.onTurnFailed(target, reason), () -> sessions.onTurnFailed(target));
}
};
}
/**
* fleetd #561: run two listener-callback halves for one lifecycle event, guaranteeing the
* SECOND always runs even when the FIRST throws. Whatever is thrown is rethrown once both
* halves have been attempted — a second failure is attached to the first via {@link
* Throwable#addSuppressed} rather than dropped. Never swallows.
*/
private static void bothMustRun(Runnable completionHalf, Runnable sessionsHalf) {
Throwable failure = null;
try {
completionHalf.run();
} catch (Throwable t) {
failure = t;
}
try {
sessionsHalf.run();
} catch (Throwable t) {
if (failure == null) {
failure = t;
} else {
failure.addSuppressed(t);
}
}
if (failure != null) {
throwUnchecked(failure);
}
}
/**
* fleetd #561: like {@link #bothMustRun}, but for {@code onTurnCompleteWithPostAction}, whose
* session half returns the value the {@link Injector} needs. The second (session) half's
* result is what this method returns; if the first (completion) half throws, the second half
* still runs and its result is still computed here, but the throw is rethrown afterward
* regardless — so that result is discarded at the {@link Injector} call site exactly as it
* already is today when this callback throws (see {@link #turnListener}'s javadoc).
*/
private static boolean bothMustRunKeepingSecondResult(Runnable completionHalf,
BooleanSupplier sessionsHalf) {
Throwable failure = null;
try {
completionHalf.run();
} catch (Throwable t) {
failure = t;
}
boolean result = false;
try {
result = sessionsHalf.getAsBoolean();
} catch (Throwable t) {
if (failure == null) {
failure = t;
} else {
failure.addSuppressed(t);
}
}
if (failure != null) {
throwUnchecked(failure);
}
return result;
}
/**
* fleetd #561: rethrow a captured {@link Throwable} without a checked-exception wrapper. The
* two callback halves above never declare a checked exception (both existing production
* halves — {@code CompletionResolver} and {@code SessionManager} — only ever throw unchecked),
* so this only ever actually rethrows a {@link RuntimeException} or {@link Error}; the generic
* cast is the standard "sneaky throw" idiom, not a claim that a checked exception is expected.
*/
@SuppressWarnings("unchecked")
private static <T extends Throwable> void throwUnchecked(Throwable t) throws T {
throw (T) t;
}
/**
* fleetd #480: construct the {@link LeadRollover} executor only when {@code leadRollover:} is
* present at startup — the same presence gate {@code leadHeartbeat:} uses just above this
@@ -2,6 +2,7 @@ package dev.ltms.fleet.inject;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.HerdrRouter;
import dev.ltms.fleet.msg.TurnToken;
import org.slf4j.Logger;
@@ -46,6 +47,17 @@ import java.util.stream.Collectors;
* turn — the pickup-grace path (a turn too fast to sample) unwedges the queue but does not fire
* completion, since without a sampled {@code working} there is no trustworthy "the worker just
* finished the task" signal to act on.
*
* <p><strong>Delivery honesty (fleetd #551).</strong> The one external send call
* ({@link AgentControl#send}) is irreversible, and its own response can still fail after the text
* has already reached the worker's pane — herdr replied (or may have), so the request was already
* processed, but the reply itself then failed to parse or carried an error. A queued entry is
* polled off the queue and marked {@code ATTEMPTED} <em>before</em> that call is made, not after,
* so a failure from the call itself can never be recorded as a confident {@code NOT_DELIVERED} for
* text that may already be sitting in the pane. {@code NOT_DELIVERED} stays reserved for the cases
* where nothing was ever attempted — the readiness grace expiring, {@link #drop}, or a herdr
* {@code *_not_found} error, which the rest of this codebase already treats as a confirmed absence
* rather than a merely inconclusive failure (see {@code StatusPoller}, {@code AgentControl}).
*/
public final class Injector {
@@ -217,8 +229,23 @@ public final class Injector {
/** The result of trying to remove an undelivered message from the injector. */
public enum Cancellation {
/** The message was still queued and this call removed it; the target saw nothing. */
CANCELLED,
/** The message reached the worker's pane and is confirmed delivered. */
DELIVERED,
/**
* fleetd #551: the injector called {@link AgentControl#send} for this message and that call
* threw before its outcome was known — the text may or may not have reached the pane. A
* caller reporting this to an operator must say "uncertain", not "definitely not
* delivered": reading it as a confident negative invites a resend of text that may already
* be sitting in the pane (a double delivery), which is worse than the ambiguity itself.
*/
ATTEMPTED,
/**
* The message never reached the worker's pane — nothing was ever attempted for it (the
* readiness grace expired, the target was dropped, or send failed with a herdr
* {@code *_not_found} error, which is a confirmed absence, not merely inconclusive).
*/
NOT_DELIVERED
}
@@ -241,7 +268,29 @@ public final class Injector {
/** A pending message and the future that completes when it has been delivered. */
private static final class Pending {
enum State { QUEUED, DELIVERED, NOT_DELIVERED, CANCELLED }
enum State {
/** Still in the target's queue, not yet attempted. */
QUEUED,
/**
* fleetd #551: {@link AgentControl#send} has been called for this entry and its outcome
* is not yet known — recorded BEFORE the call (see {@link #onStatus}), so a Throwable
* from send() can never leave the entry either still QUEUED or falsely marked
* NOT_DELIVERED. Upgraded to DELIVERED on success; left as ATTEMPTED on an ordinary
* failure, since reaching the catch does not prove the text never reached the pane.
*/
ATTEMPTED,
/** {@code send()} returned normally: the text is confirmed to have reached the pane. */
DELIVERED,
/**
* Confirmed — not merely inconclusive — that nothing was ever sent for this entry: the
* worker never became ready ({@link #onStatus}'s readiness-grace expiry), the target was
* dropped ({@link #drop}), or send failed with a herdr {@code *_not_found} error. Never
* written for an entry whose send() outcome is unknown; see ATTEMPTED.
*/
NOT_DELIVERED,
/** Removed from the queue by {@link #cancel} before it was ever attempted. */
CANCELLED
}
final String target;
final String text;
@@ -333,7 +382,14 @@ public final class Injector {
}
private static Cancellation cancellationOf(Pending p) {
return p.state == Pending.State.DELIVERED ? Cancellation.DELIVERED : Cancellation.NOT_DELIVERED;
// fleetd #551 (comment 17037): a two-way split on a now-three-way question folded ATTEMPTED
// into NOT_DELIVERED with no compiler error and no failing test — the exact defect this
// ticket exists to fix, one layer up. ATTEMPTED gets its own answer instead.
return switch (p.state) {
case DELIVERED -> Cancellation.DELIVERED;
case ATTEMPTED -> Cancellation.ATTEMPTED;
case QUEUED, NOT_DELIVERED, CANCELLED -> Cancellation.NOT_DELIVERED;
};
}
private static boolean isQuiescent(Target t) {
@@ -426,9 +482,18 @@ public final class Injector {
if (p != null && ready.test(target)) {
t.notReadySincePoll = 0;
t.notReadySinceMillis = 0;
// fleetd #551: poll and record BEFORE the irreversible send, not after.
// The entry comes off the queue and its state is set to ATTEMPTED here,
// unconditionally — so a Throwable escaping the send call below (caught or
// not) can never leave the entry QUEUED at the head of t.queue (the fleetd
// #546 hazard, since peek() alone would let the next onStatus round re-enter
// this block and send the same text again), and no path can write a
// confident DELIVERED or NOT_DELIVERED before we actually know which one
// happened.
t.queue.poll();
p.state = Pending.State.ATTEMPTED;
try {
agentsFor(target).send(target, p.text());
t.queue.poll();
p.state = Pending.State.DELIVERED;
t.awaitingPickup = true;
t.awaitingCompletion = true;
@@ -436,15 +501,24 @@ public final class Injector {
t.injectableSincePickup = 0;
sent = p;
} catch (Throwable e) {
// Delivery failed at herdr; drop the poisoned message and surface it
// rather than blocking the queue behind it. Catches Throwable, not just
// RuntimeException: fleetd #546 — an Error escaping this send (e.g. a
// NoClassDefFoundError, see #413) would otherwise leave the entry QUEUED
// at the head of t.queue. Line :378 peeks rather than polls, so the next
// onStatus round would re-enter this try and send the same text again,
// typing the same brief into the member's pane a second time.
t.queue.poll();
p.state = Pending.State.NOT_DELIVERED;
// fleetd #551: leave p.state == ATTEMPTED (recorded above, before the
// call) rather than downgrading it to NOT_DELIVERED here — reaching this
// catch does not prove the text never reached the pane. Three of the
// four HerdrException throw sites in HerdrCodec fire only after herdr
// has already replied (so it processed the request), and the fourth (a
// transport IOException) leaves it genuinely unknown whether herdr even
// received the bytes — see #551 comment 16867. The one exception is a
// herdr `*_not_found` error: that family is already read as "definitely
// absent, not merely inconclusive" everywhere else in this codebase
// (StatusPoller, AgentControl's own retry, WorkspaceControl,
// HerdrPeerLauncher, FleetApp, ReplyPushLoop) because it means the
// target pane/agent does not exist at all, so nothing could have been
// pasted anywhere — #551 keeps the new state consistent with that
// existing vocabulary rather than inventing a second one.
if (e instanceof HerdrException he && he.code() != null
&& he.code().endsWith("_not_found")) {
p.state = Pending.State.NOT_DELIVERED;
}
sent = p;
sendError = e;
}
@@ -20,6 +20,7 @@ import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.placement.BackendOutagePolicy;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.MemberSession;
@@ -108,6 +109,7 @@ public final class FleetMcp {
private final Metrics metrics; // CB-502: null → auth failures not counted
private final CapacitySource capacity;
private final HealthCoverageSource healthCoverage;
private final LoopHealthSource loopHealth;
private final QuarantineSource quarantine;
/** fleetd #201 Unit 5: SEPARATE from {@link #quarantine} — see {@link OutageSource}'s doc. */
private final OutageSource outage;
@@ -136,6 +138,15 @@ public final class FleetMcp {
/** Coverage is supplied by the health wiring, not inferred from a missing dependency. */
public record HealthCoverageSource(Supplier<String> value) { }
/** Progress states for fleetd's singleton background loops, read by {@code fleet_list} and {@code /healthz}. */
public record LoopHealthSource(Supplier<LoopWatchdog.State> statusPoller,
Supplier<LoopWatchdog.State> sessionReaper) {
/** Inert source for callers that do not wire the background loops. */
public static LoopHealthSource none() {
return new LoopHealthSource(() -> LoopWatchdog.State.STOPPED, () -> LoopWatchdog.State.STOPPED);
}
}
/**
* CB-578 stage B quarantine facts used by {@code fleet_profiles}: a profile → credential id
* lookup, plus the shared {@link BackendQuarantine} to read remaining cooldowns off.
@@ -328,11 +339,22 @@ public final class FleetMcp {
* instead of throwing. See {@link #handover}.
*/
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
this(messages, workers, sessions, identity, presence, primaryRegistry, callers, authorizationMode, metrics,
capacity, healthCoverage, LoopHealthSource.none(), quarantine, leadChannel, outage, leadSeats,
peers, leadRollover);
}
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
CapacitySource capacity, HealthCoverageSource healthCoverage, LoopHealthSource loopHealth,
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
Objects.requireNonNull(callers, "callers");
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
== AuthorizationMode.ENFORCED;
@@ -343,6 +365,7 @@ public final class FleetMcp {
this.outage = Objects.requireNonNull(outage, "outage");
this.leadSeats = Objects.requireNonNull(leadSeats, "leadSeats");
this.healthCoverage = healthCoverage;
this.loopHealth = Objects.requireNonNull(loopHealth, "loopHealth");
this.leadRollover = leadRollover;
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
@@ -474,7 +497,7 @@ public final class FleetMcp {
(exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
if (denied != null) return denied;
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, outage,
leadSeats, callers.leads(),
callerTerminal(exchange),
new CoordinationSource(leadChannel, peers),
@@ -1546,25 +1569,42 @@ public final class FleetMcp {
* @param selfTerm the calling pane's terminal id, or blank for a caller with no pane
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions,
Map<String, String> leads, String selfTerm) {
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, null, CapacitySource.none(), new HealthCoverageSource(() -> "off"),
QuarantineSource.none(), leads, selfTerm);
LoopHealthSource.none(), QuarantineSource.none(), leads, selfTerm);
}
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, leads, selfTerm,
CoordinationSource.none());
}
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm,
LoopHealthSource loopHealth, QuarantineSource quarantine,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine, leads, selfTerm,
CoordinationSource.none());
}
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
LoopHealthSource loopHealth, QuarantineSource quarantine,
Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, loopHealth, quarantine,
OutageSource.none(), LeadSeatSource.none(), leads, selfTerm, coordination, false);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none());
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, CoordinationSource.none(), false);
}
/**
@@ -1580,8 +1620,8 @@ public final class FleetMcp {
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, OutageSource.none(),
LeadSeatSource.none(), leads, selfTerm, coordination);
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, OutageSource.none(),
LeadSeatSource.none(), leads, selfTerm, coordination, false);
}
/** As above, plus fleetd #201 Unit 5 cool-off facts (see {@link OutageSource}). */
@@ -1589,8 +1629,8 @@ public final class FleetMcp {
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
Map<String, String> leads, String selfTerm, CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, coordination);
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
LeadSeatSource.none(), leads, selfTerm, coordination, false);
}
/**
@@ -1611,7 +1651,7 @@ public final class FleetMcp {
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine, outage,
leadSeats, leads, selfTerm, coordination, false);
}
@@ -1631,8 +1671,18 @@ public final class FleetMcp {
* an explicit {@code true}
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
return listFleet(workers, sessions, messages, capacity, healthCoverage, LoopHealthSource.none(), quarantine,
outage, leadSeats, leads, selfTerm, coordination, callerIsPrimary);
}
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
LoopHealthSource loopHealth,
QuarantineSource quarantine, OutageSource outage,
LeadSeatSource leadSeats, Map<String, String> leads, String selfTerm,
CoordinationSource coordination, boolean callerIsPrimary) {
try {
@@ -1656,6 +1706,9 @@ public final class FleetMcp {
Map<String, Object> result = new LinkedHashMap<>();
result.put("leads", leadRows); result.put("members", out);
result.put("healthCoverage", healthCoverage.value().get());
result.put("loopHealth", Map.of(
"statusPoller", loopHealth.statusPoller().get().name(),
"sessionReaper", loopHealth.sessionReaper().get().name()));
// fleetd #439: coordinator/coordinatorView is lead-to-lead coordination state and must
// never reach a worker or an architect -- gate BEFORE assembling it, not after, so the
// key is absent rather than present-and-empty.
@@ -2135,9 +2188,10 @@ public final class FleetMcp {
+ "cannot reliably re-identify: some backends (e.g. opencode) resolve it from the "
+ "member's working directory, which only uniquely identifies a member when it "
+ "was spawned into its own fleetd-provisioned worktree (worktree:true/<slug>); a "
+ "member spawned without one shares its directory with others and never reports "
+ "an id, however long it runs (fleetd #249). An empty 'members' "
+ "means no members are spawned; it says nothing about peers. When capacity "
+ "member spawned without one shares its directory with others and never reports "
+ "an id, however long it runs (fleetd #249). An empty 'members' "
+ "means no members are spawned; it says nothing about peers. 'loopHealth' reports "
+ "the RUNNING, STALLED, or STOPPED state of statusPoller and sessionReaper. When capacity "
+ "facts are configured, a 'capacity' row per profile reports 'free' — the "
+ "slots a fresh fleet_spawn on that profile will actually be granted right "
+ "now (max(0, maxLoad - live)), the same check the spawn gate itself runs. A "
@@ -94,16 +94,22 @@ public final class MessageService {
TIMED_OUT_WORKING,
/**
* Timed out with no confirmed delivery. Despite the name, this does not mean the message
* is sitting in a queue — on every route it will not arrive later. {@link #send} reaches
* this outcome through {@link Injector#cancel}, whose result tells the routes apart:
* is sitting in a queue. {@link #send} reaches this outcome through {@link Injector#cancel},
* whose result tells three routes apart:
* {@link Injector.Cancellation#CANCELLED} means the message was still queued and this call
* removed it, so the target saw nothing; {@link Injector.Cancellation#NOT_DELIVERED} means
* an earlier attempt already decided the message's fate — the injector's call to the
* target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) threw, or the queue
* was cleared because the target never became ready or was abandoned — and {@code cancel}
* is only reporting that pre-existing state. On the failed-attempt route the terminal may
* already hold a partial paste from before the call threw, so only the {@code CANCELLED}
* case establishes that the target saw nothing.
* removed it, so the target saw nothing and it will not arrive later;
* {@link Injector.Cancellation#NOT_DELIVERED} means nothing was ever sent — the queue was
* cleared because the target never became ready or was abandoned, or the injector's call to
* the target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) failed with a herdr
* error that this codebase already treats as a confirmed absence — so this route too
* establishes that the target saw nothing and it will not arrive later; but
* {@link Injector.Cancellation#ATTEMPTED} (fleetd #551) means that call was made and its
* outcome is unknown. {@code agent.prompt} pastes <em>and submits</em> in one call, so on
* this route the target may hold a complete, already-submitted turn and be working on it
* right now — {@link Outcome#TIMED_OUT_WORKING}'s meaning, reported here as
* {@code TIMED_OUT_QUEUED} only because this caller never observed the pickup. Only
* {@code CANCELLED} and {@code NOT_DELIVERED} establish that the target saw nothing;
* {@code ATTEMPTED} does not.
*/
TIMED_OUT_QUEUED,
/** Another send to this session was in flight for the whole window. */
@@ -313,16 +319,18 @@ public final class MessageService {
/**
* Targets whose last send timed out with {@link Outcome#TIMED_OUT_QUEUED} (CB-640) — {@link
* #send} called {@link Injector#cancel} and got back something other than {@code DELIVERED}.
* That covers more than one history: the message may have still been queued and {@code cancel}
* removed it right there, or an earlier attempt may already have failed (the call to the
* target's terminal threw) or been abandoned (the target never became ready, or was torn
* down). What holds on every route: the message will not arrive later, and it is not sitting
* in any queue. What does NOT hold on every route: that the target saw nothing — a failed
* delivery attempt can leave a partial paste behind. Set where {@link #send} already computes
* {@code wasDelivered} for that outcome; no queue is kept here, only the fact that the send
* ended with no confirmed delivery. Cleared the same way as {@link #strandedReplies}: the next
* accepted delivery for the target ({@link #send} opening a fresh waiter) or a teardown
* ({@link #abandon}).
* That covers three histories, not one: {@link Injector.Cancellation#CANCELLED} — the message
* was still queued and {@code cancel} removed it right there; {@link
* Injector.Cancellation#NOT_DELIVERED} — nothing was ever sent, because the target never became
* ready, was torn down, or the call to its terminal failed with a herdr error this codebase
* already treats as a confirmed absence; or {@link Injector.Cancellation#ATTEMPTED} (fleetd
* #551) — the call to the target's terminal was made and its outcome is unknown, so the target
* may already hold a complete, submitted turn. Only the first two mean the message will not
* arrive later and the target saw nothing; on the third it may already have arrived in full.
* Set where {@link #send} already computes {@code wasDelivered} for that outcome; no queue is
* kept here, only the fact that the send ended with no confirmed delivery. Cleared the same way
* as {@link #strandedReplies}: the next accepted delivery for the target ({@link #send} opening
* a fresh waiter) or a teardown ({@link #abandon}).
*/
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
@@ -412,14 +420,18 @@ public final class MessageService {
* with no confirmed delivery — the caller saw {@link Outcome#TIMED_OUT_QUEUED} (see the
* {@code TimeoutException} branch of {@link #send}). Despite the method's name, this is not
* proof that a message is sitting in a queue: {@link Injector#cancel} reports this outcome
* either because the message was still queued and got removed right there, or because an
* earlier attempt already failed (the call to the target's terminal threw) or was abandoned
* (the target never became ready, or was torn down). Either way the message will not arrive
* later. It does NOT follow that the target saw nothing — on the failed-attempt route the
* terminal may already hold a partial paste. Distinct from {@link Outcome#TIMED_OUT_WORKING},
* where delivery already happened and only the reply is outstanding. Cleared the next time
* this target's delivery is accepted or the target is abandoned — see
* {@link #queuedDeliveries}.
* through three routes. {@link Injector.Cancellation#CANCELLED} means the message was still
* queued and got removed right there. {@link Injector.Cancellation#NOT_DELIVERED} means
* nothing was ever sent — the target never became ready, was torn down, or the call to its
* terminal failed with a herdr error this codebase already treats as a confirmed absence.
* Only these two routes mean the message will not arrive later. {@link
* Injector.Cancellation#ATTEMPTED} (fleetd #551) means the call to the target's terminal was
* made and its outcome is unknown: {@code agent.prompt} pastes <em>and submits</em> in one
* call, so on this route the target may already hold a complete, submitted turn and be
* working on it right now — it does NOT follow that the target saw nothing. Distinct from
* {@link Outcome#TIMED_OUT_WORKING}, where delivery already happened and only the reply is
* outstanding. Cleared the next time this target's delivery is accepted or the target is
* abandoned — see {@link #queuedDeliveries}.
*/
public boolean hasQueuedDelivery(String target) {
return target != null && queuedDeliveries.containsKey(target);
@@ -974,11 +986,12 @@ public final class MessageService {
}
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
if (!wasDelivered) {
// CB-640: record that delivery did not happen for fleet health (see
// CB-640: record that delivery is not confirmed, for fleet health (see
// queuedDeliveries). Whatever injector.cancel() reported above — this call
// removed a still-queued Pending, or an earlier attempt already failed or
// was abandoned — the send ends with no confirmed delivery and will not
// arrive later.
// removed a still-queued Pending (CANCELLED), an earlier attempt already
// failed with a confirmed absence (NOT_DELIVERED), or an earlier attempt was
// made and its outcome is unknown (ATTEMPTED, fleetd #551 — the message may
// already have arrived in full) — the send ends with no confirmed delivery.
queuedDeliveries.put(target, Boolean.TRUE);
}
return recorded(new Reply(
@@ -1122,21 +1135,33 @@ public final class MessageService {
}
try {
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
// resumed turn can re-associate the async ticket with its new turnId via
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
// markAsyncQuestion silently returns null.
Task task = asyncTasksByTurn.get(turnId);
if (task != null) {
asyncTasksByWaiter.put(reply, task);
}
if (!rendezvous.answerAsk(turnId, content)) {
asyncTasksByWaiter.remove(reply);
rendezvous.close(workerSession, reply);
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
}
clearAsyncQuestion(turnId, false);
// fleetd #575: this try used to open below, AFTER the Task lookup/registration and the
// STALE_TURN early return that follows it — so that return was covered only by a
// hand-rolled copy of the finally's own cleanup pair, not the finally itself. Widening the
// try up to wrap the registration closes the gap structurally: every exit from here on,
// STALE_TURN included, now runs through the one finally below exactly once, and the
// duplicated pair is gone. This did not fix a live leak — see the ticket: neither
// rendezvous.answerAsk nor clearAsyncQuestion(turnId, false) can throw, so nothing ever
// actually left through the old gap uncovered — but #572 found this exact drift (one
// finally asserted, an identical sibling not) on this same file, and the hand-rolled copy
// was the wrong shape to keep regardless of whether it was ever exercised.
try {
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
// resumed turn can re-associate the async ticket with its new turnId via
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
// markAsyncQuestion silently returns null.
Task task = asyncTasksByTurn.get(turnId);
if (task != null) {
asyncTasksByWaiter.put(reply, task);
}
if (answerAskLapseRaceHookForTest != null) {
// Test-only (fleetd #575): see the field's own javadoc.
answerAskLapseRaceHookForTest.run();
}
if (!rendezvous.answerAsk(turnId, content)) {
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
}
clearAsyncQuestion(turnId, false);
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
// #282: this waiter can resolve with a FRESH question rather than a terminal reply —
@@ -1636,6 +1661,30 @@ public final class MessageService {
this.askTimeoutRaceHookForTest = hook;
}
/**
* Null in production; test seam for fleetd #575 — invoked from {@link #answer}, right after this
* call's own {@code Task} registration and right before its {@code rendezvous.answerAsk(turnId,
* content)} call. A test installs this to complete the SAME turnId's ask directly via {@link
* Rendezvous#answerAsk} from inside that exact window, deterministically reproducing what a
* second, concurrent {@code answer()} call racing to unblock the same ask can otherwise only
* win by timing luck: this call's own {@code askSession(turnId)} lookup at the top already saw
* the ask as open, but by the time it reaches {@code rendezvous.answerAsk} here, the other call
* already completed it (or the worker's own {@code ask()} teardown already closed it) — so this
* call must see {@code false} and return {@link Outcome#STALE_TURN}, exactly the "lapsed between
* the lookup and the unblock" case named at that call site. Proves the #575 fix (widening this
* method's try so a single finally covers this exit) does not change that outcome and still
* cleans this call's own {@code reply} up exactly once.
*/
private volatile Runnable answerAskLapseRaceHookForTest;
/**
* Test-only (fleetd #575): install {@link #answerAskLapseRaceHookForTest}. Package-private so the
* test, in the same package, can reach it without widening any production API.
*/
void setAnswerAskLapseRaceHookForTest(Runnable hook) {
this.answerAskLapseRaceHookForTest = hook;
}
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
private boolean hasAsyncQuestion(String target) {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
@@ -94,6 +94,7 @@ public final class FleetApp {
// .none() (the honest "feature not wired" view) for every constructor that does not pass one.
private final FleetMcp.QuarantineSource quarantine;
private final FleetMcp.OutageSource outage;
private final FleetMcp.LoopHealthSource loopHealth;
private final ObjectMapper mapper = new ObjectMapper();
/**
@@ -155,7 +156,8 @@ public final class FleetApp {
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials) {
this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics,
deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none());
deliverable, memberCredentials, FleetMcp.QuarantineSource.none(), FleetMcp.OutageSource.none(),
FleetMcp.LoopHealthSource.none());
}
/**
@@ -169,8 +171,18 @@ public final class FleetApp {
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage) {
this(herdr, memberHerdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, deliverable,
memberCredentials, quarantine, outage, FleetMcp.LoopHealthSource.none());
}
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
Predicate<String> deliverable, Supplier<MemberCredentialPolicyView> memberCredentials,
FleetMcp.QuarantineSource quarantine, FleetMcp.OutageSource outage,
FleetMcp.LoopHealthSource loopHealth) {
this.herdr = herdr;
this.memberHerdr = memberHerdr != null ? memberHerdr : herdr;
this.workers = workers;
@@ -183,6 +195,7 @@ public final class FleetApp {
this.memberCredentials = memberCredentials != null ? memberCredentials : MemberCredentialPolicyView::absent;
this.quarantine = quarantine != null ? quarantine : FleetMcp.QuarantineSource.none();
this.outage = outage != null ? outage : FleetMcp.OutageSource.none();
this.loopHealth = loopHealth != null ? loopHealth : FleetMcp.LoopHealthSource.none();
}
/** Wire routes onto a fresh, unstarted Javalin instance. Caller starts it. */
@@ -291,15 +304,19 @@ public final class FleetApp {
* spotted by comparing two numbers by eye.
*/
private void healthz(Context ctx) {
HealthzResponse response = healthzResponse(herdr, memberHerdr, loopHealth);
ctx.status(response.status()).json(response.body());
}
record HealthzResponse(int status, Map<String, Object> body) { }
static HealthzResponse healthzResponse(HerdrClient herdr, HerdrClient memberHerdr,
FleetMcp.LoopHealthSource loopHealth) {
JsonNode pong;
try {
pong = herdr.call("ping");
} catch (HerdrException e) {
ctx.status(503).json(Map.of(
"status", "degraded",
"herdr", "unreachable",
"detail", e.getMessage()));
return;
return degradedResponse("unreachable", e.getMessage(), loopHealth);
}
Map<String, Object> body = new LinkedHashMap<>();
body.put("status", "ok");
@@ -311,11 +328,11 @@ public final class FleetApp {
try {
memberPong = memberHerdr.call("ping");
} catch (HerdrException e) {
ctx.status(503).json(Map.of(
return new HealthzResponse(503, Map.of(
"status", "degraded",
"herdr", "member unreachable",
"detail", e.getMessage()));
return;
"detail", e.getMessage(),
"loopHealth", loopHealthView(loopHealth)));
}
int leadProtocol = pong.path("protocol").asInt();
int memberProtocol = memberPong.path("protocol").asInt();
@@ -326,7 +343,20 @@ public final class FleetApp {
body.put("protocolMismatch", true);
}
}
ctx.status(200).json(body);
body.put("loopHealth", loopHealthView(loopHealth));
return new HealthzResponse(200, body);
}
private static Map<String, String> loopHealthView(FleetMcp.LoopHealthSource loopHealth) {
return Map.of("statusPoller", loopHealth.statusPoller().get().name(),
"sessionReaper", loopHealth.sessionReaper().get().name());
}
private static HealthzResponse degradedResponse(String herdr, String detail,
FleetMcp.LoopHealthSource loopHealth) {
return new HealthzResponse(503, Map.of(
"status", "degraded", "herdr", herdr, "detail", detail,
"loopHealth", loopHealthView(loopHealth)));
}
/**
@@ -0,0 +1,286 @@
package dev.ltms.fleet;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.inject.TurnListener;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TurnToken;
import org.junit.jupiter.api.Test;
import java.util.LinkedHashSet;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #561: {@code Fleetd.turnListener} composes the completion resolver and the session
* manager into one {@link TurnListener}. Each of the four callbacks below used to be two bare,
* unguarded statements (completion first, then sessions) — a throw from the session half used to
* skip nothing <em>after</em> it (there was nothing after it), but nothing enforced that the
* completion half had to come first either, beyond call order in the source. {@code onDelivered}
* had the identical shape and was fixed by #556 (moving its registration off the fan-out
* entirely); these four callbacks cannot be fixed that way, so the fan-out itself is hardened
* instead (see {@link Fleetd#turnListener} and {@code Fleetd.bothMustRun}/{@code
* bothMustRunKeepingSecondResult}).
*
* <p>The invariant under test: <b>a throwing session-listener must not prevent the completion
* resolver from being told the turn ended.</b> Each test below builds the exact production
* composition ({@link Fleetd#turnListener}) from a real {@link CompletionResolver} and a fake
* {@code sessions} half that throws, then asserts the completion half's effect (the captured
* {@link Rendezvous} waiter resolving) happened anyway — never by inspecting call order directly.
*/
class FleetdTurnListenerCompositionTest {
/** Records which callbacks ran and can be told to throw from a chosen one. */
private static final class RecordingSessions implements TurnListener {
final Set<String> called = new LinkedHashSet<>();
private final Set<String> throwing;
RecordingSessions(String... throwingMethods) {
this.throwing = Set.of(throwingMethods);
}
private void maybeThrow(String method) {
called.add(method);
if (throwing.contains(method)) {
throw new IllegalStateException("boom: sessions." + method);
}
}
@Override
public void onTurnComplete(String target) {
maybeThrow("onTurnComplete");
}
@Override
public boolean onTurnCompleteWithPostAction(String target) {
maybeThrow("onTurnCompleteWithPostAction");
return true;
}
@Override
public void onTurnFailed(String target) {
maybeThrow("onTurnFailed");
}
@Override
public void onTurnFailed(String target, String reason) {
maybeThrow("onTurnFailedWithReason");
}
}
private static CompletionResolver newResolver(FakeHerdr herdr, Rendezvous rendezvous) {
return new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
}
/**
* A resolver with a controllable clock, and the clock itself, so a test can register a turn
* "delivered" at time 0 and then jump the clock past {@link CompletionResolver#MIN_TURN_NANOS}
* before resolving it — otherwise {@code onTurnComplete}/{@code onTurnCompleteWithPostAction}
* resolve within microseconds of registering in-test, well inside the fleetd#164 floor, and
* get classified as a too-fast crash rather than a real completion. Mirrors the {@code
* LongSupplier} clock-injection pattern fleetd#164's own tests use.
*/
private static CompletionResolver newResolverPastTheFloor(FakeHerdr herdr, Rendezvous rendezvous,
AtomicLong clock) {
LongSupplier nowNanos = clock::get;
return new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), nowNanos);
}
@Test
void onTurnCompleteResolvesTheWaiterEvenWhenTheSessionHalfThrows() throws Exception {
FakeHerdr herdr = new FakeHerdr().readText("⏺ real answer\n❯ ");
Rendezvous rendezvous = new Rendezvous();
AtomicLong clock = new AtomicLong(0);
CompletionResolver completion = newResolverPastTheFloor(herdr, rendezvous, clock);
var waiter = rendezvous.open("term_a");
completion.register("term_a", new TurnToken("term_a", waiter)); // delivered at clock=0
clock.set(CompletionResolver.MIN_TURN_NANOS + 1_000_000); // past the too-fast floor
RecordingSessions sessions = new RecordingSessions("onTurnComplete");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnComplete("term_a"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnComplete", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnComplete"), "the session half must have run");
// onTurnComplete resolves off-thread (a virtual thread) — wait on the waiter itself,
// exactly like MessageServiceTest.completionFallbackIsNeverQueued does.
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
assertEquals(Rendezvous.Kind.COMPLETION, resolution.kind(),
"the completion resolver must still have resolved the send, despite the session "
+ "half throwing");
assertEquals("real answer", resolution.text());
}
@Test
void onTurnFailedResolvesTheWaiterAsFailedEvenWhenTheSessionHalfThrows() throws Exception {
FakeHerdr herdr = new FakeHerdr().readText("⏺ crash context\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
var waiter = rendezvous.open("term_b");
completion.register("term_b", new TurnToken("term_b", waiter));
RecordingSessions sessions = new RecordingSessions("onTurnFailed");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnFailed("term_b"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnFailed", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnFailed"), "the session half must have run");
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
assertEquals(Rendezvous.Kind.FAILED, resolution.kind(),
"the completion resolver must still have failed the send, despite the session "
+ "half throwing");
}
@Test
void onTurnFailedWithReasonResolvesTheWaiterEvenWhenTheSessionHalfThrows() throws Exception {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
var waiter = rendezvous.open("term_c");
completion.register("term_c", new TurnToken("term_c", waiter));
// fleetd #561: the composed listener calls sessions.onTurnFailed(target) — the ONE-arg
// overload — for this reason-carrying event too (matching the pre-existing production
// behaviour: SessionManager never overrides the two-arg overload either), so the
// throwing key here is "onTurnFailed", not a distinct "...WithReason" one.
RecordingSessions sessions = new RecordingSessions("onTurnFailed");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnFailed("term_c", "worker unreachable"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnFailed", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnFailed"), "the session half must have run");
Rendezvous.Resolution resolution = waiter.get(5, TimeUnit.SECONDS);
assertEquals(Rendezvous.Kind.FAILED, resolution.kind(),
"the completion resolver must still have failed the send, despite the session "
+ "half throwing");
assertEquals("worker unreachable", resolution.text(),
"the explicit reason must still reach the resolved send");
}
@Test
void onTurnCompleteWithPostActionResolvesTheWaiterEvenWhenTheSessionHalfThrows() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ answer before reset\n❯ ");
Rendezvous rendezvous = new Rendezvous();
AtomicLong clock = new AtomicLong(0);
CompletionResolver completion = newResolverPastTheFloor(herdr, rendezvous, clock);
var waiter = rendezvous.open("term_d");
completion.register("term_d", new TurnToken("term_d", waiter)); // delivered at clock=0
clock.set(CompletionResolver.MIN_TURN_NANOS + 1_000_000); // past the too-fast floor
RecordingSessions sessions = new RecordingSessions("onTurnCompleteWithPostAction");
TurnListener composed = Fleetd.turnListener(completion, sessions);
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> composed.onTurnCompleteWithPostAction("term_d"),
"the session half's throw must still escape the composed listener");
assertEquals("boom: sessions.onTurnCompleteWithPostAction", thrown.getMessage());
assertTrue(sessions.called.contains("onTurnCompleteWithPostAction"), "the session half must have run");
// resolveBeforePostAction is synchronous by design (it must run before the context reset
// can erase the pane) — the waiter is already resolved by the time the throw propagates.
assertTrue(waiter.isDone(), "resolveBeforePostAction is synchronous — the send must "
+ "already be resolved once the composed call returns (by throwing)");
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind());
assertEquals("answer before reset", waiter.getNow(null).text());
}
/**
* The mirror case (comment 17041's acceptance item 2): the completion half throws, and the
* session half still ran. The production {@link CompletionResolver} is deliberately defensive
* (its scrape reads are wrapped in {@code catch (RuntimeException)}, by the same fail-open
* design {@link CompletionResolver#captureBaseline} documents) so it essentially never throws
* synchronously in normal operation — forcing it to do so needs a genuinely-real trigger, not
* a fabricated one. {@code ConcurrentHashMap.get(null)} is that trigger: passing a {@code null}
* target makes {@code onTurnComplete}'s {@code inFlight.get(target)} throw a
* {@link NullPointerException} before it ever starts its resolving thread — a real code path,
* not a contrived one.
*/
@Test
void sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronously() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
RecordingSessions sessions = new RecordingSessions(); // throws from nothing
TurnListener composed = Fleetd.turnListener(completion, sessions);
assertThrows(NullPointerException.class, () -> composed.onTurnComplete(null),
"ConcurrentHashMap.get(null) inside CompletionResolver.onTurnComplete must still "
+ "escape the composed listener");
assertTrue(sessions.called.contains("onTurnComplete"),
"the session half must still have run even though the completion half threw first");
}
/**
* The mirror case above only exercises {@code onTurnComplete}/{@code bothMustRun}. {@code
* onTurnCompleteWithPostAction} is composed through the OTHER helper,
* {@code bothMustRunKeepingSecondResult}, and nothing previously asserted that its session half
* still runs when its completion half throws — that helper could be reverted to the pre-#561
* broken shape (run the session half only if the completion half did not throw) and the suite
* would still stay green. Uses the same real, non-fabricated trigger as the test above: a
* {@code null} target makes {@code resolveBeforePostAction}'s {@code inFlight.get(target)}
* throw a {@link NullPointerException} before {@code resolve} is ever entered.
*/
@Test
void sessionHalfStillRunsWhenTheCompletionHalfThrowsSynchronouslyForPostAction() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
RecordingSessions sessions = new RecordingSessions(); // throws from nothing
TurnListener composed = Fleetd.turnListener(completion, sessions);
assertThrows(NullPointerException.class, () -> composed.onTurnCompleteWithPostAction(null),
"ConcurrentHashMap.get(null) inside CompletionResolver.resolveBeforePostAction must "
+ "still escape the composed listener");
assertTrue(sessions.called.contains("onTurnCompleteWithPostAction"),
"the session half must still have run even though the completion half threw first");
}
/**
* {@code bothMustRun} must not just run both halves — it must not DROP a second failure when
* both halves throw. Forces the completion half to throw (the same real {@code
* inFlight.get(null)} NullPointerException trigger used above) while the session half throws a
* distinct {@link IllegalStateException}, and asserts the completion half's throwable is what
* escapes while the session half's throwable survives as a suppressed exception rather than
* being silently discarded.
*/
@Test
void bothFailuresEscapeWhenBothHalvesThrowDistinctExceptions() {
FakeHerdr herdr = new FakeHerdr();
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = newResolver(herdr, rendezvous);
RecordingSessions sessions = new RecordingSessions("onTurnComplete");
TurnListener composed = Fleetd.turnListener(completion, sessions);
NullPointerException thrown = assertThrows(NullPointerException.class,
() -> composed.onTurnComplete(null),
"the completion half's throw (NPE from inFlight.get(null)) must be what escapes");
assertTrue(sessions.called.contains("onTurnComplete"), "the session half must still have run");
assertEquals(1, thrown.getSuppressed().length,
"the session half's distinct failure must be recorded as suppressed, not dropped");
assertEquals(IllegalStateException.class, thrown.getSuppressed()[0].getClass());
assertEquals("boom: sessions.onTurnComplete", thrown.getSuppressed()[0].getMessage());
}
}
@@ -720,10 +720,13 @@ class InjectorTest {
}
@Test
void anErrorFromSendRemovesTheMessageAndMarksItNotDelivered() {
void anErrorFromSendRemovesTheMessageAndMarksItAttempted() {
// fleetd #546, acceptance test 1: an Error (not a RuntimeException) escaping the send seam
// at Injector.java:383 must still be caught, the message dropped from the queue, and its
// state set to NOT_DELIVERED — never left QUEUED at the head.
// must still be caught and the message dropped from the queue, never left QUEUED at the
// head. fleetd #551 updates what it is marked: an AssertionError from the send call itself
// proves nothing about whether the text reached the pane, so it is recorded as ATTEMPTED
// (uncertain), not a confident NOT_DELIVERED — see anErrorFromSendMarksItNotDeliveredOnlyForANotFoundCode
// below for the one case that still gets NOT_DELIVERED.
ErrorOnPrompt throwing = new ErrorOnPrompt(new FakeHerdr());
Injector inj = new Injector(new AgentControl(throwing));
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
@@ -734,9 +737,9 @@ class InjectorTest {
assertTrue(delivery.completion().isCompletedExceptionally(),
"the delivery's future must surface the send failure");
assertEquals(Injector.Cancellation.NOT_DELIVERED, inj.cancel(delivery),
"the message must be dropped and marked NOT_DELIVERED, not left QUEUED at the "
+ "head of the queue");
assertEquals(Injector.Cancellation.ATTEMPTED, inj.cancel(delivery),
"fleetd #551: the message must be dropped and marked ATTEMPTED (not a confident "
+ "NOT_DELIVERED), and never left QUEUED at the head of the queue");
}
@Test
@@ -758,10 +761,14 @@ class InjectorTest {
}
@Test
void aHerdrExceptionFromSendStillProducesNotDeliveredUnchanged() {
void aHerdrExceptionFromSendStillSurfacesButNowReportsAttempted() {
// fleetd #546, acceptance test 3 (control): HerdrException extends RuntimeException, so it
// was already caught before this ticket's widening. This pins that the ordinary path is
// unchanged — still dropped, still NOT_DELIVERED, still surfaced to the caller.
// was already caught before #546's widening. The send-failure/surface behaviour is
// unchanged by #551 — still dropped, still surfaced to the caller — but fleetd #551
// deliberately changes WHAT it is marked: "send_failed" is not a herdr `*_not_found` code,
// so this is exactly the response-half case #551 exists to fix. See
// aHerdrExceptionAfterThePasteIsNeverRecordedAsConfidentlyNotDelivered for the acceptance
// test this ticket was filed for.
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
Injector inj = new Injector(new AgentControl(failing));
Injector.Delivery delivery = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
@@ -771,9 +778,90 @@ class InjectorTest {
assertTrue(delivery.completion().isCompletedExceptionally(),
"a HerdrException at the send seam must still surface to the caller, unchanged by "
+ "fleetd #546's widening");
assertEquals(Injector.Cancellation.ATTEMPTED, inj.cancel(delivery),
"fleetd #551: a non-*_not_found HerdrException must now report ATTEMPTED, not a "
+ "confident NOT_DELIVERED");
}
// --- fleetd #551: the injector records the delivery attempt BEFORE the irreversible send, not
// after, so a failure in the send's response window can never be written down as a confident
// NOT_DELIVERED for text that may already be sitting in the worker's pane. ---
@Test
void aHerdrExceptionAfterThePasteIsNeverRecordedAsConfidentlyNotDelivered() {
// fleetd #551, the ticket's acceptance test 1, and comment 17037's added test ("a test
// pinning that a caller cannot receive a confident NOT_DELIVERED for a message that reached
// the paste"): a HerdrException thrown from the response half of the send seam — herdr
// replied, with an error, so it definitely processed the request — must not leave a record
// claiming the text was never delivered. Red before the fix: on main at ba2f4d1 this
// asserted (and got) Injector.Cancellation.NOT_DELIVERED.
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
Injector inj = new Injector(new AgentControl(failing));
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
assertTrue(delivery.completion().isCompletedExceptionally(),
"the caller must still see the send failure");
assertEquals(Injector.Cancellation.ATTEMPTED, inj.cancel(delivery),
"fleetd #551: a HerdrException from the response half of the send seam must not be "
+ "recorded as a confident NOT_DELIVERED — the text may already be sitting "
+ "in the pane");
}
@Test
void ordinarySuccessStillReportsDeliveredExactlyOnce() {
// fleetd #551, the ticket's acceptance test 2: the record-before-send reorder must not
// change the ordinary success path — exactly one send, and the caller-facing Cancellation
// for it is still DELIVERED, never left at the new ATTEMPTED value.
Injector.Delivery delivery = injector.enqueue(T, "hello", TestTurnTokens.inert(T));
injector.onStatus(T, AgentStatus.IDLE);
assertEquals(List.of("hello"), sent(), "the text must be sent exactly once");
assertTrue(delivery.completion().isDone() && !delivery.completion().isCompletedExceptionally(),
"the ordinary success path must still complete normally");
assertEquals(Injector.Cancellation.DELIVERED, injector.cancel(delivery),
"fleetd #551: the ordinary success path must still report DELIVERED, unaffected by "
+ "the record-before-send reorder");
}
@Test
void transportDownStillSurfacesTheErrorAndDropsTheEntry() {
// fleetd #551, the ticket's acceptance test 3: the ordinary transport-down path (herdr
// unreachable — a plain HerdrException with no error code) must be unaffected by the #551
// reorder — the failure still surfaces to the caller, and the entry is not left sitting in
// the queue for a second onStatus round to resend.
FakeHerdr unreachable = new FakeHerdr().healthy(false);
Injector inj = new Injector(new AgentControl(unreachable));
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
assertTrue(delivery.completion().isCompletedExceptionally(),
"a transport-down failure must still surface to the caller");
assertTrue(inj.activeTargets().isEmpty(),
"the entry must not be left queued for a second round to resend");
}
@Test
void aNotFoundCodeFromSendStillReportsAConfidentNotDelivered() {
// fleetd #551 (comment 17037, point 3): a herdr `*_not_found` error means the target
// pane/agent does not exist at all, so nothing could have been pasted anywhere — the
// codebase already treats this family as a confirmed absence everywhere else (StatusPoller,
// AgentControl's own retry). #551's new ATTEMPTED state must not swallow this case: it stays
// a confident NOT_DELIVERED.
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("agent_not_found");
Injector inj = new Injector(new AgentControl(failing));
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
assertTrue(delivery.completion().isCompletedExceptionally(),
"an agent_not_found failure must still surface to the caller");
assertEquals(Injector.Cancellation.NOT_DELIVERED, inj.cancel(delivery),
"a HerdrException must still be dropped and marked NOT_DELIVERED, unchanged by the "
+ "wider Throwable catch");
"fleetd #551: a herdr *_not_found error is a confirmed absence, not merely "
+ "inconclusive — it must stay NOT_DELIVERED, not the new ATTEMPTED");
}
// --- fleetd #553: a throwable from any listener callback in onStatus must not skip the
@@ -8,6 +8,7 @@ import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
@@ -1053,6 +1054,36 @@ class FleetMcpTest {
assertFalse(out.contains("quarantinedForSeconds"), out);
}
@Test
void loopHealthReportsStalledStatusPoller() {
String out = loopHealth(LoopWatchdog.State.STALLED, LoopWatchdog.State.RUNNING);
assertTrue(out.contains("\"statusPoller\":\"STALLED\""),
"fleet_list must report a stalled StatusPoller: " + out);
}
@Test
void loopHealthReportsStoppedSessionReaperAsStopped() {
String out = loopHealth(LoopWatchdog.State.RUNNING, LoopWatchdog.State.STOPPED);
assertTrue(out.contains("\"sessionReaper\":\"STOPPED\""),
"fleet_list must report a deliberately stopped SessionReaper as STOPPED, not an alarm: " + out);
}
@Test
void loopHealthReportsRunningStatusPoller() {
String out = loopHealth(LoopWatchdog.State.RUNNING, LoopWatchdog.State.STOPPED);
assertTrue(out.contains("\"statusPoller\":\"RUNNING\""),
"fleet_list must report a running StatusPoller: " + out);
}
private static String loopHealth(LoopWatchdog.State statusPoller, LoopWatchdog.State sessionReaper) {
FakeHerdr herdr = new FakeHerdr();
return textOf(FleetMcp.listFleet(workerService(herdr, "http://gx00.gw:8000", Set.of("gx00.gw")),
new SessionManager(workerService(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"))), null,
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
new FleetMcp.LoopHealthSource(() -> statusPoller, () -> sessionReaper),
FleetMcp.QuarantineSource.none(), Map.of(), ""));
}
@Test
void capacityIncludesConfiguredProfileWithoutMembers() {
FakeHerdr h = new FakeHerdr();
@@ -19,6 +19,7 @@ import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -207,6 +208,31 @@ class LeadMailboxTest {
}
}
/**
* {@link LeadMailbox#inspect} opens a third channel after the mailbox's consume and publish
* channels. Limit this connection to three channels, then require a replacement channel after
* the successful inspect. If inspect leaves its probe open, the broker refuses that replacement.
*/
@Test
void inspectClosesItsSuccessfulProbeChannel() throws Exception {
var factory = LeadMailbox.connectionFactory(uri());
factory.setRequestedChannelMax(3);
Connection connection = factory.newConnection();
try (LeadMailbox mailbox = new LeadMailbox(connection, coordId("lead-inspect-probe-close"))) {
LeadChannel.MailboxState state = mailbox.inspect(mailbox.selfCoordId());
assertTrue(state.exists(), "the owned mailbox must be found before checking the probe channel");
Channel replacement = connection.createChannel();
assertNotNull(replacement,
"inspect must close its successful probe channel; the replacement channel was null");
try {
assertTrue(replacement.isOpen(), "the replacement channel must be open after inspect returns");
} finally {
replacement.close();
}
}
}
/**
* fleetd #440: {@code heldDurable()} must be derived from what {@link LeadMailbox#own} actually
* did against the real broker — a durable queue declare plus a manual-ack consumer — not a
@@ -506,6 +506,52 @@ class MessageServiceTest {
"an answer to a turn that never existed (or already lapsed) is stale, not a hang");
}
/**
* fleetd #575. {@code answer()}'s STALE_TURN return for "the ask lapsed between the lookup and
* the unblock" ({@code rendezvous.answerAsk(turnId, content)} returning {@code false} even though
* this call's own {@code rendezvous.askSession(turnId)} check at the top saw the ask as open) used
* to be covered only by a hand-rolled copy of the cleanup pair its own finally already runs, not
* the finally itself — a structural gap, closed by widening the try up to cover the registration
* above it. This test pins the exact race with {@link
* MessageService#setAnswerAskLapseRaceHookForTest}, which fires right before this call's own
* {@code rendezvous.answerAsk} call and completes the same turnId's ask directly — reproducing
* what a second, concurrent {@code answer()} winning that race could otherwise only do by timing
* luck. Proves the fix changed no behaviour on this path: it still returns {@code STALE_TURN},
* and this call's own forward waiter is still closed exactly once (not left open, and not closed
* twice — there is now only one cleanup site left to run).
*/
@Test
void answerLosingTheRaceToAnAlreadyAnsweredAskStillReturnsStaleTurnAndCleansUpOnce() throws Exception {
String ticket = messages.sendAsync(T, "long task");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
String turnId = asking.turnId();
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
messages.setAnswerAskLapseRaceHookForTest(() -> rendezvous.answerAsk(turnId, "raced in first"));
try {
MessageService.Reply r = messages.answer(turnId, "too late", 500);
assertEquals(MessageService.Outcome.STALE_TURN, r.outcome(),
"an ask already answered by the race must be seen as lapsed, not double-delivered");
} finally {
messages.setAnswerAskLapseRaceHookForTest(null);
}
// Cleanup ran exactly once: the forward waiter THIS call opened is closed, not leaked.
assertNull(rendezvous.currentWaiter(T),
"the forward waiter this answer() call opened must be closed after a STALE_TURN return");
// The worker's own ask() call, unblocked by the hook's direct answerAsk, still completes
// normally — the race this test simulates does not strand it.
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.AskOutcome.ANSWERED, a.outcome());
assertEquals("raced in first", a.answer());
}
// --- timeout, answer, poll, and lock-contention edges ----------------------------------
@Test
@@ -8,6 +8,7 @@ import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
@@ -165,6 +166,29 @@ class FleetAppTest {
assertEquals("degraded", mapper.readTree(res.body()).get("status").asText());
}
@Test
void healthzKeepsOkStatusAndReportsLoopHealthInItsBody() {
FleetMcp.LoopHealthSource loops = new FleetMcp.LoopHealthSource(
() -> LoopWatchdog.State.RUNNING, () -> LoopWatchdog.State.STOPPED);
FleetApp.HealthzResponse ok = FleetApp.healthzResponse(new FakeHerdr(), new FakeHerdr(), loops);
assertEquals(200, ok.status(), "a healthy herdr must keep /healthz at 200 regardless of loop states");
assertEquals(Map.of("statusPoller", "RUNNING", "sessionReaper", "STOPPED"), ok.body().get("loopHealth"),
"the /healthz body must report each loop state without making STOPPED an alarm");
}
@Test
void healthzKeepsDegradedStatusAndReportsLoopHealthInItsBody() {
FleetMcp.LoopHealthSource loops = new FleetMcp.LoopHealthSource(
() -> LoopWatchdog.State.RUNNING, () -> LoopWatchdog.State.STOPPED);
FleetApp.HealthzResponse degraded = FleetApp.healthzResponse(new FakeHerdr().healthy(false),
new FakeHerdr(), loops);
assertEquals(503, degraded.status(),
"an unreachable herdr must keep /healthz at 503 regardless of loop states");
assertEquals(Map.of("statusPoller", "RUNNING", "sessionReaper", "STOPPED"),
degraded.body().get("loopHealth"), "the degraded /healthz body must retain loop states");
}
@Test
void sessionsMapsWorkspaceList() throws Exception {
int port = startHealthy();