Compare commits

...

24 Commits

Author SHA1 Message Date
Dai Ha 2e349139e9 #568: add hunter member role
CI / shell-tests (pull_request) Successful in 6s
CI / contract (pull_request) Successful in 1m21s
CI / build (pull_request) Successful in 1m44s
2026-09-19 15:09:07 +07:00
Dai Ha a639969a9a CLAUDE.md: a blocked lead consults architects, not the operator
CI / shell-tests (push) Successful in 6s
CI / contract (push) Successful in 1m32s
CI / build (push) Successful in 1m42s
The operator set this rule on 2026-09-19: when a decision blocks a lead, it
consults one or more architect members, who are authorized to agree on one
decision and unblock. The operator is not asked. Escalation stays open only
for things outside the fleet's authority -- money, credentials, or a promise
made to someone else.

The paragraph also carries the reason the ticket record is mandatory rather
than optional. The operator's old notification channel was the block itself:
work stopped, so they found out. Taking the operator out of the loop removes
that signal with it, so the decision goes on the ticket, which reaches them
whether or not they are at a terminal when it is made.

The rule has a second half aimed at architects, which lives in
fleet.charters.architect and is applied per daemon -- filed as #591, because
a charter can never reach a lead and charters do not travel between hosts.
The missing notification event is #592.

Block verified byte-identical with wiki 7-Use-Cases.md at 1d9bd1b.
2026-09-19 14:57:04 +07:00
Dai Ha 17c3a69c57 docs: CB-591 gateway page had two claims that went stale
CI / shell-tests (push) Successful in 5s
CI / contract (push) Successful in 1m23s
CI / build (push) Failing after 1m37s
The gateway's chat model is served under the stable alias `acoder`, and the
model behind that alias changed on 2026-08-28 — it is Qwen3.8-27B now, not
DeepSeek-V4-Flash. The old name is still served, so nothing broke, but it
names a model this is not.

Two claims on the page were wrong as a result, and both were written as
current facts rather than dated measurements:

- `/v1/models` returns exactly `["deepseek-v4-flash"]` — it returns 6 ids
  now. This sat under a heading saying it needs no re-testing.
- the status banner said `local` and `gx` are both at `weight: 100` —
  `local` is at 0.

Measured today against the live gateway: /v1/models returns acoder,
qwen3.8-27b-nvfp4, deepseek-v4-flash and three embedding names; a completion
sent as `deepseek-v4-flash` comes back reporting `"model": "acoder"`, which
is the alias in plain sight. /v1/deployment reports generation
2026-08-28-qwen3.8-27b-nvfp4.

§2 and §3 are left alone. They are the August plan, and rewriting them would
destroy the record of the migration.

fleetd.yaml moved to `acoder` in the same change. It is not tracked here.
2026-09-13 07:01:15 +07:00
ltms 49a5875586 Merge #583: fleetd #582 — assert pending message-id cleanup at every publish cleanup site
CI / shell-tests (push) Successful in 9s
CI / contract (push) Successful in 48s
CI / build (push) Failing after 2m1s
All ten assertions proven live: six by the implementer, the last four by the lead.
One contract build with four deleted removal lines produced exactly four named failures,
one per site, with the total unchanged at 1825.
2026-09-12 15:52:25 +02:00
Dai Ha 634d33b50b Merge worker/562-loop-health-wiring-test-99611c-5
CI / shell-tests (push) Successful in 9s
CI / contract (push) Successful in 1m17s
CI / build (push) Successful in 1m45s
2026-09-12 20:28:14 +07:00
Dai Ha 1db79bcaa9 Merge worker/581-completionresolver-cas-sites-0542b7-6 2026-09-12 20:28:14 +07:00
Dai Ha 4507bc5a70 Merge worker/571-attempted-outcome-5739f7-2 2026-09-12 20:28:14 +07:00
Dai Ha d7239ed23b fleetd #571: pin FleetMcp.formatReply's TIMED_OUT_UNCONFIRMED wording
CI / shell-tests (pull_request) Successful in 5s
CI / contract (pull_request) Successful in 1m14s
CI / build (pull_request) Successful in 1m42s
CORRECTION 5 on the ticket: mutating the new arm's message text to the
queued/working arm's text survived every existing test, because nothing
asserted the specific wording. This adds one test that asserts the
unconfirmed-delivery message and asserts it does NOT carry the
queued/working arm's retry invitation — the distinction #571 exists for.

No production code changes; formatReply's TIMED_OUT_UNCONFIRMED arm was
already correct.
2026-09-12 20:19:37 +07:00
Dai Ha dfeb9340b4 fleetd #581: cover completion CAS removals
CI / shell-tests (pull_request) Successful in 5s
CI / build (pull_request) Failing after 1m31s
CI / contract (pull_request) Successful in 1m35s
2026-09-12 20:19:20 +07:00
Dai Ha 1513d4f260 fleetd #562 follow-up: extract loopHealthSource factory, pin its wiring
CI / shell-tests (pull_request) Successful in 6s
CI / contract (pull_request) Successful in 59s
CI / build (pull_request) Successful in 1m43s
PR #579's inline `new FleetMcp.LoopHealthSource(poller::health, ...)` in
Fleetd.main had nothing a test could call directly. Measured: replacing
poller::health with a constant () -> RUNNING compiled clean and left all
1771 tests green (see issue #562 comment "HOLD on PR #579").

Extracts the inline construction to a package-private Fleetd.loopHealthSource
factory, the same style as the sibling capacitySource/healthCoverageSource
factories, and adds FleetdLoopHealthSourceWiringTest with three separate
assertions: the statusPoller half, the sessionReaper half, and the
reaper == null branch (still STOPPED).
2026-09-12 20:16:23 +07:00
Dai Ha 4ca7d72303 #582: assert pending message-id cleanup
CI / shell-tests (pull_request) Successful in 5s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 2m32s
2026-09-12 20:15:48 +07:00
Dai Ha c0545d003d fleetd #571: make FleetApp.writeReply's inner Outcome switch exhaustive, no default
CI / shell-tests (pull_request) Successful in 8s
CI / contract (pull_request) Successful in 1m16s
CI / build (pull_request) Successful in 1m48s
Ticket comments (17126, 17127) corrected the original acceptance criterion after this
unit was already in flight: a hand-listed grep for the enum's constant names goes stale
silently the moment a new constant lands, so the compiler must be the enumeration instead.
sendOutcomeLabel (MessageService.java) and formatReply (FleetMcp.java) were already
default-free switch expressions. The one gap was writeReply's inner "status" switch, which
had `default -> "done"` — the exact value that would have lied about TIMED_OUT_UNCONFIRMED.
Remove the default and list every Outcome constant explicitly; REPLIED, COMPLETED_UNREPLIED,
QUESTION and STALE_TURN get an arm too even though the outer switch always dispatches them
first, so the inner switch stays exhaustive on its own. The outer switch (a statement, not
an expression) keeps its own default — Java does not require exhaustiveness there regardless,
and "everything not terminal is a 202" is an intentional catch-all.

Verified with the proof the ticket asked for: added a scratch 11th Outcome constant after
deleting all default arms and confirmed all three switch-expression sites (and no test file)
fail to compile without an arm for it, one at a time, then removed the scratch constant.
2026-09-12 19:28:59 +07:00
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 c1e06c9e12 fleetd #571: add TIMED_OUT_UNCONFIRMED so an ATTEMPTED delivery is not reported as never-arriving
MessageService.send's TimeoutException branch collapsed Injector.Cancellation.ATTEMPTED
(fleetd #551 — the send call was made but its outcome is unknown) into
Outcome.TIMED_OUT_QUEUED, which promises the caller the message will never arrive. On this
route agent.prompt may already have pasted and submitted the text, so a caller's natural
recovery (resend) risks a double delivery.

Add Outcome.TIMED_OUT_UNCONFIRMED and route ATTEMPTED to it. Update the three readers found
by searching for the enum's constant names (not `Outcome.`, which misses FleetMcp's
unqualified `case REPLIED ->` switches and would false-positive on ConfigRef's unrelated
Outcome record):
 - MessageService.sendOutcomeLabel: add it to the "timeout" metric label group.
 - FleetMcp.formatReply: its own case, warning against a blind retry (distinct from the
   generic "retry or poll status" message the other timeouts get).
 - FleetApp.writeReply: its own "unconfirmed" status and detail text, so it no longer falls
   through the switch's default -> "done" arm, which would have reported "the delegation
   completed" for the one case where delivery is unconfirmed.
2026-09-12 19:21:10 +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
Dai Ha a4dbc8f8b7 fleetd #572: pin answer()'s session-lock release across all four exits
CI / shell-tests (pull_request) Successful in 7s
CI / contract (pull_request) Successful in 47s
CI / build (pull_request) Successful in 2m10s
MessageService.answer() releases its per-session lock in an outer finally
(MessageService.java:1218) that mutation testing showed was covered but
unasserted: removing that line left all 1750 existing tests green, because
every existing test on this path checks answer()'s return value, never that
the lock it took is actually reacquirable afterward. If it leaked, a session
would be wedged forever with no exception and no log line.

Adds four tests, one per exit of answer() (normal REPLIED reply,
TIMED_OUT_WORKING, ExecutionException rethrow, InterruptedException
rethrow), each proving the lock is reacquirable via a bounded (300ms)
follow-up send on the same session rather than merely checking answer()'s
own outcome. No production change.
2026-09-12 18:23:54 +07: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 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 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
25 changed files with 1879 additions and 175 deletions
+23
View File
@@ -0,0 +1,23 @@
---
name: hunter
description: Sweep one assigned scope for defects and report ranked findings without changes.
---
<!-- CB-617: The model comes from fleetd.yaml because the launch flag overrides model here on both backends. -->
You sweep the assigned package or scope for real defects. Read the full assigned scope before you
judge it. Report several ranked findings when the evidence supports them. Change nothing: do not
edit code, commit, push, or open a pull request.
You may run the build or tests to check a finding. Read the complete output and report the real
result. Do not hide failures with a pipe. State only checks you actually ran. The primary's IDE
tools are not yours. A mounted forge tool may use a blocked credential and fail by design.
Do only the assigned scope. Note anything outside it in one line and do not investigate it further.
Use `fleet_ask{question}` only when a decision belongs to the lead, such as an unclear requirement
or two defensible fixes. Do not ask about something you can decide by reading more code.
Your handoff must name the files you read, each ranked finding or `NO FINDINGS`, the checks you ran,
and any caveat for review.
The launcher provides the required bridge reply instructions for every member.
+12 -1
View File
@@ -139,12 +139,21 @@ prefer `wait:false` + `fleet_poll` for anything non-trivial: a blocking `fleet_s
**Delegating does not delegate responsibility.** Workers open PRs; you are the gate. Never delegate
the merge — and merging on a reviewer's word is delegating it by proxy.
**When a decision blocks you, consult architects — not the operator.** Spawn one or more architect
members, give them the question and the evidence you have, and act on what they agree. They are
authorized to settle it, not only to advise. If two of them still disagree after two rounds, they
return both positions and you decide. Go to the operator only for something outside the fleet's
authority: money, credentials, or a promise made to someone else. **Then write the decision on the
ticket.** Taking the operator out of the loop also removes the signal they used to get, because
that signal was the block itself — work stopped, so they found out. A ticket comment replaces it,
and it reaches them whether or not they are at a terminal when you decide.
| Intent | Tool |
|---|---|
| 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` |
@@ -275,6 +284,8 @@ must obey belongs in the charter, not here.
it for a multi-finding sweep hands the worker two contradictory output contracts. That has
already cost three workers' turns: each wrote a good report to its terminal and ended the turn
with no `fleet_reply`, and the scrape returned the tail of the brief instead.
Spawn `implementer` with role `dev`, `reviewer` with role `reviewer`, and `hunter` with role
`hunter`.
- **Primary-side skills** (not delegation playbooks — a worker cannot use them):
`port-to-opencode` (make an OpenCode session a participant in this workspace),
`fleets-status` (report every fleet that shares one LavinMQ instance),
+19 -1
View File
@@ -5,6 +5,21 @@
here took a revert and two upstream fixes — see §7.1, which is the useful part of this document. One
risk is **accepted rather than solved**: a stream cut by any mid-response timer arrives as HTTP 200
with no terminator, and our third-party members cannot detect it (§7.2).
> **Superseded in part — 2026-09-13.** Two claims on this page are no longer true of the live fleet.
> I measured both on this host today.
>
> 1. **The model is named `acoder` now, not `deepseek-v4-flash`.** `acoder` is a stable alias, and
> the model behind it changed on 2026-08-28: it is Qwen3.8-27B, not DeepSeek. The old name is
> still served, so nothing broke — the gateway answers it and reports `"model": "acoder"` in the
> reply, which is how you can see for yourself that it is an alias. `fleetd.yaml` moved to
> `acoder` on 2026-09-13. Do not guess behaviour from the name; ask the gateway's own manifest,
> `GET https://llm.ltms.dev/v1/deployment`, and read its `generation` field.
> 2. **`local` sits at `weight: 0`, not 100.** Only `gx` is auto-selected today.
>
> §2 and §3 below are the plan as written in August. They are the record of the migration, so they
> stay as they are. If this note stops matching `fleetd.yaml`, re-measure and rewrite the note.
· **Upstream:** [systems/vms wiki → LLM and MCP Gateway](https://git.ltms.dev/systems/vms/wiki/LLM-and-MCP-Gateway)
· **Upstream issue:** [systems/vms#31](https://git.ltms.dev/systems/vms/issues/31)
@@ -380,7 +395,10 @@ one turn; this one costs the whole task and is indistinguishable from a slow wor
- Token accepted on both surfaces. **Unauthenticated → 401**, so the Caddy proxy really does gate —
the wiki's "SecurityPolicy fails open" warning is about the gateway itself, not the edge.
- `/v1/models` returns exactly `["deepseek-v4-flash"]`, so trap 3 is clear.
- `/v1/models` returned exactly `["deepseek-v4-flash"]` **on 2026-08-15**, so trap 3 was clear
then. It returns 6 ids now — `acoder`, `qwen3.8-27b-nvfp4`, `deepseek-v4-flash` and three
embedding names — measured on this host 2026-09-13. The exact-name rule still holds; the
one-item list does not.
- **Reasoning survives both surfaces** — see §3b above.
- The launcher's generated opencode provider block is correct, carrying a real 48-character `llmk-`
key rather than the `fleetd-local-noauth` placeholder.
+12 -6
View File
@@ -560,7 +560,7 @@ placement: weighted
# older keys: `leaders:`, `members:`, `leadScan:` and `defaultProfile:`.
#
# A member is anything a lead spawns, and every member has two INDEPENDENT attributes:
# role — which contract: architect, dev or reviewer. It picks the launch charter, the role
# role — which contract: architect, dev, hunter or reviewer. It picks the launch charter, the role
# file, the playbook skill and the authz row.
# profile — which backend: one of the `profiles:` keys above (model, CLI adapter, cost).
# They vary on their own. A reviewer may run on the same profile as the dev whose diff it reads,
@@ -572,13 +572,13 @@ placement: weighted
#
# Each pool lists the profiles that role MAY run on — these are pools, not identities. That is also
# what replaced `defaultProfile:`: an unqualified spawn names a role, and that role's pool supplies
# the candidates, in definition order. A dev and a reviewer staying anonymous is exactly compatible
# with being listed here; the entry key just names the entry.
# the candidates, in definition order. A dev, hunter and reviewer staying anonymous is exactly
# compatible with being listed here; the entry key just names the entry.
fleet:
# Optional launch-charter text, keyed only by the singular role wire names: architect, dev,
# reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put secrets
# here: a later launch step writes this text to a world-readable temp file, and ${ENV} interpolation
# is deliberately not supported.
# hunter, reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put
# secrets here: a later launch step writes this text to a world-readable temp file, and ${ENV}
# interpolation is deliberately not supported.
charters:
architect: |-
You are an architect in this fleet. You refine work before anyone builds it:
@@ -589,6 +589,9 @@ fleet:
dev: |-
You implement the one unit you were given, and nothing else. You test it,
commit it, and open your own pull request. You never merge.
hunter: |-
You sweep the assigned scope for real defects. You may run the build or tests
to check a finding. You change nothing, and report several ranked findings.
reviewer: |-
You review the diff you were given. You report bugs, risks and missing tests.
You do not change code.
@@ -666,6 +669,9 @@ fleet:
developers:
gx10:
profile: gx10
# hunters:
# gx10:
# profile: gx10 # a hunt may run checks, but never changes code
# reviewers:
# gx10:
# profile: gx10 # the same backend may serve two roles; that is the point
+185 -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,12 @@ public final class Fleetd {
return configured == null ? null : configured.effectiveCredentialId();
}, outagePolicy);
FleetMcp.LoopHealthSource loopHealth = loopHealthSource(poller, reaper);
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 +763,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);
@@ -1069,6 +1041,31 @@ public final class Fleetd {
});
}
/**
* fleetd #562 follow-up: package-private factory for {@code fleet_list}'s and {@code
* /healthz}'s {@code loopHealth} source, extracted out of {@code main} for the same reason
* {@link #capacitySource} and {@link #healthCoverageSource} were. Before this ticket the
* {@link FleetMcp.LoopHealthSource} was built inline with a bare {@code new}, so there was
* nothing a test could call directly — measured: replacing {@code poller::health} with a
* constant {@code () -> LoopWatchdog.State.RUNNING} at the call site compiled clean and left
* the full suite green, meaning the daemon could report the {@link StatusPoller} as always
* {@code RUNNING} even while it was actually stalled. That is a false negative on the exact
* signal this ticket exists to surface, and is the mirror of a false positive muting a real
* monitoring component — worse, because there is no noise for anyone to notice and then
* silence. {@link FleetdLoopHealthSourceWiringTest} calls this factory directly and pins both
* halves separately, plus the {@code reaper == null} branch below.
*
* <p>{@code reaper} may be {@code null} — a {@link SessionReaper} is only constructed when
* {@code lifecycle.idleTtlSeconds} is configured (see the {@code reaper} local above) — and
* this factory preserves the existing behaviour of reporting {@link LoopWatchdog.State#STOPPED}
* in that case, rather than a {@code NullPointerException} on the first {@code fleet_list} or
* {@code /healthz} call.
*/
static FleetMcp.LoopHealthSource loopHealthSource(StatusPoller poller, SessionReaper reaper) {
return new FleetMcp.LoopHealthSource(poller::health,
() -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health());
}
/**
* fleetd #248: package-private factory for the member worktree/branch lookup {@link
* CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241).
@@ -1096,6 +1093,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
@@ -221,7 +221,7 @@ public final class CallerResolver {
// The config/live binding names this pane as an architect slot's own. Same
// unforgeable pane mapping; the live binding, never a request argument, decides.
// Check the slot role too: this defence in depth prevents a bad lifecycle bind from
// escalating a dev or reviewer into an architect. Checked before the worker fallback.
// escalating a dev, hunter, or reviewer into an architect. Checked before the worker fallback.
return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid());
}
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
@@ -42,7 +42,7 @@ public interface MemberLifecycle {
* Try to bind a newly spawned {@code terminal} into the role it was granted.
*
* @return the role this session actually holds: {@code role} unchanged for a role with no
* live slot-binding semantics (dev, reviewer), or when the bind succeeded; a fallback
* live slot-binding semantics (dev, hunter, reviewer), or when the bind succeeded; a fallback
* role — never {@code role} — when a slot-bound role (architect) could not be bound.
* Callers must record THIS value on the session, never the requested {@code role}, so
* a later roster read never reports a role the session does not hold (CB-619). In
@@ -20,7 +20,8 @@ import java.util.function.Supplier;
*
* <p>Two halves, split by who owns each:
* <ul>
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code reviewers}
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code hunters}/
* {@code reviewers}
* (see {@link #slots()}), each carrying the {@code profile} reference the spawn lifecycle
* reads when it stands the slot up. <strong>Live, since fleetd #424</strong>: {@link #live}
* re-reads {@code fleet:} on every call, through a supplier the same shape as
@@ -322,7 +323,7 @@ public final class MemberRegistry implements MemberLifecycle {
* CB-619 / fleetd #123: refuse an architect acquire before anything spawns when no configured
* slot carries {@code profile} — the config-gap case from the original defect report (a spawn
* asked for {@code role=architect, profile=sonnet}, and {@code fleet.architects} carried only
* {@code opus} and {@code sol}). A dev/reviewer acquire is always a no-op: those pools are
* {@code opus} and {@code sol}). A dev/hunter/reviewer acquire is always a no-op: those pools are
* placement candidates only (see {@code CompositePeerLauncher}), never a live identity binding,
* so there is nothing here to refuse — an explicit profile outside the pool for those roles is a
* documented operator override, not a defect.
@@ -31,7 +31,7 @@ import java.util.function.Supplier;
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Both are
* read through a supplier on {@code CompositePeerLauncher}, which is what makes them hot —
* not the fact that they are config. Most of {@code fleet:} — every role pool
* ({@code architects}/{@code developers}/{@code reviewers}), {@code charters}, and
* ({@code architects}/{@code developers}/{@code hunters}/{@code reviewers}), {@code charters}, and
* {@code tabLabel} — is read the same live way, through the same supplier
* ({@code () -> config.get().fleet()}). {@code architects} in particular is hot for
* <strong>two independent consumers</strong> (fleetd #424): {@code CompositePeerLauncher}
@@ -605,7 +605,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "opened once and needs a restart; the broker URI env-var name kept out of a "
+ "member's environment is read live on every spawn and already applied");
}
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, reviewers,
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, hunters, reviewers,
// charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
// ThroughGet and aCharterChangeIsHotAndReachesTheLiveConfig prove it reaches the live config
// with no restart note. `architects` is hot too, and — since fleetd #424 — hot for BOTH of
@@ -630,7 +630,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
+ "identity map and to auto-launch leads, and neither is rebuilt on reload, so a "
+ "lead added, removed, or given a new tab: label needs a restart — until then it "
+ "stays unrecognised, and a caller from its new tab resolves as a worker, not a "
+ "lead; the rest of fleet: (developers, reviewers, charters, tabLabel) is read "
+ "lead; the rest of fleet: (developers, hunters, reviewers, charters, tabLabel) is read "
+ "live through the supplier on CompositePeerLauncher, and architects is read "
+ "live through that same supplier for placement AND through a separate supplier "
+ "on MemberRegistry for spawn-time identity — both already applied");
@@ -62,7 +62,7 @@ import java.util.regex.PatternSyntaxException;
* @param fleet who the daemon may run and under which role (CB-557). One block replacing
* the former {@code leaders:}, {@code members:}, {@code leadScan:} and
* {@code defaultProfile:}. Role is the containing key — {@code leaders},
* {@code architects}, {@code developers}, {@code reviewers} — and each entry
* {@code architects}, {@code developers}, {@code hunters}, {@code reviewers} — and each entry
* names the {@code profiles:} backend it runs on. See {@link Fleet}
* @param leadHeartbeat opt-in idle-lead heartbeat (CB-551); {@code null} ⇒ off, and an upgraded
* daemon never nudges an idle lead on its own initiative
@@ -1200,6 +1200,7 @@ public record FleetConfig(
* @param leaders panes that orchestrate rather than are orchestrated, keyed by lead name
* @param architects profiles the {@code architect} role may run on
* @param developers profiles the {@code dev} role may run on
* @param hunters profiles the {@code hunter} role may run on
* @param reviewers profiles the {@code reviewer} role may run on
* @param charters optional launch-charter text keyed by singular role wire name
* @param tabLabel template for a member tab's label; {@code {role}}, {@code {profile}},
@@ -1210,6 +1211,7 @@ public record FleetConfig(
public record Fleet(Map<String, Leader> leaders,
Map<String, Slot> architects,
Map<String, Slot> developers,
Map<String, Slot> hunters,
Map<String, Slot> reviewers,
Map<String, String> charters,
String tabLabel) {
@@ -1226,6 +1228,7 @@ public record FleetConfig(
leaders = unmodifiableOrEmpty(leaders);
architects = unmodifiableOrEmpty(architects);
developers = unmodifiableOrEmpty(developers);
hunters = unmodifiableOrEmpty(hunters);
reviewers = unmodifiableOrEmpty(reviewers);
charters = unmodifiableOrEmpty(charters);
tabLabel = (tabLabel == null || tabLabel.isBlank()) ? DEFAULT_TAB_LABEL : tabLabel;
@@ -1240,9 +1243,15 @@ public record FleetConfig(
* constructor: the launcher reads {@code fleet.charters()} from the live config. Jackson
* binds the canonical constructor, so this one cannot swallow an operator's YAML.
*/
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers,
Map<String, String> charters, String tabLabel) {
this(leaders, architects, developers, null, reviewers, charters, tabLabel);
}
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
Map<String, Slot> developers, Map<String, Slot> reviewers, String tabLabel) {
this(leaders, architects, developers, reviewers, null, tabLabel);
this(leaders, architects, developers, null, reviewers, null, tabLabel);
}
/**
@@ -1264,6 +1273,7 @@ public record FleetConfig(
return switch (role) {
case ARCHITECT -> architects;
case DEV -> developers;
case HUNTER -> hunters;
case REVIEWER -> reviewers;
};
}
@@ -1872,7 +1882,7 @@ public record FleetConfig(
/** The {@code fleet:} child blocks whose direct children are slot names. */
private static final Set<String> FLEET_POOL_KEYS =
Set.of("leaders", "architects", "developers", "reviewers");
Set.of("leaders", "architects", "developers", "hunters", "reviewers");
/**
* Reject a {@code fleet:} role pool whose slot names repeat (CB-548, re-homed by CB-557).
@@ -1882,7 +1892,7 @@ public record FleetConfig(
* daemon would never know. Jackson's YAML parser does not fail on duplicate mapping keys by
* default, so duplicates are caught here, at parse time, before the map is built.
*
* <p>Only the four pools <em>directly under the top-level {@code fleet:}</em> are considered,
* <p>Only the five pools <em>directly under the top-level {@code fleet:}</em> are considered,
* and only their direct child keys (the slot names). A nested field elsewhere, even one also
* named {@code developers:}, is ignored, so parsing of the rest of the config is unaffected.
*
@@ -2050,7 +2060,7 @@ public record FleetConfig(
+ " and that role's pool supplies the candidate profiles",
"architects", "'fleet.architects'",
"members", "a role pool under 'fleet:' — 'fleet.architects', 'fleet.developers' or"
+ " 'fleet.reviewers'; the role is the containing key, not a 'role:' field",
+ " 'fleet.hunters' or 'fleet.reviewers'; the role is the containing key, not a 'role:' field",
"leaders", "'fleet.leaders'",
"leadScan", "'fleet.leaders.<name>.tabPrefix' and '.scanIntervalSeconds' — lead"
+ " discovery is now configured on the lead it discovers");
@@ -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),
@@ -805,6 +828,12 @@ public final class FleetMcp {
+ "answered (turnId stale)");
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> text("[no reply within " + timeout + "ms — worker "
+ r.outcome().name().toLowerCase().replace("timed_out_", "") + "; retry or poll status]");
// fleetd #571: delivery is unknown here — agent.prompt pastes and submits in one call,
// so the message may already be sitting in the pane. Do not invite a blind retry the way
// the case above does; a resend on this route can double-deliver the same brief.
case TIMED_OUT_UNCONFIRMED -> text("[no reply within " + timeout + "ms — delivery unconfirmed; "
+ "the message may already have reached the worker, so a retry risks sending it "
+ "twice — poll status before resending]");
};
}
@@ -1546,25 +1575,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 +1626,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 +1635,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 +1657,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 +1677,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 +1712,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.
@@ -2075,8 +2134,9 @@ public final class FleetMcp {
return tool(FleetTool.SPAWN.wireName(),
"Spawn a new off-subscription member session. A member has two independent attributes: "
+ "role (what it is for) and profile (which backend it runs on). Pass role to pick "
+ "the contract — 'dev' implements a unit and opens its own PR, 'reviewer' reviews a "
+ "diff it did not write, 'architect' refines a ticket before anyone builds it; omit "
+ "the contract — 'dev' implements a unit and opens its own PR, 'hunter' sweeps a "
+ "scope without changing it, 'reviewer' reviews a diff it did not write, 'architect' "
+ "refines a ticket before anyone builds it; omit "
+ "it for 'dev'. Pass profile (from fleet_profiles) to pick the backend, or omit it "
+ "for the default. The two are independent: a reviewer may run on the same profile "
+ "as the dev it reviews. The member opens your current directory by default; pass "
@@ -2093,7 +2153,7 @@ public final class FleetMcp {
+ "one. Returns the member's sessionId (use with fleet_send) and paneId (use with "
+ "fleet_stop).",
objectSchema(Map.of(
"role", stringProp("What the member is for: architect, dev or reviewer (default dev)"),
"role", stringProp("What the member is for: architect, dev, hunter, or reviewer (default dev)"),
"profile", stringProp("Which backend to run it on (omit for the default profile)"),
"cwd", stringProp("Working directory for the member (omit to inherit yours)"),
"worktree", Map.of("type", "string", "description", "'true' or a ticket slug — requests an isolated git worktree"),
@@ -2135,9 +2195,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 "
@@ -93,25 +93,30 @@ public final class MessageService {
/** Timed out after the message was delivered — the worker is still working. */
TIMED_OUT_WORKING,
/**
* Timed out with no confirmed delivery. Despite the name, this does not mean the message
* 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 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 with no confirmed delivery, and the target saw nothing — the message will not
* arrive later, so a caller may resend. {@link #send} reaches this outcome through {@link
* Injector#cancel} reporting one of two routes: {@link Injector.Cancellation#CANCELLED}
* means the message was still queued and this call removed it; {@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. A third route,
* {@link Injector.Cancellation#ATTEMPTED}, used to be folded into this same outcome
* (fleetd #571) — it no longer is; see {@link #TIMED_OUT_UNCONFIRMED}.
*/
TIMED_OUT_QUEUED,
/**
* Timed out with delivery unknown. {@link #send} reaches this outcome when {@link
* Injector#cancel} reports {@link Injector.Cancellation#ATTEMPTED} (fleetd #551): the call
* to the target's terminal ({@link dev.ltms.fleet.herdr.AgentControl#send}) was made, but
* this caller never observed whether it reached the pane. {@code agent.prompt} pastes
* <em>and submits</em> in one call, so the target may already hold a complete, submitted
* turn and be working on it right now — the same reality as {@link #TIMED_OUT_WORKING},
* just not confirmed. The message may or may not have arrived. Treat this as neither a
* confirmed delivery nor a confirmed absence: a caller that resends on this outcome risks a
* double delivery — the same brief typed into the pane twice (fleetd #571).
*/
TIMED_OUT_UNCONFIRMED,
/** Another send to this session was in flight for the whole window. */
BUSY,
/**
@@ -317,20 +322,21 @@ public final class MessageService {
*/
private final ConcurrentHashMap<String, Boolean> strandedReplies = new ConcurrentHashMap<>();
/**
* 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 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}).
* Targets whose last send timed out with no confirmed delivery (CB-640) — {@link #send} called
* {@link Injector#cancel} and got back something other than {@code DELIVERED}. 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 — and the caller sees a different outcome for it
* ({@link Outcome#TIMED_OUT_UNCONFIRMED}, fleetd #571) than for the first two ({@link
* Outcome#TIMED_OUT_QUEUED}). 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();
@@ -417,21 +423,22 @@ public final class MessageService {
/**
* Read-only delegation fact for fleet views (CB-640): {@code target}'s last send timed out
* 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
* 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
* with no confirmed delivery — the caller saw {@link Outcome#TIMED_OUT_QUEUED} or {@link
* Outcome#TIMED_OUT_UNCONFIRMED} (fleetd #571; 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 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, and both report {@code TIMED_OUT_QUEUED}. {@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}.
* working on it right now — it does NOT follow that the target saw nothing, and this route
* reports {@code TIMED_OUT_UNCONFIRMED} instead. 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);
@@ -646,7 +653,7 @@ public final class MessageService {
return switch (o) {
case REPLIED -> "replied";
case COMPLETED_UNREPLIED -> "completion_fallback";
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, BUSY -> "timeout";
case TIMED_OUT_WORKING, TIMED_OUT_QUEUED, TIMED_OUT_UNCONFIRMED, BUSY -> "timeout";
case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
case STALE_TURN, QUESTION -> null; // not a completed delegation
@@ -975,6 +982,7 @@ public final class MessageService {
} catch (TimeoutException e) {
boolean wasDelivered = delivery.completion().isDone()
&& !delivery.completion().isCompletedExceptionally();
Injector.Cancellation cancellation = null;
if (!wasDelivered) {
if (timeoutCancellationRaceHookForTest != null) {
// Test-only (fleetd #345): see the field's own javadoc.
@@ -982,20 +990,27 @@ public final class MessageService {
}
// The target monitor makes cancellation atomic with onStatus picking this
// Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
cancellation = injector.cancel(delivery);
wasDelivered = cancellation == Injector.Cancellation.DELIVERED;
}
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
if (!wasDelivered) {
Outcome outcome;
if (wasDelivered) {
outcome = Outcome.TIMED_OUT_WORKING;
} else if (cancellation == Injector.Cancellation.ATTEMPTED) {
// fleetd #571: the call to the target's terminal was made and its outcome is
// unknown — the message may already have arrived in full, so this must not
// be reported as TIMED_OUT_QUEUED, which promises it never will.
outcome = Outcome.TIMED_OUT_UNCONFIRMED;
} else {
// 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 (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). cancellation is CANCELLED (this call removed a
// still-queued Pending) or NOT_DELIVERED (an earlier attempt already failed
// with a confirmed absence) — both mean the target saw nothing.
queuedDeliveries.put(target, Boolean.TRUE);
outcome = Outcome.TIMED_OUT_QUEUED;
}
return recorded(new Reply(
wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
return recorded(new Reply(outcome, null));
} catch (ExecutionException e) {
Throwable cause = e.getCause();
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
@@ -1135,21 +1150,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 —
@@ -1649,6 +1676,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));
@@ -29,8 +29,8 @@ public enum MemberRole {
* <p>Reads the repo and writes analysis. Never commits code and never opens a pull request —
* an architect that starts implementing has stopped doing the job that makes it useful.
*
* <p>Architects are the one member kind declared in config, because a lead addresses the same
* slots across many tickets and needs a stable name for them.
* <p>Architects are the one member kind with live slot binding, because a lead addresses the
* same slots across many tickets and needs a stable name for them.
*/
ARCHITECT,
@@ -43,6 +43,14 @@ public enum MemberRole {
*/
DEV,
/**
* Sweeps an assigned package for defects and reports several ranked findings.
*
* <p>Never changes code, commits, or opens a pull request. A hunt gathers evidence, which can
* include running the build, but leaves every fix to a later implementation unit.
*/
HUNTER,
/**
* Reviews a diff it did not write and reports one structured finding.
*
@@ -59,7 +67,7 @@ public enum MemberRole {
/**
* The {@code fleet:} block that holds this role's pool — {@code architects},
* {@code developers}, {@code reviewers}.
* {@code developers}, {@code hunters}, {@code reviewers}.
*
* <p>Plural, and not always the wire name: the pool of things a {@code dev} may run on reads
* naturally as {@code developers:}. The wire name stays the singular {@code dev}, because that
@@ -69,6 +77,7 @@ public enum MemberRole {
return switch (this) {
case ARCHITECT -> "architects";
case DEV -> "developers";
case HUNTER -> "hunters";
case REVIEWER -> "reviewers";
};
}
@@ -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)));
}
/**
@@ -640,15 +670,31 @@ public final class FleetApp {
}
default -> ctx.status(202).json(Map.of(
"sessionId", id,
// fleetd #571 (ticket comment 17126): no `default` here on purpose. This switch
// is an expression, so the compiler already demands every Outcome constant have
// an arm — adding an 11th constant to Outcome is a compile error here, not a
// silent fall-through. That is exactly the bug this ticket exists to fix:
// `default -> "done"` used to sit here and would have told a REST caller the
// delegation completed for TIMED_OUT_UNCONFIRMED, the one outcome where delivery
// is unknown. REPLIED, COMPLETED_UNREPLIED, QUESTION and STALE_TURN can never
// actually reach this inner switch — the outer switch above always dispatches
// them first — but they still need an arm to keep this switch exhaustive.
"status", switch (reply.outcome()) {
case TIMED_OUT_WORKING -> "working";
case TIMED_OUT_QUEUED -> "queued";
// Delivery here is unknown, not merely still queued — see
// Outcome#TIMED_OUT_UNCONFIRMED's own javadoc.
case TIMED_OUT_UNCONFIRMED -> "unconfirmed";
case BUSY -> "busy";
case WORKER_FAILED -> "failed";
case BACKEND_EXHAUSTED -> "backend_exhausted";
default -> "done"; // unreachable (terminal outcomes handled above)
case REPLIED, COMPLETED_UNREPLIED, QUESTION, STALE_TURN -> "done"; // unreachable
},
"detail", (reply.outcome() == MessageService.Outcome.WORKER_FAILED
"detail", reply.outcome() == MessageService.Outcome.TIMED_OUT_UNCONFIRMED
? "no reply within " + timeout + "ms; delivery is unconfirmed — the "
+ "message may already have reached the worker, so a resend "
+ "risks sending it twice; poll status first"
: (reply.outcome() == MessageService.Outcome.WORKER_FAILED
|| reply.outcome() == MessageService.Outcome.BACKEND_EXHAUSTED)
&& reply.text() != null
? reply.text()
@@ -0,0 +1,174 @@
package dev.ltms.fleet;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.SpawnRequest;
import dev.ltms.fleet.placement.PlacementDecision;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.SessionReaper;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Set;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* fleetd #562 follow-up (issue comment "HOLD on PR #579"): {@code Fleetd.main}'s {@code loopHealth}
* local used to be a bare {@code new FleetMcp.LoopHealthSource(poller::health, ...)} built inline,
* with nothing a test could call directly. Measured on that shape: replacing {@code
* poller::health} with a constant {@code () -> LoopWatchdog.State.RUNNING} at the call site
* compiled with 0 errors and left all 1771 existing tests green — the daemon could be changed to
* always report the {@link StatusPoller} as {@code RUNNING}, so the watchdog could never fire and
* a stalled poller would be invisible, while every test stayed green. That is exactly the false
* negative this ticket exists to prevent.
*
* <p>The five tests PR #579 added ({@code FleetMcpTest}, {@code FleetAppTest}) all build their own
* {@link FleetMcp.LoopHealthSource} directly with fixed lambdas — they prove the seam ({@code
* LoopHealthSource} reports what it is given) and nothing about what {@code Fleetd.main} actually
* gives it. This is the same hand-built-vs-config-wired shape as fleetd #561/#248/#426.
*
* <p>The fix extracts the inline {@code new} into {@link Fleetd#loopHealthSource}, a package-private
* factory in the same style as {@link Fleetd#capacitySource} and {@link Fleetd#healthCoverageSource}
* — which is exactly what makes it directly callable here. This test calls that factory with real
* {@link StatusPoller}/{@link SessionReaper} instances (never started, so no herdr or git I/O
* happens) and pins each half separately, plus the {@code reaper == null} branch: one invariant
* wired at three places needs three assertions, not one combined check whose non-zero total could
* hide a gap at any single place.
*/
class FleetdLoopHealthSourceWiringTest {
@Test
@DisplayName("the statusPoller half reports the real poller's health, not a hardcoded state")
void statusPollerHalfReflectsThePollersRealHealth() {
// Stopped without ever being started — stop() still marks the watchdog STOPPED. A poller
// that has never reported RUNNING is the discriminating case: if Fleetd.loopHealthSource
// ever hardcoded RUNNING (the exact mutation this test exists to catch), this would fail.
StatusPoller stoppedPoller = freshPoller();
stoppedPoller.stop();
SessionReaper unusedReaper = freshReaper(); // present only to satisfy the signature
FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(stoppedPoller, unusedReaper);
assertEquals(LoopWatchdog.State.STOPPED, source.statusPoller().get(),
"the statusPoller supplier must delegate to the real poller's health() — "
+ "replacing poller::health with a constant () -> RUNNING at the "
+ "Fleetd.loopHealthSource call site must fail this assertion");
}
@Test
@DisplayName("the sessionReaper half reports the real reaper's health, not a hardcoded state")
void sessionReaperHalfReflectsTheReapersRealHealth() {
StatusPoller unusedPoller = freshPoller(); // present only to satisfy the signature
SessionReaper stoppedReaper = freshReaper();
stoppedReaper.stop();
FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(unusedPoller, stoppedReaper);
assertEquals(LoopWatchdog.State.STOPPED, source.sessionReaper().get(),
"the sessionReaper supplier must delegate to the real reaper's health() — "
+ "replacing reaper.health() with a constant at the "
+ "Fleetd.loopHealthSource call site must fail this assertion");
}
@Test
@DisplayName("a null reaper (idle ttl not configured) still reports STOPPED, not a crash")
void nullReaperStillReportsStopped() {
// SessionReaper is only constructed when lifecycle.idleTtlSeconds is configured (see the
// `reaper` local in Fleetd.main) — a real deployment routinely passes null here. That null
// check is real behaviour, not a simplification to delete: it must keep reporting STOPPED
// rather than throwing a NullPointerException on the first fleet_list/healthz call.
StatusPoller runningPoller = freshPoller();
FleetMcp.LoopHealthSource source = Fleetd.loopHealthSource(runningPoller, null);
assertEquals(LoopWatchdog.State.STOPPED, source.sessionReaper().get(),
"reaper == null must still report STOPPED, exactly like an intentionally-stopped "
+ "reaper would — do not delete this null check to simplify the wiring");
}
/** Never started, so no herdr call is ever made; freshly constructed reports RUNNING. */
private static StatusPoller freshPoller() {
AgentControl agents = new AgentControl(new FakeHerdr());
return new StatusPoller(agents, new Injector(agents), 1000);
}
/** Never started, so no git/session I/O is ever made; freshly constructed reports RUNNING. */
private static SessionReaper freshReaper() {
return new SessionReaper(new SessionManager(new NeverSpawnsLauncher()), 60, 1000);
}
/**
* Same minimal shape as {@code FleetdBackendErrorSinkTest.NeverSpawnsLauncher} — every method
* throws or returns an empty/no-op value, since a {@link SessionReaper} that is only ever
* constructed and then stopped (never started) never calls any of them.
*/
private static final class NeverSpawnsLauncher implements PeerLauncher {
@Override
public Set<Capability> capabilities() {
return Set.of();
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return Set.of();
}
@Override
public PeerHandle spawn(SpawnRequest req) {
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
}
@Override
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
}
@Override
public Set<String> profiles() {
return Set.of();
}
@Override
public String defaultProfile() {
return null;
}
@Override
public String effectiveCwd(SpawnRequest req) {
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
}
@Override
public List<String> parityOverlay(String profileName) {
return List.of();
}
@Override
public List<?> list() {
return List.of();
}
@Override
public int reapOrphanWorkers() {
return 0;
}
@Override
public void stop(String id) {
}
@Override
public boolean clearContext(String id) {
return false;
}
}
}
@@ -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());
}
}
@@ -345,7 +345,7 @@ class FleetConfigTest {
IllegalStateException unknownError = assertThrows(IllegalStateException.class,
() -> FleetConfig.load(unknown).validateCharters());
assertTrue(unknownError.getMessage().contains("architetc"));
assertTrue(unknownError.getMessage().contains("[architect, dev, reviewer]"));
assertTrue(unknownError.getMessage().contains("[architect, dev, hunter, reviewer]"));
}
/**
@@ -812,13 +812,17 @@ class FleetConfigTest {
reviewers:
b:
profile: sonnet
hunters:
c:
profile: sonnet
""");
FleetConfig cfg = FleetConfig.load(f);
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.DEV));
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.HUNTER));
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.REVIEWER));
assertTrue(cfg.fleet().profilesFor(MemberRole.ARCHITECT).isEmpty());
assertEquals(List.of(MemberRole.DEV, MemberRole.REVIEWER), cfg.fleet().rolesConfigured());
assertEquals(List.of(MemberRole.DEV, MemberRole.HUNTER, MemberRole.REVIEWER), cfg.fleet().rolesConfigured());
}
/** The case the two axes exist for: one backend, two roles, and neither is a duplicate. */
@@ -625,6 +625,149 @@ class CompletionResolverTest {
assertEquals(Rendezvous.Kind.REPLY, waiterB.getNow(null).kind());
}
@Test
void aSupersededDoneTurnMustNotEvictItsSuccessorsRegistration() {
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(new FakeHerdr()), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
assertTrue(rendezvous.resolve("term_a", "A replied"));
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a done turn must not evict B from resolve()'s early return");
}
@Test
void aSupersededExhaustedTurnMustNotEvictItsSuccessorsRegistration() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ usage limit has been reached\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
target -> Pattern.compile("usage limit has been reached"), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"an exhausted turn must not evict B from resolve()'s exhausted branch");
}
@Test
void aSupersededBackendErrorTurnMustNotEvictItsSuccessorsRegistration() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ API Error: 400 invalid request body\n❯ ");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a backend-error turn must not evict B from resolve()'s error branch");
}
@Test
void aSupersededRawExhaustedTurnMustNotEvictItsSuccessorsRegistration() {
FakeHerdr herdr = new FakeHerdr().readText("╭────\nusage limit has been reached");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
target -> Pattern.compile("usage limit has been reached"), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a raw exhausted turn must not evict B from the raw-scrape exhausted branch");
}
@Test
void aSupersededRawBackendErrorTurnMustNotEvictItsSuccessorsRegistration() {
FakeHerdr herdr = new FakeHerdr().readText("╭────\nAPI Error: 400 invalid request body");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a raw backend-error turn must not evict B from the raw-scrape error branch");
}
@Test
void aSupersededDoneFailedTurnMustNotEvictItsSuccessorsRegistration() {
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(new FakeHerdr()), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none());
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null);
assertTrue(rendezvous.resolve("term_a", "A replied"));
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
resolver.fail("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a done failed turn must not evict B from fail()'s early return");
}
@Test
void aSupersededTooFastBackendErrorTurnMustNotEvictItsSuccessorsRegistration() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ API Error: 400 invalid request body\n❯ ");
Rendezvous rendezvous = new Rendezvous();
long[] clock = {10_000_000_000L};
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
var waiterA = rendezvous.open("term_a");
var turnA = new CompletionResolver.InFlight(waiterA, null, clock[0]);
rendezvous.close("term_a", waiterA);
var waiterB = rendezvous.open("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiterB));
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1;
resolver.resolve("term_a", turnA);
assertSuccessorRegistrationSurvives(resolver, rendezvous, waiterB,
"a too-fast backend-error turn must not evict B from failTooFast()");
}
private static void assertSuccessorRegistrationSurvives(CompletionResolver resolver, Rendezvous rendezvous,
Object waiterB, String message) {
CompletionResolver.InFlight afterA = resolver.inFlight("term_a");
assertNotNull(afterA, message + " — a one-arg remove(target) would remove B");
assertEquals(waiterB, afterA.waiter(), message + " — the surviving record must belong to B");
assertTrue(rendezvous.resolve("term_a", "B replied"), message + " — B must still resolve normally");
}
// --- CB-578 stage A: backend-exhausted classification ---------------------------------
@Test
@@ -7,7 +7,9 @@ import dev.ltms.fleet.auth.Role;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.Injector;
@@ -57,9 +59,10 @@ class FleetMcpTest {
private final FakeHerdr herdr = new FakeHerdr();
private final AgentControl agents = new AgentControl(herdr);
private final Injector injector = new Injector(agents);
private final Rendezvous rendezvous = new Rendezvous();
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
private final MessageService messages = new MessageService(agents, new Injector(agents), rendezvous, inbox);
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
@BeforeEach
void setUp() {
@@ -298,6 +301,35 @@ class FleetMcpTest {
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
}
/**
* fleetd #571 (ticket CORRECTION 5): {@code formatReply}'s {@code TIMED_OUT_UNCONFIRMED} arm is
* the one message whose whole job is to stop a caller retrying a delivery that may already have
* arrived. Pin that its wording is actually distinct from the queued/working arm's retry
* invitation — a mutation that swapped this arm's text for that one still passed every other
* test in this suite, because nothing asserted the specific wording.
*/
@Test
void sendTimesOutWithAnUnconfirmedNoteNotARetryInvitation() throws Exception {
herdr.agentSendFailsWith("send_failed");
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> FleetMcp.send(messages, T, "hi", 150L, null, Set.of()));
long deadline = System.currentTimeMillis() + 2000;
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt -> ATTEMPTED
McpSchema.CallToolResult res = send.get(5, TimeUnit.SECONDS);
String text = textOf(res);
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(text.contains("delivery unconfirmed"), "got: " + text);
assertFalse(text.contains("retry or poll status"),
"an unconfirmed delivery must not carry the queued/working arm's retry invitation — "
+ "a resend here can double-deliver the same brief: got " + text);
}
@Test
void sendRejectsMissingArgs() {
assertTrue(FleetMcp.send(messages, null, "hi", null, null, Set.of()).isError());
@@ -1053,6 +1085,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();
@@ -1875,7 +1937,7 @@ class FleetMcpTest {
null, null, null, null, null, null);
assertEquals(Boolean.TRUE, res.isError());
assertTrue(textOf(res).contains("architect, dev, reviewer"), textOf(res));
assertTrue(textOf(res).contains("architect, dev, hunter, reviewer"), textOf(res));
}
// ── CB-619 / fleetd #123: a spawn asking for a role its profile has no slot for must be
@@ -556,6 +556,31 @@ class ClaudeCodeLauncherTest {
"no --agent flag when the role has no agent-definition file");
}
@Test
void hunterRoleUsesItsAgentFileAndStopsUsingItWhenRemoved(@TempDir Path cwd) throws Exception {
FakeHerdr herdr = new FakeHerdr();
Path agentFile = Files.createDirectories(cwd.resolve(".claude/agents")).resolve("hunter.md");
Files.writeString(agentFile, "---\nname: hunter\n---\nSweep for defects.");
FleetConfig.Profile cfg = new FleetConfig.Profile(
"sonnet", "http://gx00.gw:8000", null, null, "FLEETD_WORKER_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
svc.spawn(new SpawnRequest("sonnet", cwd.toString(), null, null, null, MemberRole.HUNTER));
List<String> args = spawnedArgs(herdr);
int flag = args.indexOf("--agent");
assertTrue(flag >= 0, "the hunter role reaches its agent-definition file: " + args);
assertEquals("hunter", args.get(flag + 1));
Files.delete(agentFile);
svc.spawn(new SpawnRequest("sonnet", cwd.toString(), null, null, null, MemberRole.HUNTER));
assertFalse(spawnedArgs(herdr).contains("--agent"),
"the hunter role no longer gets an agent when its file is removed");
}
private ClaudeCodeLauncher multiProfile(FakeHerdr herdr) {
FleetConfig.Profile gx10 = new FleetConfig.Profile("gx10", "http://gx10.gw:8000", "coder",
null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
@@ -8,7 +8,10 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Method;
import java.lang.reflect.Proxy;
import java.io.IOException;
import java.util.Map;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
@@ -20,6 +23,7 @@ import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* CB-528 follow-up: {@link AmqpReplyInbox#failPendingPublishesOnRecovery()} must not fail a publish
@@ -197,11 +201,78 @@ class AmqpReplyInboxRecoveryRaceTest {
+ elapsedMillis.get() + "ms");
}
@Test
void publishIOExceptionRemovesThePendingMessageId() throws Exception {
AtomicLong seqCounter = new AtomicLong();
Channel failing = fakeChannel(seqCounter, new CopyOnWriteArrayList<>(), new CopyOnWriteArrayList<>(),
new AtomicReference<>(), new AtomicReference<>(), true);
AmqpReplyInbox inbox = new AmqpReplyInbox(fakeConnection(failing, failing), AmqpReplyInbox.DEFAULT_PREFETCH);
try {
org.junit.jupiter.api.Assertions.assertThrows(IllegalStateException.class,
() -> inbox.publish("worker", "catch", "body"));
assertEquals(0, pendingByMsgId(inbox).size(), "publish IOException must remove its msgId entry");
} finally {
inbox.close();
}
}
@Test
void interruptedPublishRemovesThePendingMessageIdInFinally() throws Exception {
InboxFixture fixture = new InboxFixture();
Thread publish = fixture.startPublish("finally");
fixture.awaitPublished("finally");
publish.interrupt();
publish.join(5_000);
assertEquals(0, pendingByMsgId(fixture.inbox).size(), "publish finally must remove its msgId entry");
fixture.inbox.close();
}
@Test
void confirmResolutionRemovesThePendingMessageId() throws Exception {
AmqpReplyInbox inbox = new InboxFixture().inbox;
try {
seedPending(inbox, 1, "confirm");
invoke(inbox, "resolveConfirm", new Class<?>[] {long.class, boolean.class, boolean.class}, 1L, false, true);
assertEquals(0, pendingByMsgId(inbox).size(), "confirm resolution must remove its msgId entry");
} finally {
inbox.close();
}
}
@Test
void recoverySweepRemovesThePendingMessageId() throws Exception {
AmqpReplyInbox inbox = new InboxFixture().inbox;
try {
seedPending(inbox, 1, "recovery");
inbox.failPendingPublishesOnRecovery();
assertEquals(0, pendingByMsgId(inbox).size(), "recovery sweep must remove its msgId entry");
} finally {
inbox.close();
}
}
@Test
void closeRemovesThePendingMessageId() throws Exception {
AmqpReplyInbox inbox = new InboxFixture().inbox;
seedPending(inbox, 1, "close");
inbox.close();
assertEquals(0, pendingByMsgId(inbox).size(), "close must remove its msgId entry");
}
/** A {@link Proxy}-backed {@link Channel}: only the calls {@link AmqpReplyInbox} actually makes
* are meaningfully implemented; everything else returns a harmless default. */
private static Channel fakeChannel(AtomicLong seqCounter, List<Long> seqOrder, List<String> msgIdOrder,
AtomicReference<ConfirmCallback> ackCallback,
AtomicReference<ConfirmCallback> nackCallback) {
AtomicReference<ConfirmCallback> ackCallback,
AtomicReference<ConfirmCallback> nackCallback) {
return fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback, false);
}
private static Channel fakeChannel(AtomicLong seqCounter, List<Long> seqOrder, List<String> msgIdOrder,
AtomicReference<ConfirmCallback> ackCallback,
AtomicReference<ConfirmCallback> nackCallback, boolean failPublish) {
InvocationHandler handler = (proxy, method, args) -> {
String name = method.getName();
if (name.equals("getNextPublishSeqNo")) {
@@ -210,6 +281,9 @@ class AmqpReplyInboxRecoveryRaceTest {
return value;
}
if (name.equals("basicPublish")) {
if (failPublish) {
throw new IOException("test publish failure");
}
AMQP.BasicProperties props = (AMQP.BasicProperties) args[3];
msgIdOrder.add(props.getMessageId());
return null;
@@ -285,4 +359,60 @@ class AmqpReplyInboxRecoveryRaceTest {
}
return 0;
}
@SuppressWarnings("unchecked")
private static Map<String, Object> pendingByMsgId(AmqpReplyInbox inbox) throws Exception {
var field = AmqpReplyInbox.class.getDeclaredField("pendingByMsgId");
field.setAccessible(true);
return (Map<String, Object>) field.get(inbox);
}
@SuppressWarnings("unchecked")
private static void seedPending(AmqpReplyInbox inbox, long seq, String msgId) throws Exception {
Class<?> pendingType = Class.forName(AmqpReplyInbox.class.getName() + "$Pending");
var constructor = pendingType.getDeclaredConstructor(String.class);
constructor.setAccessible(true);
Object pending = constructor.newInstance(msgId);
var seqField = AmqpReplyInbox.class.getDeclaredField("pendingBySeq");
seqField.setAccessible(true);
((Map<Long, Object>) seqField.get(inbox)).put(seq, pending);
pendingByMsgId(inbox).put(msgId, pending);
}
private static void invoke(AmqpReplyInbox inbox, String name, Class<?>[] types, Object... args) throws Exception {
Method method = AmqpReplyInbox.class.getDeclaredMethod(name, types);
method.setAccessible(true);
method.invoke(inbox, args);
}
private static final class InboxFixture {
final AtomicLong seqCounter = new AtomicLong();
final List<Long> seqOrder = new CopyOnWriteArrayList<>();
final List<String> msgIdOrder = new CopyOnWriteArrayList<>();
final AtomicReference<ConfirmCallback> ackCallback = new AtomicReference<>();
final AtomicReference<ConfirmCallback> nackCallback = new AtomicReference<>();
final AmqpReplyInbox inbox = new AmqpReplyInbox(
fakeConnection(fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback),
fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback)),
AmqpReplyInbox.DEFAULT_PREFETCH);
Thread startPublish(String msgId) {
Thread thread = Thread.ofVirtual().start(() -> {
try {
inbox.publish("worker", msgId, "body");
} catch (IllegalStateException ignored) {
// Interrupting the confirm wait is the path under test.
}
});
return thread;
}
void awaitPublished(String msgId) throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
while (!msgIdOrder.contains(msgId) && System.nanoTime() < deadline) {
Thread.sleep(10);
}
assertTrue(msgIdOrder.contains(msgId), "publish did not register " + msgId);
}
}
}
@@ -12,13 +12,18 @@ import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;
import java.io.IOException;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Method;
import java.lang.reflect.Proxy;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;
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 +212,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
@@ -347,6 +377,72 @@ class LeadMailboxTest {
() -> "expected AlreadyClosedException, got: " + thrown);
}
@Test
void publishIOExceptionRemovesThePendingMessageId() throws Exception {
LeadMailbox mailbox = newMailbox(true);
try {
assertThrows(IllegalStateException.class,
() -> mailbox.publish("target", new LeadMessage("catch", "from", "target", "body")));
assertEquals(0, pendingByMsgId(mailbox).size(), "publish IOException must remove its msgId entry");
} finally {
mailbox.close();
}
}
@Test
void interruptedPublishRemovesThePendingMessageIdInFinally() throws Exception {
LeadMailbox mailbox = newMailbox(false);
Thread publish = Thread.ofVirtual().start(() -> {
try {
mailbox.publish("target", new LeadMessage("finally", "from", "target", "body"));
} catch (IllegalStateException ignored) {
// Interrupting the confirm wait is the path under test.
}
});
awaitPending(mailbox, "finally");
publish.interrupt();
publish.join(5_000);
try {
assertEquals(0, pendingByMsgId(mailbox).size(), "publish finally must remove its msgId entry");
} finally {
mailbox.close();
}
}
@Test
void confirmResolutionRemovesThePendingMessageId() throws Exception {
LeadMailbox mailbox = newMailbox(false);
try {
seedPending(mailbox, 1, "confirm");
invoke(mailbox, "resolveConfirm", new Class<?>[] {long.class, boolean.class, boolean.class}, 1L, false, true);
assertEquals(0, pendingByMsgId(mailbox).size(), "confirm resolution must remove its msgId entry");
} finally {
mailbox.close();
}
}
@Test
void recoverySweepRemovesThePendingMessageId() throws Exception {
LeadMailbox mailbox = newMailbox(false);
try {
seedPending(mailbox, 1, "recovery");
mailbox.failPendingPublishesOnRecovery();
assertEquals(0, pendingByMsgId(mailbox).size(), "recovery sweep must remove its msgId entry");
} finally {
mailbox.close();
}
}
@Test
void closeRemovesThePendingMessageId() throws Exception {
LeadMailbox mailbox = newMailbox(false);
seedPending(mailbox, 1, "close");
mailbox.close();
assertEquals(0, pendingByMsgId(mailbox).size(), "close must remove its msgId entry");
}
/** Poll peek until at least one message is held, or ~10s elapse (broker delivery is async). */
@SuppressWarnings("BusyWait")
private static List<LeadMessage> awaitPeek(LeadMailbox inbox) throws InterruptedException {
@@ -372,4 +468,108 @@ class LeadMailboxTest {
}
return state;
}
private static LeadMailbox newMailbox(boolean failPublish) {
AtomicLong sequence = new AtomicLong();
Channel consume = fakeChannel(sequence, false);
Channel publish = fakeChannel(sequence, failPublish);
return new LeadMailbox(fakeConnection(consume, publish), "self");
}
private static Channel fakeChannel(AtomicLong sequence, boolean failPublish) {
InvocationHandler handler = (proxy, method, args) -> {
if (method.getName().equals("getNextPublishSeqNo")) {
return sequence.incrementAndGet();
}
if (method.getName().equals("basicPublish") && failPublish) {
throw new IOException("test publish failure");
}
if (method.getName().equals("equals")) {
return proxy == args[0];
}
if (method.getName().equals("hashCode")) {
return System.identityHashCode(proxy);
}
return defaultValue(method.getReturnType());
};
return (Channel) Proxy.newProxyInstance(LeadMailboxTest.class.getClassLoader(), new Class<?>[] {Channel.class}, handler);
}
private static Connection fakeConnection(Channel first, Channel second) {
AtomicLong calls = new AtomicLong();
InvocationHandler handler = (proxy, method, args) -> {
if (method.getName().equals("createChannel") && (args == null || args.length == 0)) {
return calls.getAndIncrement() == 0 ? first : second;
}
if (method.getName().equals("equals")) {
return proxy == args[0];
}
if (method.getName().equals("hashCode")) {
return System.identityHashCode(proxy);
}
return defaultValue(method.getReturnType());
};
return (Connection) Proxy.newProxyInstance(LeadMailboxTest.class.getClassLoader(), new Class<?>[] {Connection.class}, handler);
}
@SuppressWarnings("unchecked")
private static Map<String, Object> pendingByMsgId(LeadMailbox mailbox) throws Exception {
var field = LeadMailbox.class.getDeclaredField("pendingByMsgId");
field.setAccessible(true);
return (Map<String, Object>) field.get(mailbox);
}
@SuppressWarnings("unchecked")
private static void seedPending(LeadMailbox mailbox, long seq, String msgId) throws Exception {
Class<?> pendingType = Class.forName(LeadMailbox.class.getName() + "$Pending");
var constructor = pendingType.getDeclaredConstructor(String.class);
constructor.setAccessible(true);
Object pending = constructor.newInstance(msgId);
var seqField = LeadMailbox.class.getDeclaredField("pendingBySeq");
seqField.setAccessible(true);
((Map<Long, Object>) seqField.get(mailbox)).put(seq, pending);
pendingByMsgId(mailbox).put(msgId, pending);
}
private static void invoke(LeadMailbox mailbox, String name, Class<?>[] types, Object... args) throws Exception {
Method method = LeadMailbox.class.getDeclaredMethod(name, types);
method.setAccessible(true);
method.invoke(mailbox, args);
}
private static void awaitPending(LeadMailbox mailbox, String msgId) throws Exception {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
while (!pendingByMsgId(mailbox).containsKey(msgId) && System.nanoTime() < deadline) {
Thread.sleep(10);
}
assertTrue(pendingByMsgId(mailbox).containsKey(msgId), "publish did not register " + msgId);
}
private static Object defaultValue(Class<?> type) {
if (!type.isPrimitive() || type == void.class) {
return null;
}
if (type == boolean.class) {
return Boolean.FALSE;
}
if (type == long.class) {
return 0L;
}
if (type == short.class) {
return (short) 0;
}
if (type == byte.class) {
return (byte) 0;
}
if (type == char.class) {
return (char) 0;
}
if (type == double.class) {
return 0.0d;
}
if (type == float.class) {
return 0.0f;
}
return 0;
}
}
@@ -22,6 +22,7 @@ import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
@@ -505,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
@@ -554,6 +601,30 @@ class MessageServiceTest {
}
}
/**
* fleetd #571 (the acceptance test the ticket was filed for). The worker is idle so the injector
* attempts delivery, but the {@code agent.prompt} call itself fails with a herdr error that is
* not a confirmed absence (not a {@code *_not_found} code) — {@link Injector} marks the Pending
* {@code ATTEMPTED} (fleetd #551), meaning the call was made and whether it reached the pane is
* unknown. Before this fix, {@code send}'s {@code TimeoutException} branch collapsed
* {@code ATTEMPTED} into {@code TIMED_OUT_QUEUED} — a promise that the message will never arrive,
* which may already be false: {@code agent.prompt} pastes and submits in one call.
*/
@Test
void sendTimesOutWithAttemptedDeliveryReportsUnconfirmedNotQueued() throws Exception {
herdr.agentSendFailsWith("send_failed");
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> messages.send(T, "brief", 150));
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // triggers the failing delivery attempt → ATTEMPTED
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.TIMED_OUT_UNCONFIRMED, r.outcome(),
"an ATTEMPTED delivery must not collapse into TIMED_OUT_QUEUED — the message may "
+ "already have arrived in full, and TIMED_OUT_QUEUED promises it never will");
assertNull(r.text());
}
@Test
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
@@ -579,6 +650,183 @@ class MessageServiceTest {
assertEquals("config.yaml", a.answer());
}
// --- fleetd #572: answer() must release the session lock on EVERY exit, not just return it -----
//
// answer()'s session lock is released in an outer `finally` (MessageService.java:1217-1219) that
// wraps its whole body. Removing that one line survives the entire suite: coverage is not the
// gap, assertion is — every existing answer() test checks the RETURN VALUE, never that the lock
// it took is reacquirable afterward. If it is not, the session is wedged forever with no
// exception and no log line. Each test below drives answer() through one of its four exits
// (normal reply, TIMED_OUT_WORKING, ExecutionException, InterruptedException) and then proves
// reacquisition the only way that actually proves it: a bounded follow-up `send` on the SAME
// session must not come back BUSY. A `send` only reports BUSY when `tryLock` itself timed out
// (MessageService.java:926-928) — every other early return in `send` still passes through its own
// lock-acquired `try`, so a non-BUSY probe result is specifically evidence the lock was free.
/** The REPLIED exit (the happy path) — the resumed worker's real {@code fleet_reply} arrives. */
@Test
void answerReleasesTheSessionLockAfterANormalReply() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
ask.get(5, TimeUnit.SECONDS); // worker resumed with the answer
awaitWaiting(); // the answering call has (re)opened its own forward waiter
assertTrue(rendezvous.resolve(T, "done"), "the worker's final reply resolves the answering send");
MessageService.Reply done = answer.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.REPLIED, done.outcome());
MessageService.Reply probe = messages.send(T, "probe after normal reply", 300);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released after a normal REPLIED answer(), or this bounded "
+ "follow-up send would come back BUSY instead of timing out on its own work");
}
/**
* The {@code TIMED_OUT_WORKING} exit. Same setup as {@link #answerTimesOutWhenTheResumedWorkerNeverReplies}
* (which only checks {@code answer()}'s return value — exactly the assertion the fleetd #572
* mutation survives), plus the reacquisition proof that test does not make.
*
* <p>{@code answer()} MUST run on its own thread here, not the test's main thread: {@code lock}
* is a {@link java.util.concurrent.locks.ReentrantLock}, so a probe {@code send} issued by the
* SAME thread that (under the mutation) still "holds" it would reenter for free and report a
* false pass — reentrancy, not release. Measured while writing this test: with {@code answer()}
* called inline, this test stayed green under the mutation while its three siblings correctly
* went red.
*/
@Test
void answerReleasesTheSessionLockAfterATimedOutWorkingReturn() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
// The primary answers, unblocking the worker; the worker never sends its follow-up
// fleet_reply, so the answering call rides out its short window as still-working.
CompletableFuture<MessageService.Reply> answer =
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 200));
MessageService.Reply answered = answer.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answered.outcome(),
"an answered worker that never replies times out as still working");
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
MessageService.Reply probe = messages.send(T, "probe after timed-out-working", 300);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released after a TIMED_OUT_WORKING answer(), or this bounded "
+ "follow-up send would come back BUSY instead of timing out on its own work");
}
/**
* The {@code ExecutionException} exit. Forces it directly on answer()'s own reopened forward
* waiter — {@code rendezvous.currentWaiter(T)} is the exact {@code CompletableFuture} its
* {@code reply.get(...)} is blocked on — rather than trying to make a real worker fail, since the
* failure mode under test is in answer()'s own wait, not in how it got triggered.
*/
@Test
void answerReleasesTheSessionLockWhenTheReplyFutureFailsExceptionally() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
String turnId = q.turnId();
java.util.concurrent.atomic.AtomicReference<Throwable> caught =
new java.util.concurrent.atomic.AtomicReference<>();
Thread answerer = new Thread(() -> {
try {
messages.answer(turnId, "config.yaml", 5000);
caught.set(new AssertionError("expected answer() to throw"));
} catch (Throwable t) {
caught.set(t);
}
});
answerer.start();
awaitWaiting(); // answer() re-opened its forward waiter and is about to block on it
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(T);
assertNotNull(waiter, "answer() should have registered a forward waiter for the worker session");
waiter.completeExceptionally(new RuntimeException("boom"));
answerer.join(5000);
assertFalse(answerer.isAlive(), "answer() should have left by throwing once its reply future failed");
Throwable thrown = caught.get();
assertNotNull(thrown, "answer() should have thrown");
assertTrue(thrown instanceof RuntimeException, "a RuntimeException cause is rethrown as-is: " + thrown);
assertEquals("boom", thrown.getMessage());
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
MessageService.Reply probe = messages.send(T, "probe after exceptional failure", 300);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released when answer()'s reply future fails exceptionally, "
+ "or this bounded follow-up send would come back BUSY instead of timing out on its own work");
}
/**
* The {@code InterruptedException} exit. Same shape as the {@code ExecutionException} case above,
* but the thread blocked in {@code reply.get(...)} is interrupted instead of the future failing.
*/
@Test
void answerReleasesTheSessionLockWhenInterrupted() throws Exception {
CompletableFuture<MessageService.Reply> send = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
String turnId = q.turnId();
java.util.concurrent.atomic.AtomicReference<Throwable> caught =
new java.util.concurrent.atomic.AtomicReference<>();
Thread answerer = new Thread(() -> {
try {
messages.answer(turnId, "config.yaml", 5000);
caught.set(new AssertionError("expected answer() to throw"));
} catch (Throwable t) {
caught.set(t);
}
});
answerer.start();
awaitWaiting(); // answer() re-opened its forward waiter and is about to block on it
answerer.interrupt();
answerer.join(5000);
assertFalse(answerer.isAlive(), "answer() should have left by throwing once its thread was interrupted");
Throwable thrown = caught.get();
assertNotNull(thrown, "answer() should have thrown");
assertTrue(thrown instanceof IllegalStateException,
"an interrupted wait is wrapped in IllegalStateException: " + thrown);
ask.get(5, TimeUnit.SECONDS); // drain: the worker resumed with the answer
MessageService.Reply probe = messages.send(T, "probe after interruption", 300);
assertNotEquals(MessageService.Outcome.BUSY, probe.outcome(),
"the session lock must be released when answer()'s wait is interrupted, or this bounded "
+ "follow-up send would come back BUSY instead of timing out on its own work");
}
@Test
void pollReturnsNullForAnUnknownTicket() {
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
@@ -10,12 +10,13 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
class MemberRoleTest {
@Test
void theThreeRolesAreArchitectDevAndReviewer() {
assertEquals(3, MemberRole.values().length,
void theFourRolesAreArchitectDevHunterAndReviewer() {
assertEquals(4, MemberRole.values().length,
"a new role changes the charter, the role file, the skill and the authz row — "
+ "adding one is a deliberate act, so this count is meant to fail first");
assertEquals("architect", MemberRole.ARCHITECT.wireName());
assertEquals("dev", MemberRole.DEV.wireName());
assertEquals("hunter", MemberRole.HUNTER.wireName());
assertEquals("reviewer", MemberRole.REVIEWER.wireName());
}
@@ -30,6 +31,7 @@ class MemberRoleTest {
void parseIsCaseInsensitiveAndTrimsSurroundingSpace() {
assertSame(MemberRole.ARCHITECT, MemberRole.parse("Architect"));
assertSame(MemberRole.DEV, MemberRole.parse(" DEV "));
assertSame(MemberRole.HUNTER, MemberRole.parse("HuNtEr"));
assertSame(MemberRole.REVIEWER, MemberRole.parse("ReViEwEr"));
}
@@ -38,7 +40,7 @@ class MemberRoleTest {
IllegalArgumentException e =
assertThrows(IllegalArgumentException.class, () -> MemberRole.parse("archtiect"));
assertTrue(e.getMessage().contains("archtiect"), e.getMessage());
assertTrue(e.getMessage().contains("architect, dev, reviewer"),
assertTrue(e.getMessage().contains("architect, dev, hunter, reviewer"),
"a typo in config should be fixable from the message alone: " + e.getMessage());
}
@@ -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();
@@ -623,6 +647,28 @@ class FleetAppTest {
assertTrue(herdr.called("agent.prompt"), "message was injected");
}
/**
* fleetd #571: the worker is idle, so the poller attempts delivery, but the {@code agent.prompt}
* call itself fails with a herdr error that is not a confirmed absence — {@link
* dev.ltms.fleet.inject.Injector} marks this {@code ATTEMPTED}, meaning the call was made and
* whether it reached the pane is unknown. {@code writeReply}'s default arm must map this to its
* own {@code "unconfirmed"} status, not silently fall through to {@code "done"} (which would
* claim the delegation completed) nor collapse into {@code "queued"} (which would claim the
* message will never arrive, when it may already be sitting in the pane).
*/
@Test
void messageTimesOutUnconfirmedWhenDeliveryAttemptFails() throws Exception {
FakeHerdr herdr = new FakeHerdr().agentStatus("idle").agentSendFailsWith("send_failed");
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
HttpResponse<String> res = postMessage(port, "{\"content\":\"hi\",\"timeoutMs\":250}");
assertEquals(202, res.statusCode());
JsonNode body = mapper.readTree(res.body());
assertEquals("unconfirmed", body.get("status").asText(),
"an ATTEMPTED delivery must report its own status, not \"queued\" or \"done\"");
assertTrue(herdr.called("agent.prompt"), "delivery must have been attempted");
}
@Test
void messageRejectsBlankContent() throws Exception {
int port = startHealthy();