Compare commits

..

21 Commits

Author SHA1 Message Date
Dai Ha 158a2a84b5 Merge PR #632: fleetd #612 ranks 6+7 + #630 — behavioural pins for lifecycle wirings
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 1m13s
CI / build (push) Failing after 2m10s
Pins three call sites the assembly owns and nothing observed:
- healthFailTarget (FleetdAssembly:429) — inert, a dead member's waiting ticket sits
  PENDING for the full 30-minute async timeout instead of failing immediately.
- releaseCleanup (:447) — inert, every teardown leaks three things: a stuck rendezvous
  waiter, an unreleased reply-inbox consumer, and a stale lead binding.
- requireOperatorConfirm (:402/:409, fleetd #630) — dropping the 14th constructor
  argument selects #621's 13-arg overload, which hardcodes true, silently reverting
  the operator's fix on a host that set requireOperatorConfirm: false. Pinned in both
  directions, plus an assertion that the two notice strings differ, so no constant
  satisfies both.

Verified by the lead beyond the worker's proof: its releaseCleanup mutation killed all
three cleanups at once and so proved only the first assertion had teeth. Starving them
one at a time — abandon kept, release starved; then abandon and release kept, forget
starved — each fails its own named assertion. All three leaks are pinned independently.

A review pass found these three tests assembled real schedulers and never tore them
down, by any route: no close(), no shutdownHook, no @AfterEach, no finally. Surefire
runs one JVM fork for the whole suite, so those loops outlived their tests. Fixed by
capturing the hook and running it in a finally, with an assertion on FakeHerdr.closed
so the teardown itself is pinned rather than assumed.

MERGE RESOLUTION BY THE LEAD, the same collision as #633. These 3 test files each add
a ResourcePorts fake, and #633 landed first making herdrPollWait() abstract with no
default. Git reported a clean merge that did not compile. Added the override to all 3,
matching the established convention for an always-healthy fake — a Runnable that
throws, verified first that none of the three uses healthy(false), so the tripwire can
only fire if the test's herdr behaviour changes.

Full suite on the resolved merge: 1892 tests, 0 failures (1889 + this branch's 3).
2026-10-01 16:52:39 +02:00
Dai Ha 41ebc9cf69 Merge PR #633: fleetd #629 + #625 — the two ResourcePorts seams
CI / shell-tests (push) Failing after 12s
CI / contract (push) Successful in 53s
CI / build (push) Failing after 2m57s
#629: the herdr boot wait went through ports.nanoClock() but hardcoded
Fleetd::sleepHerdrPoll, so a test assembling against an unhealthy lead herdr burned
30 real seconds whatever clock it injected. The poll sleep now goes through
ResourcePorts.herdrPollWait(). FleetdAssemblyFleetAppTest's lead-down test drops from
30.276s to 0.062s. It also gains @Timeout(10, SEPARATE_THREAD) at class level: the
pin's failure mode is otherwise an infinite hang, because the test's fake clock only
advances when herdrPollWait() is called. SAME_THREAD cannot interrupt a real
Thread.sleep, so the thread mode is load-bearing, not decoration.

#625: guard.assertPrimaryClean(System.getenv()) could be deleted with a fully green
suite — the check behind charter invariant 1, which keeps the primary on the
operator's subscription. main() now delegates to main(String[], ResourcePorts) and
the guard reads ports.environment(), so a test can taint the environment without
touching the real process env. The guard still runs before cfg.validateAll() and
before any socket, broker or HTTP work; two tests with different fixtures pin that
ordering as two independently falsifiable claims, not one.

MERGE RESOLUTION BY THE LEAD. ResourcePorts.herdrPollWait() is abstract with no
default, by design (#629 keeps ResourcePorts free of a none() default). This branch
patched the 8 implementations that existed when it forked. PRs #631 and #634 merged
ahead of it and added 3 more fakes, so git reported a clean 14-file merge that did
not compile:

  FleetdAssemblyAmqpOpenersTest.RecordingPorts is not abstract and does not
  override abstract method herdrPollWait() in dev.ltms.fleet.ResourcePorts

(plus the two ControllableResourcePorts in #634's tests). I added the override to
those 3, matching this branch's own convention for an always-healthy fake: return a
Runnable that throws, so if one of those assemblies ever does start polling herdr it
fails loudly instead of sleeping quietly. All 13 implementations now carry it.

Full suite on the resolved merge: 1889 tests, 0 failures (1886 + this branch's 3).
2026-10-01 16:50:10 +02:00
Dai Ha 619769bb52 fleetd #632: tear down the three assembly behavioural tests' background loops
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 1m6s
CI / build (pull_request) Failing after 2m39s
Each of FleetdAssemblyHealthFailTargetBehaviouralTest,
FleetdAssemblyReleaseCleanupBehaviouralTest and
FleetdAssemblyRequireOperatorConfirmBehaviouralTest called
FleetdAssembly.assembleAndStart without ever tearing it down: no close(),
no shutdownHook, no @AfterEach, no finally. Surefire runs the whole suite
in one JVM fork, so every scheduler/loop these tests started kept running
for the rest of the suite.

Capture the shutdown hook in each test's fake ResourcePorts (the existing
pattern from FleetdAssemblyLifecycleTest et al.) and run it in a finally
block, on the failure path too. RequireOperatorConfirmBehaviouralTest
assembles twice in one method, so assembleHeartbeat now returns both the
loop and its ResourcePorts so each assembly gets its own teardown.

Proof the teardown actually runs: assert ports.herdr.closed after running
the hook (FakeHerdr.close() only flips that flag from inside the real
close chain). Verified the assertion is load-bearing by temporarily
removing one shutdownHook.run() call and confirming the test then fails.

No assertion, test name or reflection changed. Diff is test-only.
2026-10-01 16:49:14 +02:00
Dai Ha 180c953c42 Merge PR #634: fleetd #612 ranks 1+2 — behavioural pins for the exhaustion wirings
CI / shell-tests (push) Failing after 8s
CI / build (push) Failing after 2m20s
CI / contract (push) Successful in 2m22s
Pins all four exhaustion call sites in FleetdAssembly: liveExhaustedPatterns (:299),
exhaustedPatternLookup (:300), publishExhaustionSink (:327) and the independent
OpenCode forwardingExhaustionSink (:164). Rank 1 is the worst consequence in the
#612 sweep — an inert lookup hands a genuine usage-limit refusal back to a waiting
caller as real completed work instead of BACKEND_EXHAUSTED.

Two separate tests, so rank 2's OpenCode half is pinned independently: mutating
:164 fails only the forwarding test, which is the independence the ticket asserts.

Verified by the lead beyond the worker's own proof: wiring publishExhaustionSink to
a throwaway BackendQuarantine — every symbol kept at the call site, only the
collaborator identity changed — is caught by both tests. So these pins survive
mis-wiring, not just deletion.

Test-only; no production change. Tears down each assembly via the captured
shutdown hook in a finally.
2026-10-01 16:45:20 +02:00
Dai Ha 4af919af0c fleetd #629 follow-up: bound FleetdAssemblyFleetAppTest with a SEPARATE_THREAD @Timeout
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 1m6s
CI / build (pull_request) Failing after 2m25s
The mutation cycle for #629 proved deleting the fix (reverting FleetdAssembly's awaitHerdr
call back to a hardcoded Fleetd::sleepHerdrPoll) doesn't just make a test fail — it hangs
forever, because the test's fake nanoClock() only advances when ports.herdrPollWait() is
actually called. SAME_THREAD @Timeout (JUnit's default) can't catch that: it only measures
elapsed time after the test method returns on its own, which never happens here. A class-level
@Timeout(10s, SEPARATE_THREAD) does, since it runs the test on its own thread and interrupts it
on timeout. Verified by re-running the same mutation: the suite now fails fast with a named
TimeoutException instead of hanging indefinitely.
2026-10-01 16:41:40 +02:00
Dai Ha 09c37061c1 Merge PR #631: fleetd #612 ranks 3+8 — behavioural pins for the AMQP assembly openers
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 59s
CI / build (push) Failing after 3m14s
Pins FleetdAssembly's replyInboxOpener (:357) and leadMailboxOpener (:363) so an
inert opener can no longer pass a green suite. Both directions covered: a durable
opener must leave the runtime owning the exact object it returned, and a failing
one must leave the in-memory fallback with coordination off — with the startup
report agreeing with the real state in each case.

Verified by the lead beyond the worker's own proof: a mutation that still CALLS
ports.replyInboxOpener() and discards the result is caught by assertSame, which is
rank 3's live defect (opener called, result thrown away, log still printing
'reply inbox: AMQP broker (durable)').

Test-only; no production change.
2026-10-01 16:38:16 +02:00
Dai Ha 2be287ea03 fleetd #612 step 4 (ranks 1&2): behavioural assembly tests for the exhaustion wirings
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 3m5s
Replaces nothing (no existing test covered these call sites through the real
assembly); adds two new FleetdAssembly-driven tests, matching Unit A's pattern
of inspecting FleetdRuntime's real, assembled objects rather than a copy.

- FleetdExhaustedPatternAssemblyTest pins the liveExhaustedPatterns/
  exhaustedPatterns call sites (rank 1 — "the worst consequence in the whole
  sweep": a genuine usage-limit refusal handed back as real completed work)
  together with publishExhaustionSink (rank 2, non-OpenCode half): drives a
  real CompletionResolver through a pane scrape matching a configured
  exhaustedPattern and asserts BACKEND_EXHAUSTED classification plus a real
  BackendQuarantine credential quarantine.

- FleetdOpenCodeExhaustionForwardingAssemblyTest pins forwardingExhaustionSink
  (rank 2, OpenCode half — independent of publishExhaustionSink per the
  ticket) by reflectively reaching the real, assembled OpenCodeLauncher's
  exhaustionSink field (SessionManager.launcher -> CompositePeerLauncher.
  byProfile -> OpenCodeLauncher.exhaustionSink) and proving it forwards into
  the same production BackendQuarantine.

All four call sites (FleetdAssembly.java:164,299,300,327) were each put
through grep-anchor -> line-anchored sed mutation -> mvn -o compile -> full
mvn -o test (named test RED) -> restore -> full mvn -o test (1885/0 GREEN).
Mutating line 327 alone also fails the OpenCode test, confirming the
documented construction-order dependency (forwardingExhaustionSink reads
exhaustionSinkRef, which publishExhaustionSink sets) without weakening either
site's independent pin.
2026-10-01 16:38:00 +02:00
Dai Ha c3b0406826 fleetd #612: isolate fallback AMQP reports
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Failing after 2m20s
2026-10-01 16:33:45 +02:00
Dai Ha e76fa1660b fleetd #629/#625: inject the herdr poll wait and the startup env read through ResourcePorts
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 1m31s
CI / build (pull_request) Failing after 2m42s
#629: ResourcePorts gains herdrPollWait() (SystemResourcePorts: Fleetd::sleepHerdrPoll).
FleetdAssembly's awaitHerdr call now takes its poll wait from ports instead of a hardcoded
Thread.sleep, so a test's fake clock can actually reach the deadline without burning real
wall-clock time. FleetdAssemblyFleetAppTest's down-lead case drops from ~30s to well under 1s.
All other ResourcePorts fakes get a trivial "must not be called" override since their herdr is
always healthy and never polls.

#625: Fleetd.main(String[]) now delegates to a new package-private
main(String[], ResourcePorts) overload that reads the startup environment via
ports.environment() instead of System.getenv(), reusing the same ports instance for
FleetdAssembly.assembleAndStart. The guard call stays at its original point, before
cfg.validateAll() and before any assembly/socket/broker/HTTP work. New
FleetdSubscriptionGuardOrderingTest drives the real main() with a tainted fake environment and
pins both presence and ordering against validateAll() and against assembly's first port call,
plus a clean-env control proving the guard only blocks on an actual taint.
2026-10-01 16:33:15 +02:00
Dai Ha 9dca376604 fleetd #612 ranks 6/7 + #630: behavioural pins for healthFailTarget, releaseCleanup, requireOperatorConfirm
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Failing after 2m58s
Three FleetdAssembly.java call sites had no test that drives the real
assembly and inspects what it actually built, so each could be swapped
for an inert variant and the suite would stay green:

- FleetdAssembly.java:429 (healthFailTarget, #612 rank 7): a no-op
  BiConsumer leaves a dead member's waiting ticket PENDING for the full
  30-minute async timeout instead of failing it immediately.
- FleetdAssembly.java:447 (releaseCleanup, #612 rank 6): a no-op
  onRelease listener leaks a stuck rendezvous waiter, an unreleased
  reply-inbox consumer, and a stale lead binding on every teardown.
- FleetdAssembly.java:402/:409 (requireOperatorConfirm, #630): dropping
  the 14th LeadHeartbeatLoop constructor argument selects the
  13-argument overload, which hardcodes true (fleetd #621) regardless
  of leadRollover.requireOperatorConfirm — silently reverting an
  operator's own config choice.

Each existing "wiring" test for these sites (FleetdHealthFailTargetWiringTest,
FleetdReleaseCleanupWiringTest) calls the Fleetd.* factory method directly
and never drives FleetdAssembly.assembleAndStart, so none of them can see
whether the real call site still passes the real, assembled collaborators.

The three new tests here assemble the real daemon via
FleetdAssembly.assembleAndStart, reach the real wired object (reflection,
same technique StatusPollerResilienceTest already uses — the fields and
the requireOperatorConfirm overload resolution point are package-private),
and assert the real behavioural effect against the real collaborators
FleetdRuntime exposes. Each is mutation-tested: a line-anchored sed to the
inert form compiles clean and turns the new test RED; restoring the
original line turns it GREEN again. requireOperatorConfirm is proven in
both directions (false and true produce genuinely different notice text).

Full suite: 1886 tests, 0 failures, 0 errors (baseline 1883 + 3 new).
2026-10-01 16:30:23 +02:00
Dai Ha 91be5079d6 fleetd #612: pin AMQP assembly openers
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m44s
CI / build (pull_request) Failing after 3m41s
2026-10-01 16:26:02 +02:00
ltms 26f198675c Merge pull request 'fleetd #612 Unit A: extract main's boot composition into FleetdAssembly/FleetdRuntime' (#620) from worker/fleetd-612-unita-87807e-1 into main
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 50s
CI / build (push) Failing after 2m30s
2026-09-22 07:43:39 +02:00
lead 72f46d7c0c fleetd #612: merge main into Unit A, carrying #621's requireOperatorConfirm into the assembly
CI / shell-tests (pull_request) Failing after 11s
CI / contract (pull_request) Successful in 48s
CI / build (pull_request) Failing after 2m59s
Resolves the one conflict in Fleetd.java. main's side is the inline boot block
that Unit A had already moved into FleetdAssembly.assembleAndStart, so the
resolution keeps Unit A's single assembly call.

That resolution is not purely mechanical. #622 (fleetd #621) landed on main
AFTER Unit A forked, and it added

    boolean requireOperatorConfirm = cfg.leadRollover() == null
            || cfg.leadRollover().requireOperatorConfirm();

plus a 14th argument to the LeadHeartbeatLoop constructor, inside the very block
Unit A moved. Taking Unit A's side alone would have dropped both and silently
reverted the operator's #621 fix: the 13-argument overload still exists and
delegates with `true`, so the daemon would go back to telling every lead to ask
the operator before a context roll. Both are carried into FleetdAssembly here.

Measured: with the carried line removed, the full suite is
`Tests run: 1883, Failures: 0, Errors: 0` — nothing pins it. That is the fleetd
#612 defect shape applied to #621's own wiring, and it is filed separately
rather than fixed here, because this commit is a merge resolution and must not
also introduce new tests.

Full suite on this resolved tree: Tests run: 1883, Failures: 0, Errors: 0.
2026-09-22 12:43:22 +07:00
ltms fa61dc587c Merge pull request '#608 replace MessageService timing sleeps' (#623) from worker/608-sleeps-3a64ff-3 into main
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 3h14m55s
2026-09-22 06:38:29 +02:00
Dai Ha b6b006c651 #608 replace MessageService timing sleeps
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 1m13s
CI / build (pull_request) Failing after 2m27s
2026-09-22 11:33:48 +07:00
ltms 63eec8a0da Merge pull request 'fleetd #621: make the context-roll notice obey requireOperatorConfirm' (#622) from worker/621-b4520b-1 into main
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 1m11s
CI / build (push) Failing after 1m52s
2026-09-22 06:31:02 +02:00
Dai Ha cbe872b538 fleetd #621: make the context-roll notice obey requireOperatorConfirm
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m0s
CI / build (pull_request) Failing after 1m53s
contextNotice hardcoded 'ask the operator' and 'Only the operator can
approve the roll', so setting leadRollover.requireOperatorConfirm to
false stopped the daemon refusing the roll but never stopped the lead
being told to ask. Thread the effective config value into
contextNotice: when true the text stays byte-identical, when false it
tells the lead to confirm on its own judgement against the three
handover-file checks instead.

LeadRollover.confirm's own enforcement is untouched — this is the
message only.
2026-09-22 11:27:17 +07:00
Dai Ha 7d9a807243 handover skill: the rollover bootstrap is proven, and requireOperatorConfirm is per-host
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 53s
CI / build (push) Failing after 2m33s
Two bullets in the handover skill were telling every outgoing lead something
that is no longer true.

1. The skill said the bootstrap prompt "has never yet landed, and the fix is
   unproven (fleetd #489)", and told the lead to warn the operator it may fail.
   Measured today from fleetd/fleetd.out:

     grep -c "lead-rollover: rolled" -> 4
     grep -c "lead-rollover:"        -> 16   (positive control)
     grep -c "Unknown command"       -> 0

   Three of the four rolls ran on 2026-09-22 (10:01:43, 10:38:28, 11:15:47).
   Each cleared the old lead and bootstrapped a fresh one against the handover
   file. The old "Unknown command: /clearFresh" failure does not appear at all.
   The paragraph now carries the measured result and the three re-measure
   commands, including the control line, because a broken grep pattern returns
   a clean 0 that reads like good news.

   It also records that the "/clear was never observed as WORKING ... releasing
   rather than wedging the roll" WARN accompanies every successful roll. That
   is the safe branch, not a failure, and it was being misread as one.

2. The skill said requireOperatorConfirm "defaults to true and this is the only
   thing standing between a judgement call and a wiped session", which reads as
   if asking is always required. The default is still true
   (FleetConfig.java:1426), but this host set it to false on 2026-09-22 on the
   operator's explicit grant. The bullet now says to read the live value rather
   than assume, and notes the key is deferred, not hot.

   It also warns that until fleetd #621 merges, LeadHeartbeatLoop.contextNotice()
   still hardcodes "ask the operator" and takes no config, so the nudge text and
   the config disagree. Trust the config. That warning names the ticket that
   removes it.

Documentation only. No code or test changes.
2026-09-22 11:24:03 +07:00
ltms 8915e40c7d Merge pull request 'fleetd #618: state the measured auto-compact precedence' (#619) from worker/618-b83894-2 into main
CI / shell-tests (push) Failing after 11s
CI / contract (push) Successful in 48s
CI / build (push) Failing after 2m12s
2026-09-22 05:53:35 +02:00
Dai Ha 6cb31a10e4 fleetd #618: fix the third stale spot the brief missed
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m30s
CI / build (pull_request) Failing after 2m6s
The method-level javadoc on FleetConfig.warnConflictingAutoCompactWindows
(above the log.warn call) still claimed the autoCompactWindow vs
CLAUDE_CODE_AUTO_COMPACT_WINDOW precedence was 'intentionally not
asserted' and cited fleetd.yaml's now-corrected comment as evidence the
question was open. Replace it with the measured answer from #618: the
env var wins, so autoCompactWindow is inert on a profile that sets both.
Kept the WARN-not-throw rationale paragraph above it untouched (#601)
and kept the ClaudeCodeArguments cross-reference, which now points to an
agreeing claim instead of a contradicting one. No behaviour change.
2026-09-22 10:50:59 +07:00
Dai Ha 8368a274a0 fleetd #618: state the measured auto-compact precedence, not 'unverified'
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m4s
CI / build (pull_request) Failing after 1m57s
ClaudeCodeArguments.withAutoCompactWindow's javadoc and FleetConfig's
warnConflictingAutoCompactWindows WARN text both used to say the
precedence between --autocompact and CLAUDE_CODE_AUTO_COMPACT_WINDOW was
not verified. fleetd #618 measured it: the env var wins, so the flag has
no effect when both are set. Update both texts to say so, name #618, and
warn that deleting the env var to resolve the conflict LOWERS the live
window rather than fixing anything. No behaviour change; the WARN still
fires on the same condition and stays a WARN (per #601).
2026-09-22 10:46:31 +07:00
26 changed files with 1859 additions and 86 deletions
+51 -11
View File
@@ -167,19 +167,59 @@ fails.
- **There is no terminal or session parameter, on purpose.** The pane is always your own, resolved
from your connection, so you can only ever roll yourself.
- **`operatorConfirmed` is your report of what a human told you.** Do not pass `true` because you
are confident. Ask, wait for the answer, then pass what they said. `requireOperatorConfirm`
defaults to `true` and this is the only thing standing between a judgement call and a wiped
session.
are confident. Ask, wait for the answer, then pass what they said.
- **Whether you must ask at all depends on `leadRollover.requireOperatorConfirm`. Check it; do not
assume.** The default is `true` (`FleetConfig.java:1426`), and then `confirm` refuses unless you
also pass `operatorConfirmed: true`. **This host set it to `false` on 2026-09-22**, on the
operator's explicit grant, because they do not want to approve routine context rolls. Where it is
`false`, the three handover-file checks are the whole gate: the file must exist, be fresher than
`maxDocAgeSeconds`, and have been modified after the open request.
Read the live value rather than trusting this line:
```bash
grep -A1 'requireOperatorConfirm' fleetd/fleetd.yaml
```
No match means the key is unset, so the default `true` applies and you must ask. The key is
**deferred, not hot** — it is read once at boot, so an edit does nothing until the daemon is
redeployed.
**Until fleetd #621 merges, the nudge text will tell you to ask the operator even where the
daemon no longer requires it.** `LeadHeartbeatLoop.contextNotice()` hardcodes "ask the operator"
and takes no config, so it cannot know. Trust the config value over the nudge text. Once #621 is
merged and deployed, the nudge matches the config and this warning can be deleted.
- **The roll can still refuse after `confirm` returns**, and by then there is no caller to tell.
Those outcomes are logged only, as `lead-rollover:` lines in the daemon log.
- **The bootstrap prompt has never yet landed, and the fix is unproven (fleetd #489).** The first
real rollover, on 2026-09-12, joined `/clear` and the bootstrap text into one line and Claude Code
refused it as `Unknown command: /clearFresh`. The pane was never cleared and no context was lost,
so the failure was safe — the roll simply did nothing. PR #490 fixed the cause and is deployed,
but no roll has bootstrapped a fresh session end to end yet. **Assume it may still fail, and tell
the operator so before you confirm.** The recovery is the same either way: the file is already
written, so the operator starts a session and points it at the file. That is why you write the
file before you confirm, and never the other way round.
- **The bootstrap prompt works end to end. Measured 2026-09-22.** This used to say the fix was
unproven (fleetd #489) and told you to expect a failure. That is no longer true. The daemon log
now holds four `lead-rollover: rolled` lines, and three of them ran on 2026-09-22 at 10:01:43,
10:38:28 and 11:15:47. Each one cleared the old lead and started a fresh session against the
handover file, with the configured `bootstrapText` arriving as its first message. No context was
lost. The old `Unknown command: /clearFresh` failure from 2026-09-12 does not appear in the log
at all. Re-measure both numbers with:
```bash
grep -c "lead-rollover: rolled" fleetd/fleetd.out # successful rolls
grep -c "lead-rollover:" fleetd/fleetd.out # positive control: must be larger
grep -c "Unknown command" fleetd/fleetd.out # the old failure: expect 0
```
Run the control line too. A broken pattern returns a clean `0` that reads exactly like good news.
If the first number stops growing across rolls, or `Unknown command` returns anything above 0,
the bootstrap has regressed and this paragraph is stale again.
**You still write the file before you confirm, and never the other way round.** That order is not
about the bootstrap being unreliable. It is what the daemon checks: the handover file must have
been modified *after* the open request, or `confirm` refuses it as stale.
- **One warning in the log is normal and is not a failure.** Every one of the three rolls above also
logged `/clear on term_… was never observed as WORKING after 8 consecutive IDLE/DONE polls —
releasing rather than wedging the roll`. The daemon could not see the pane go WORKING after
`/clear`, so it released instead of hanging. The roll then succeeded anyway. That is the safe
branch behaving correctly. Do not report it as a broken roll.
## Writing style
@@ -131,6 +131,27 @@ public final class Fleetd {
}
static void main(String[] args) {
main(args, ResourcePorts.system());
}
/**
* fleetd #625: package-private so a test can drive the literal startup sequence — not a copy
* of it — against a non-production {@link ResourcePorts} whose {@link
* ResourcePorts#environment()} is tainted, without ever touching the real process environment.
* The real process environment is exactly what a test cannot taint from inside the JVM, which
* is why nothing could pin {@link SubscriptionGuard#assertPrimaryClean} running at this call
* site before this ticket.
*
* <p>{@link #main(String[])} is the one production caller, passing {@link
* ResourcePorts#system()}. Every statement below, in the same order, is otherwise unchanged
* from before this ticket — in particular the guard still runs here, before {@code
* cfg.validateAll()} and before {@link FleetdAssembly#assembleAndStart} ever touches a socket,
* a broker, or HTTP, exactly as it always has. {@code ports} is reused for the guard check and
* then handed on to the assembly, rather than a second instance being constructed there, so a
* test's fake backs the whole boot path with one consistent view — see {@code
* FleetdSubscriptionGuardOrderingTest}.
*/
static void main(String[] args, ResourcePorts ports) {
Path configPath = args.length > 0 ? Path.of(args[0]) : chooseDefaultConfigFile(Path.of(""));
// The config file was renamed bridged.yaml -> fleetd.yaml. Name the file we actually
// loaded, whichever of the two names it carries.
@@ -166,7 +187,7 @@ public final class Fleetd {
// The primary/host env that launched fleetd must not be tainted.
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
guard.assertPrimaryClean(System.getenv());
guard.assertPrimaryClean(ports.environment());
// Every FleetConfig.validateXxx() the operator's config can fail — CB-501's auth-exposure
// check, CB-531's lead-tab-prefix check, CB-542's subscription-profile check, the charter
@@ -189,7 +210,9 @@ public final class Fleetd {
// after validateAll() (this line) puts both back under test, in the same relative order,
// before either one does any I/O — see FleetdAssembly's javadoc for the full boot-order
// contract this preserves exactly.
FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ResourcePorts.system());
// fleetd #625: the same `ports` the guard check above just used, not a second
// ResourcePorts.system() instance — see this method's own javadoc.
FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
}
/**
@@ -194,7 +194,7 @@ final class FleetdAssembly {
// CB-504: under supervision (launchd/systemd) fleetd can start before herdr's socket
// exists. Wait, then degrade rather than die: serving with /healthz reporting "degraded" is
// strictly more useful than exiting.
Fleetd.HerdrAwaitOutcome herdrOutcome = Fleetd.awaitHerdr(herdr, ports.nanoClock(), Fleetd::sleepHerdrPoll);
Fleetd.HerdrAwaitOutcome herdrOutcome = Fleetd.awaitHerdr(herdr, ports.nanoClock(), ports.herdrPollWait());
boolean herdrUp = Fleetd.logHerdrWaitOutcomeAndShouldReap(herdrOutcome);
if (herdrUp) {
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died
@@ -392,13 +392,21 @@ final class FleetdAssembly {
if (cfg.leadHeartbeat() != null) {
var hb = cfg.leadHeartbeat();
var leadContextGauge = new LeadContextGauge();
// fleetd #621: the context-high notice's own wording must track this same effective
// value — LeadRollover.confirm(...) already gates the roll on it (LeadRollover.java:480),
// and absent `leadRollover:` entirely the roll is unusable regardless (NOT_CONFIGURED),
// so `true` (the FleetConfig.LeadRollover default) is the safe, byte-identical fallback.
// Carried in from Fleetd.main when #612 Unit A merged main: #622 added this line to the
// block Unit A had already moved here, so the merge would otherwise have silently
// dropped it — with a fully green suite, because nothing pins it (see the follow-up issue).
boolean requireOperatorConfirm = cfg.leadRollover() == null || cfg.leadRollover().requireOperatorConfirm();
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
pushLoop, heartbeatScheduler, ports.nanoClock(),
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
metrics,
Fleetd.leadContextSource(leadContextGauge, router.leadAgents(), leads,
Fleetd.leadConfigDirLookup(() -> config.get().profiles(), leaders)),
Boolean.TRUE.equals(hb.contextHighNudge()));
Boolean.TRUE.equals(hb.contextHighNudge()), requireOperatorConfirm);
heartbeat.start();
} else {
heartbeat = null;
@@ -43,6 +43,19 @@ public interface ResourcePorts {
/** A monotonic elapsed-time clock. Production: {@link System#nanoTime()}. */
LongSupplier nanoClock();
/**
* fleetd #629: the per-poll wait {@code FleetdAssembly#assembleAndStart} passes to {@code
* Fleetd#awaitHerdr} while polling for herdr's socket. Production: {@link
* Fleetd#sleepHerdrPoll()} — a real {@code Thread.sleep}. {@link #nanoClock()} alone is not
* enough to make {@code awaitHerdr}'s deadline controllable: the old call site passed {@code
* Fleetd::sleepHerdrPoll} directly, hardcoded, so a test that injected a fake clock still had
* to wait out the real sleep between each poll to ever reach the deadline — the clock looked
* injected and was not actually controllable. A test supplies a no-op that advances its own
* injected {@link #nanoClock()} instead, so the deadline becomes reachable without any real
* wall-clock time passing.
*/
Runnable herdrPollWait();
/**
* A wall-clock reading, in nanoseconds. Production: {@code System.currentTimeMillis()}
* converted to nanoseconds. Kept separate from {@link #nanoClock()} because {@link
@@ -44,6 +44,11 @@ final class SystemResourcePorts implements ResourcePorts {
return System::nanoTime;
}
@Override
public Runnable herdrPollWait() {
return Fleetd::sleepHerdrPoll;
}
@Override
public LongSupplier wallClockNanos() {
return () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis());
@@ -2240,11 +2240,11 @@ public record FleetConfig(
* information — which profiles, and now both values, so they can fix it without reading the
* source — without ever taking the fleet down.
*
* <p>Which of the two inputs Claude Code actually follows when they disagree is intentionally
* <em>not</em> asserted here. {@code ClaudeCodeArguments}'s javadoc used to state the
* environment variable always wins; nobody had measured that, and this host's own
* {@code fleetd.yaml} asserts the opposite in a comment. This method only detects and reports
* the disagreement — see {@link dev.ltms.fleet.launch.ClaudeCodeArguments}.
* <p>fleetd #618 measured which of the two inputs Claude Code actually follows when they
* disagree: the environment variable wins, so {@code autoCompactWindow} is inert on a profile
* that also sets the env var. This method only detects and reports the disagreement — it does
* not correct it — see {@link dev.ltms.fleet.launch.ClaudeCodeArguments} for the full measured
* precedence.
*
* <p>Equal values never warn: either input then produces the same session window, so there is
* nothing to reconcile.
@@ -2285,9 +2285,11 @@ public record FleetConfig(
names.sort(String::compareTo);
detail.sort(String::compareTo);
log.warn("Claude Code profile(s) {} set disagreeing autoCompactWindow and env."
+ "CLAUDE_CODE_AUTO_COMPACT_WINDOW — the daemon starts anyway. Fix by "
+ "removing one key or setting equal values on each: {}. Which input Claude "
+ "Code actually follows when they disagree is not verified here.",
+ "CLAUDE_CODE_AUTO_COMPACT_WINDOW — the daemon starts anyway: {}. fleetd "
+ "#618 measured that CLAUDE_CODE_AUTO_COMPACT_WINDOW wins, so "
+ "autoCompactWindow is inert on these profiles. Set equal values on each "
+ "to resolve this — do not just delete the env var, since that LOWERS the "
+ "live window to autoCompactWindow's value rather than fixing anything.",
names, String.join(", ", detail));
}
@@ -15,13 +15,15 @@ public final class ClaudeCodeArguments {
* Append the configured Claude Code auto-compaction window when the profile opts in.
*
* <p>This flag and the environment variable {@code CLAUDE_CODE_AUTO_COMPACT_WINDOW} can
* disagree. Which one Claude Code actually follows when they do is NOT verified here — this
* javadoc used to claim the environment variable always wins, but nobody had measured that, and
* this host's own {@code fleetd.yaml} asserts the opposite in a comment. So this javadoc no
* longer picks a side. {@link FleetConfig#load(java.nio.file.Path)} only WARNS when a Claude
* Code profile sets both to different values (see {@code
* disagree, and fleetd #618 measured which one Claude Code actually follows: the environment
* variable wins, ahead of this {@code --autocompact} flag, ahead of the settings file, ahead of
* clientdata, the experiment, and the model default. So when a profile sets both, the flag this
* method appends has NO effect — Claude Code reads {@code CLAUDE_CODE_AUTO_COMPACT_WINDOW}
* first and never consults the flag. {@link FleetConfig#load(java.nio.file.Path)} only WARNS
* when a Claude Code profile sets both to different values (see {@code
* FleetConfig.warnConflictingAutoCompactWindows}) — it does not stop the daemon from starting,
* and a launched session may end up honouring either window.
* and the launched session honours the env var, not this flag. Measured against Claude Code
* 2.1.278 (fleetd #618) — a later version could reorder this precedence.
*/
public static List<String> withAutoCompactWindow(List<String> argv, FleetConfig.Profile profile) {
if (profile.autoCompactWindow() == null) {
@@ -76,6 +76,7 @@ public final class LeadHeartbeatLoop {
private final Metrics metrics; // CB-512 pattern: nullable — no registry in unit tests
private final LeadContextSource contextSource; // fleetd #609
private final boolean contextHighNudge; // fleetd #609: opt-in, like the loop itself
private final boolean requireOperatorConfirm; // fleetd #621: mirrors leadRollover.requireOperatorConfirm
/** When the current idle stretch began (nanos), or {@link #NOT_IDLE}. Single scheduler thread only. */
private long idleSinceNanos = NOT_IDLE;
@@ -106,12 +107,32 @@ public final class LeadHeartbeatLoop {
* fleetd #609: as above, plus the lead's own context source and whether a HIGH reading should
* append a hand-over notice to the loop's nudge. Pass {@link LeadContextSource#none()} and
* {@code false} to keep the pre-#609 behaviour exactly (both existing public constructors do).
*
* <p>fleetd #621: delegates to the full constructor with {@code requireOperatorConfirm=true} —
* the pre-#621 wording ("ask the operator ... only the operator can approve the roll") assumed
* the config default, so every caller of this overload keeps that text byte-identical.
*/
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
ScheduledExecutorService scheduler, LongSupplier clock,
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
LeadContextSource contextSource, boolean contextHighNudge) {
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock,
idleAfterNanos, backoffMs, quietNudgeCap, metrics, contextSource, contextHighNudge, true);
}
/**
* fleetd #621: as above, plus the daemon's effective {@code leadRollover.requireOperatorConfirm}
* value — threaded into {@link #contextNotice(boolean, LeadContextGauge.Reading, boolean, boolean)}
* so the notice's wording tracks the config the daemon actually enforces (see {@code
* LeadRollover.confirm}) instead of always asserting the operator gate is on.
*/
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
ScheduledExecutorService scheduler, LongSupplier clock,
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
LeadContextSource contextSource, boolean contextHighNudge,
boolean requireOperatorConfirm) {
this.primaryRegistry = primaryRegistry;
this.agents = agents;
this.inbox = inbox;
@@ -125,6 +146,7 @@ public final class LeadHeartbeatLoop {
this.metrics = metrics;
this.contextSource = contextSource;
this.contextHighNudge = contextHighNudge;
this.requireOperatorConfirm = requireOperatorConfirm;
}
/**
@@ -344,7 +366,7 @@ public final class LeadHeartbeatLoop {
// d.contextNotified() is the value to persist once delivery is confirmed, not the value the text
// itself should be built from. Otherwise a HIGH stretch that is still latched would never see the
// notice at all, defeating the very check this fixes.
String notice = contextNotice(contextHighNudge, reading, contextNotified);
String notice = contextNotice(contextHighNudge, reading, contextNotified, requireOperatorConfirm);
var lead = primaryRegistry.primaryTerminal();
boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
// The latch becomes true only when all three hold: decide() chose to notify, a notice was
@@ -381,8 +403,9 @@ public final class LeadHeartbeatLoop {
/**
* fleetd #609: the text appended to a nudge when the lead's own context is full — {@code ""}
* whenever the notice does not apply, so callers can unconditionally append this without an extra
* branch. Wording stays plain (CEFR B1) and honest that only the operator approves a roll — this
* loop only ever prints text, it never calls {@code fleet_handover} itself.
* branch. Wording stays plain (CEFR B1) and honest about who actually gates the roll — see the
* {@code requireOperatorConfirm} overload (fleetd #621) for which check that is. This loop only
* ever prints text, it never calls {@code fleet_handover} itself.
*
* @param enabled the {@code leadHeartbeat.contextHighNudge} config flag
* @param reading the lead's current {@link LeadContextGauge} reading
@@ -401,9 +424,33 @@ public final class LeadHeartbeatLoop {
* closing sentence ("You will not be told again until your context reads ok.") false. {@link
* #injectNudge} is the only caller that passes a non-default {@code alreadyNotified}.
*
* <p>fleetd #621: delegates with {@code requireOperatorConfirm=true} — the pre-#621 default and the
* value every existing caller of this overload (including every test written before #621) already
* assumed, so the text this overload returns stays byte-identical.
*
* @param alreadyNotified whether the lead has already been told about the current HIGH stretch
*/
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified) {
return contextNotice(enabled, reading, alreadyNotified, true);
}
/**
* fleetd #621: as {@link #contextNotice(boolean, LeadContextGauge.Reading, boolean)}, but the closing
* instructions also track the daemon's effective {@code leadRollover.requireOperatorConfirm} value,
* instead of always asserting that only the operator can approve the roll.
*
* <p>{@code LeadRollover.confirm(...)} already honours this flag: when it is {@code false}, the daemon
* itself gates the roll on the three handover-file checks alone (exists, modified after the {@code
* open()} request, and no older than {@code maxDocAgeSeconds}) and never consults {@code
* operatorConfirmed}. Before this parameter existed, this notice told the lead to ask the operator
* regardless — so a lead that followed its own instructions asked anyway, and setting the config knob
* to {@code false} stopped the daemon refusing the roll without stopping the operator being
* interrupted. This parameter is how the text is kept honest about which gate is actually live.
*
* @param requireOperatorConfirm the effective {@code leadRollover.requireOperatorConfirm} value
*/
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified,
boolean requireOperatorConfirm) {
if (!enabled || alreadyNotified || reading.state() != LeadContextGauge.State.HIGH) {
return "";
}
@@ -420,10 +467,18 @@ public final class LeadHeartbeatLoop {
sb.append(" (").append(reading.compactions()).append(' ').append(compactionWord)
.append(" so far).");
}
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
+ "again until your context reads ok.");
if (requireOperatorConfirm) {
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
+ "again until your context reads ok.");
} else {
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
+ "write the file it names, then call fleet_handover(action=\"confirm\", token). Decide for "
+ "yourself when to confirm: the roll goes through if the handover file exists, was "
+ "changed after you opened it, and is not older than maxDocAgeSeconds. You will not be "
+ "told again until your context reads ok.");
}
return sb.toString();
}
@@ -0,0 +1,183 @@
package dev.ltms.fleet;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.msg.InMemoryReplyInbox;
import dev.ltms.fleet.msg.LeadChannel;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadMessage;
import dev.ltms.fleet.msg.ReplyInbox;
import io.javalin.Javalin;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 step 4 ranks 3 and 8: the assembled daemon must use the AMQP openers from {@link
* ResourcePorts}, and its startup report must describe the object the runtime actually owns. These
* fakes never open a socket.
*/
class FleetdAssemblyAmqpOpenersTest {
private static final String COORD_ID = "assembly-test";
private static final class DurableReplyInbox implements ReplyInbox {
@Override public void own(String target) { }
@Override public void release(String target) { }
@Override public void publish(String target, String msgId, String content) { }
@Override public List<InboxMessage> peek(String target) { return List.of(); }
@Override public boolean ack(String target, String msgId) { return false; }
}
private static final class DurableLeadMailbox implements LeadChannelHandle {
@Override public void publish(String toCoordId, LeadMessage message) { }
@Override public List<LeadMessage> peek() { return List.of(); }
@Override public void ack(String msgId) { }
@Override public String selfCoordId() { return COORD_ID; }
@Override public boolean heldDurable() { return true; }
@Override public MailboxState inspect(String coordId) { return MailboxState.unknown(coordId); }
@Override public void close() { }
}
private static final class RecordingPorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
final DurableReplyInbox replyInbox = new DurableReplyInbox();
final DurableLeadMailbox leadMailbox = new DurableLeadMailbox();
final AtomicInteger replyOpenCalls = new AtomicInteger();
final AtomicInteger mailboxOpenCalls = new AtomicInteger();
final boolean openSucceeds;
RecordingPorts(boolean openSucceeds) {
this.openSucceeds = openSucceeds;
}
@Override public Map<String, String> environment() { return Map.of(); }
@Override public HerdrClient connectHerdr(Path socketPath) { return herdr; }
@Override public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> {
replyOpenCalls.incrementAndGet();
if (!openSucceeds) throw new IllegalStateException("fake reply broker is down");
return replyInbox;
};
}
@Override public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfId, prefetch) -> {
mailboxOpenCalls.incrementAndGet();
if (!openSucceeds) throw new IllegalStateException("fake coordination broker is down");
return leadMailbox;
};
}
@Override public LongSupplier nanoClock() { return System::nanoTime; }
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
@Override public LongSupplier wallClockNanos() { return System::nanoTime; }
@Override public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override public void addShutdownHook(Runnable hook) { }
@Override public void startHttp(Javalin app, String host, int port) { }
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Files.createDirectories(dir);
Path config = dir.resolve("fleetd.yaml");
Files.writeString(config, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
broker:
uri: "amqp://fake-reply-broker/vh"
coordinator:
uri: "amqp://fake-coordination-broker/vh"
selfId: "assembly-test"
""");
return FleetConfig.load(config);
}
private static FleetdRuntime assemble(Path dir, RecordingPorts ports) throws Exception {
FleetConfig cfg = writeConfig(dir);
return FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, new ConfigRef(dir.resolve("fleetd.yaml"), cfg),
new SubscriptionGuard(cfg.guard().hostSet())), ports);
}
private static boolean reportContains(ListAppender<ILoggingEvent> appender, String text) {
return appender.list.stream().map(ILoggingEvent::getFormattedMessage).anyMatch(message -> message.contains(text));
}
@Test
void assembledAmqpOpenersAndTheirReportsAgreeOnDurableAndFallbackStates(@TempDir Path dir) throws Exception {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
Level oldLevel = logger.getLevel();
ListAppender<ILoggingEvent> reports = new ListAppender<>();
reports.start();
logger.setLevel(Level.INFO);
logger.addAppender(reports);
try {
RecordingPorts durablePorts = new RecordingPorts(true);
FleetdRuntime durable = assemble(dir.resolve("durable"), durablePorts);
try {
// Control: this fails loudly if the assembly did not run or used an inert opener.
assertEquals(1, durablePorts.replyOpenCalls.get(), "assembly must call replyInboxOpener once");
assertEquals(1, durablePorts.mailboxOpenCalls.get(), "assembly must call leadMailboxOpener once");
assertSame(durablePorts.replyInbox, durable.replyInbox(),
"the durable reply report must describe the exact inbox the runtime owns");
assertSame(durablePorts.leadMailbox, durable.leadMailbox(),
"the coordination-on report must describe the exact mailbox the runtime owns");
assertNotNull(durable.leadCoordLoop(), "a durable mailbox must start lead coordination");
assertTrue(reportContains(reports, "reply inbox: AMQP broker (durable)"));
assertTrue(reportContains(reports, "lead coordination: ON as coord-id " + COORD_ID));
} finally {
durable.close();
}
reports.list.clear();
RecordingPorts fallbackPorts = new RecordingPorts(false);
FleetdRuntime fallback = assemble(dir.resolve("fallback"), fallbackPorts);
try {
assertEquals(1, fallbackPorts.replyOpenCalls.get(), "assembly must call the failing reply opener once");
assertEquals(1, fallbackPorts.mailboxOpenCalls.get(), "assembly must call the failing mailbox opener once");
assertTrue(fallback.replyInbox() instanceof InMemoryReplyInbox,
"a failed reply opener must make the runtime own the in-memory fallback");
assertNull(fallback.leadMailbox(), "a failed mailbox opener must leave coordination off");
assertNull(fallback.leadCoordLoop(), "coordination must not start without a mailbox");
assertTrue(reportContains(reports, "reply inbox: in-memory (soft-state)"));
assertTrue(reportContains(reports, "lead-to-lead messaging is OFF"));
} finally {
fallback.close();
}
} finally {
logger.detachAppender(reports);
logger.setLevel(oldLevel);
}
}
}
@@ -117,6 +117,14 @@ class FleetdAssemblyConnectionIdentityTest {
public void startHttp(Javalin app, String host, int port) {
// Deliberately never bind — this test never issues a real HTTP request.
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
}
private FleetdRuntime runtime;
@@ -140,6 +140,14 @@ class FleetdAssemblyCoordinatorLifecycleTest {
public void startHttp(Javalin app, String host, int port) {
// No real HTTP bind in a unit test.
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
}
private static final class SentinelReplyInbox implements ReplyInbox {
@@ -8,6 +8,7 @@ import dev.ltms.fleet.herdr.HerdrClient;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import org.junit.jupiter.api.io.TempDir;
import java.net.URI;
@@ -21,6 +22,8 @@ import java.util.Map;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -69,13 +72,33 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
* below is what this class relies on for CB-185's {@code FleetApp} half; {@code
* FleetAppTwoDaemonTest} remains the full behavioural proof that {@code FleetApp} itself merges
* {@code /sessions} correctly once handed two clients.
*
* <p><strong>fleetd #629 follow-up.</strong> The fix below (see {@link TwoHerdrResourcePorts})
* makes {@link #healthzGoesRedWhenTheLeadDaemonIsDownEvenThoughTheMemberIsUp}'s fake {@code
* nanoClock()} frozen unless {@code herdrPollWait()} itself advances it. That is a sharper pin
* than an assertion — if a future edit to {@code FleetdAssembly} ever bypasses {@code
* ports.herdrPollWait()} again (e.g. reverting to a hardcoded {@code Thread.sleep}), the clock
* never advances, {@code Fleetd#awaitHerdr}'s deadline is never reached, and this test hangs
* forever instead of failing — proven by deliberately reintroducing that exact regression while
* fixing this ticket. {@code @Timeout} turns that silent hang into a bounded, named test failure:
* {@code SEPARATE_THREAD} so JUnit's timeout governor can actually interrupt a thread stuck in a
* real {@code Thread.sleep} loop (the default {@code SAME_THREAD} mode cannot — it only measures
* elapsed time after the test method returns on its own, which never happens here). 10 seconds is
* roughly 150x the real passing times measured here (~0.06s), so a slow CI machine has no reason
* to flake, and it is still 3x faster than discovering the regression by burning a CI job's whole
* wall-clock budget.
*/
@Timeout(value = 10, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
class FleetdAssemblyFleetAppTest {
private static final class TwoHerdrResourcePorts implements ResourcePorts {
final Map<Path, HerdrClient> herdrsBySocket = new LinkedHashMap<>();
final CopyOnWriteArrayList<ScheduledExecutorService> schedulers = new CopyOnWriteArrayList<>();
// fleetd #629: a fake, advanceable clock — NOT System::nanoTime. awaitHerdr's poll wait
// (herdrPollWait() below) advances this on every poll instead of sleeping for real, so the
// down-lead test below reaches awaitHerdr's deadline without burning real wall-clock time.
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); // arbitrary non-zero start
Runnable shutdownHook;
@Override
@@ -107,7 +130,7 @@ class FleetdAssemblyFleetAppTest {
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
return nowNanos::get;
}
@Override
@@ -115,6 +138,13 @@ class FleetdAssemblyFleetAppTest {
return System::nanoTime;
}
@Override
public Runnable herdrPollWait() {
// fleetd #629: advance the fake clock instead of a real Thread.sleep, so awaitHerdr's
// deadline is reached in real time regardless of the configured poll interval.
return () -> nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(1));
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
@@ -0,0 +1,198 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.health.FleetHealthMonitor;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.ReplyInbox;
import io.javalin.Javalin;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Field;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.BiConsumer;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 rank 7 — {@code FleetdAssembly.java:429} wires {@link FleetHealthMonitor}'s {@code
* failTarget} callback with {@code Fleetd.healthFailTarget(messages)}. {@link
* FleetdHealthFailTargetWiringTest} already pins that the FACTORY itself delegates to {@code
* messages::abandon}, but it calls {@code Fleetd.healthFailTarget} directly — it never drives {@code
* FleetdAssembly.assembleAndStart} and so cannot see whether the real call site at {@code :429}
* still passes it the real, assembled {@link MessageService}. Swapping that argument for a no-op
* {@code (a, b) -> {}} compiles clean and leaves the whole suite — including the factory-level test
* — green: a dead member's waiting ticket then sits {@code PENDING} for the full 30-minute async
* timeout instead of failing immediately.
*
* <p>This test assembles the real daemon with {@code health.enabled: true}, pulls the REAL {@code
* failTarget} {@link BiConsumer} out of the REAL, assembled {@link FleetHealthMonitor} (via
* reflection — the field is package-private to {@code dev.ltms.fleet.health}, and nothing public
* exposes it; {@code StatusPollerResilienceTest} already uses the same technique in this suite), and
* invokes it directly against the REAL {@link MessageService} {@link FleetdRuntime#messages()}
* returns. A no-op lambda swapped in at the call site leaves the ticket {@code PENDING} forever,
* which this test catches; the real one fails it.
*/
class FleetdAssemblyHealthFailTargetBehaviouralTest {
private static final String TARGET = "term_a";
private static final class RecordingResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> {
throw new UnsupportedOperationException("no broker: block is configured");
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException("no coordinator: block is configured");
};
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
health:
enabled: true
intervalSeconds: 30
""");
return FleetConfig.load(f);
}
@SuppressWarnings("unchecked")
@Test
@DisplayName("[BEHAVIOURAL] the real assembled FleetHealthMonitor's failTarget reaches the real "
+ "MessageService.abandon, not a no-op")
void assembledHealthFailTargetReachesRealMessages(@TempDir Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
RecordingResourcePorts ports = new RecordingResourcePorts();
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
// Surefire runs the whole suite in one JVM fork, so the scheduler/loops this assembly starts
// (SessionReaper, StatusPoller, the health monitor) must be torn down here, on the failure
// path too — hence the try/finally, not just a statement at the end of the happy path.
try {
FleetHealthMonitor healthMonitor = runtime.healthMonitor();
assertNotNull(healthMonitor, "health.enabled: true in this test's config, so "
+ "FleetdAssembly.assembleAndStart must have built a real FleetHealthMonitor");
Field field = FleetHealthMonitor.class.getDeclaredField("failTarget");
field.setAccessible(true);
BiConsumer<String, String> failTarget = (BiConsumer<String, String>) field.get(healthMonitor);
assertNotNull(failTarget, "FleetHealthMonitor's failTarget must never be null — the "
+ "constructor itself requires it");
MessageService messages = runtime.messages();
// --- loud control: prove the assembled MessageService is actually wired up and a ticket is
// genuinely PENDING before failTarget ever runs. If this fails, the test below would pass
// vacuously on a MessageService that never got a ticket in the first place. TARGET has no
// live agent behind it (no session was ever acquired), so nothing resolves this ticket on
// its own — it stays PENDING until failTarget (or a timeout) ends it.
String ticket = messages.sendAsync(TARGET, "long task");
MessageService.TaskView before = messages.poll(ticket);
assertEquals(MessageService.Phase.PENDING, before.phase(),
"control: the async ticket must be PENDING before failTarget runs");
failTarget.accept(TARGET, "member unreachable (health monitor)");
MessageService.TaskView after = awaitTerminal(messages, ticket);
assertEquals(MessageService.Phase.FAILED, after.phase(),
"FleetdAssembly.java:429 must pass Fleetd.healthFailTarget(messages) built from the "
+ "SAME assembled MessageService — a no-op BiConsumer at that call site leaves "
+ "this ticket PENDING for the full 30-minute async timeout instead of failing it");
assertTrue(after.detail() != null && after.detail().contains("member unreachable"),
"the failure reason passed to failTarget.accept must reach MessageService.abandon and "
+ "end up in the ticket's detail");
} finally {
// Proof the teardown actually ran, not just an assurance that a finally was added: the
// captured shutdown hook's close order (FleetdAssemblyLifecycleTest) closes the herdr
// client last, so ports.herdr.closed flips to true only if this hook really executed.
ports.shutdownHook.run();
assertTrue(ports.herdr.closed, "the captured shutdown hook must have run and closed herdr — "
+ "proof this test's assembled background loops/scheduler were torn down");
}
}
private static MessageService.TaskView awaitTerminal(MessageService messages, String ticket)
throws InterruptedException {
long deadline = System.currentTimeMillis() + 5000;
MessageService.TaskView view = messages.poll(ticket);
while (view.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(10);
view = messages.poll(ticket);
}
return view;
}
}
@@ -126,6 +126,14 @@ class FleetdAssemblyLifecycleTest {
ledger.add("startHttp");
this.startedApp = app;
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
}
/** A fake {@link ReplyInbox} that is also {@link AutoCloseable}, so the ledger can prove it closes. */
@@ -0,0 +1,232 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.msg.ReplyInbox;
import dev.ltms.fleet.msg.ReplyPushLoop;
import dev.ltms.fleet.session.SessionManager;
import io.javalin.Javalin;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Field;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.Consumer;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 rank 6 — {@code FleetdAssembly.java:447} wires {@code
* sessions.onRelease(Fleetd.releaseCleanup(messages, replyInbox, primaryRegistry))}. {@link
* FleetdReleaseCleanupWiringTest} already pins that the FACTORY {@code Fleetd.releaseCleanup}
* itself reaches all three collaborators — but it calls the factory directly, never {@code
* FleetdAssembly.assembleAndStart}, so it cannot see whether the real call site at {@code :447}
* still registers it (as opposed to a no-op {@code detail -> { }}) or still passes it the REAL,
* assembled {@code messages}/{@code replyInbox}/{@code primaryRegistry}. Swapping the registered
* listener for a no-op at that call site compiles clean and leaves the whole suite — including the
* factory-level test — green: EVERY teardown then leaks a stuck rendezvous waiter, an unreleased
* reply-inbox consumer, and a stale lead binding, all three at once.
*
* <p>This test assembles the real daemon with {@code idleSleepGuard.enabled: false} — the ONLY
* other {@code onRelease} registration in {@code FleetdAssembly} (see {@code
* dev.ltms.fleet.power.IdleSleepGuard}'s own wiring at {@code FleetdAssembly.java:228}) — so the
* real {@link SessionManager}'s release-listener list holds exactly the one listener this call site
* registers. It pulls that REAL listener out via reflection (the list itself is private, like
* {@code StatusPollerResilienceTest}'s use of the same technique elsewhere in this suite), invokes
* it directly, and asserts all three collaborator effects against the REAL, assembled {@link
* MessageService} ({@link FleetdRuntime#messages()}), the REAL {@link ReplyInbox} ({@link
* FleetdRuntime#replyInbox()}), and the REAL {@link PrimaryRegistry} — reached through {@link
* FleetdRuntime#pushLoop()}, the only other accessor that was handed the same {@code
* primaryRegistry} instance ({@code FleetdAssembly.java:382}), since {@code FleetMcp} never exposes
* it.
*/
class FleetdAssemblyReleaseCleanupBehaviouralTest {
private static final String TARGET = "term_a";
private static final class RecordingResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> {
throw new UnsupportedOperationException("no broker: block is configured");
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException("no coordinator: block is configured");
};
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
""");
return FleetConfig.load(f);
}
@SuppressWarnings("unchecked")
@Test
@DisplayName("[BEHAVIOURAL] the real assembled release listener reaches messages.abandon, "
+ "replyInbox.release, AND primaryRegistry.forgetDelegation — all three leaks at once")
void assembledReleaseListenerReachesAllThreeCollaborators(@TempDir Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
RecordingResourcePorts ports = new RecordingResourcePorts();
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
// Surefire runs the whole suite in one JVM fork, so the scheduler/loops this assembly starts
// must be torn down here, on the failure path too — hence the try/finally, not just a
// statement at the end of the happy path.
try {
// --- reach into SessionManager's private release-listener list. idleSleepGuard.enabled:
// false above means FleetdAssembly.java:228 never registers, so this list must hold EXACTLY
// the one listener :447 registers.
Field listenersField = SessionManager.class.getDeclaredField("releaseListeners");
listenersField.setAccessible(true);
List<Consumer<SessionManager.ReleaseDetail>> releaseListeners =
(List<Consumer<SessionManager.ReleaseDetail>>) listenersField.get(runtime.sessions());
assertEquals(1, releaseListeners.size(), "control: with idleSleepGuard.enabled: false, "
+ "FleetdAssembly.java:447 must be the ONLY onRelease registration — a different "
+ "count means this test is no longer isolating the call site it claims to pin");
Consumer<SessionManager.ReleaseDetail> releaseListener = releaseListeners.get(0);
MessageService messages = runtime.messages();
ReplyInbox replyInbox = runtime.replyInbox();
// primaryRegistry is never exposed by FleetdRuntime directly — ReplyPushLoop is the other
// collaborator FleetdAssembly.java:382 hands the SAME instance to, so reach it from there.
Field primaryRegistryField = ReplyPushLoop.class.getDeclaredField("primaryRegistry");
primaryRegistryField.setAccessible(true);
PrimaryRegistry primaryRegistry = (PrimaryRegistry) primaryRegistryField.get(runtime.pushLoop());
assertNotNull(primaryRegistry, "control: the assembled ReplyPushLoop must hold a real "
+ "PrimaryRegistry instance");
// --- loud controls: set up the "before" state each collaborator's effect is measured
// against, against the REAL assembled objects. If any of these three fails, the test below
// would pass vacuously because the subject it claims to observe never existed in the first
// place.
String ticket = messages.sendAsync(TARGET, "long task");
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
"control: the async ticket must be PENDING before the release listener runs");
replyInbox.own(TARGET);
replyInbox.publish(TARGET, "msg-1", "hello");
assertEquals(1, replyInbox.peek(TARGET).size(),
"control: the reply inbox must own TARGET and hold one message before the release "
+ "listener runs");
primaryRegistry.recordDelegation(TARGET, "lead-1");
assertEquals("lead-1", primaryRegistry.nudgeTargetFor(TARGET).orElse(null),
"control: the delegation must be recorded before the release listener runs");
// --- the one call under test: invoke the REAL, assembled release listener directly, the
// same way SessionManager.release(...) would on a real teardown.
releaseListener.accept(new SessionManager.ReleaseDetail(TARGET, null, null, null, null));
MessageService.TaskView after = awaitTerminal(messages, ticket);
assertEquals(MessageService.Phase.FAILED, after.phase(),
"FleetdAssembly.java:447 must register a listener that calls messages.abandon(...) "
+ "on the SAME assembled MessageService — an inert listener leaves this "
+ "ticket PENDING for the full 30-minute async timeout");
assertTrue(after.detail() != null && after.detail().contains("released"),
"the abandon reason must say the worker session was released");
assertTrue(replyInbox.peek(TARGET).isEmpty(),
"FleetdAssembly.java:447 must register a listener that calls replyInbox.release(...) "
+ "— an inert listener leaves the inbox still owning TARGET with its message");
assertTrue(primaryRegistry.nudgeTargetFor(TARGET).isEmpty(),
"FleetdAssembly.java:447 must register a listener that calls "
+ "primaryRegistry.forgetDelegation(...) — an inert listener leaves the stale "
+ "delegation in place");
} finally {
// Proof the teardown actually ran, not just an assurance that a finally was added: the
// captured shutdown hook's close order (FleetdAssemblyLifecycleTest) closes the herdr
// client last, so ports.herdr.closed flips to true only if this hook really executed.
ports.shutdownHook.run();
assertTrue(ports.herdr.closed, "the captured shutdown hook must have run and closed herdr — "
+ "proof this test's assembled background loops/scheduler were torn down");
}
}
private static MessageService.TaskView awaitTerminal(MessageService messages, String ticket)
throws InterruptedException {
long deadline = System.currentTimeMillis() + 5000;
MessageService.TaskView view = messages.poll(ticket);
while (view.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
Thread.sleep(10);
view = messages.poll(ticket);
}
return view;
}
}
@@ -0,0 +1,241 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.lead.LeadContextGauge;
import dev.ltms.fleet.msg.LeadHeartbeatLoop;
import io.javalin.Javalin;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Field;
import java.lang.reflect.Method;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #630, and fleetd #612's ranks for this call site. {@code FleetdAssembly.java:402}
* computes {@code requireOperatorConfirm} from the effective {@code leadRollover.requireOperatorConfirm}
* config, and {@code :409} threads it as the 14th argument into the full {@link LeadHeartbeatLoop}
* constructor. Measured on 26f1986: dropping that one argument so the 13-argument overload is
* selected instead (it delegates with {@code true} hardcoded — see that overload's own javadoc,
* fleetd #621) compiles with 0 errors and leaves all 1883 tests green, both with and without the
* argument. In production this means the daemon keeps starting and keeps nudging, but the
* context-high notice silently goes back to telling EVERY lead to ask the operator before a
* context roll — on a host that set {@code requireOperatorConfirm: false} specifically so it would
* not have to. That is the operator's own fix silently reverting, with a fully green suite.
*
* <p>{@code LeadHeartbeatLoopTest} already proves {@link LeadHeartbeatLoop}'s package-private
* {@code contextNotice(boolean, LeadContextGauge.Reading, boolean, boolean)} branches correctly on
* its own {@code requireOperatorConfirm} argument — that the METHOD works. It says nothing about
* which value {@code FleetdAssembly} actually passes into the constructed loop, so it is not
* reused here as coverage for the call site.
*
* <p>This test assembles the real daemon TWICE — once with {@code leadRollover.requireOperatorConfirm:
* false}, once with {@code true} — pulls the REAL {@code requireOperatorConfirm} field out of the
* REAL, assembled {@link LeadHeartbeatLoop} each time (reflection: the field, and {@code
* contextNotice} itself, are package-private to {@code dev.ltms.fleet.msg}, and nothing public
* exposes either — the same technique {@code StatusPollerResilienceTest} already uses in this
* suite), and calls the REAL {@code contextNotice} method with that field's value to produce the
* actual notice text the assembled loop would append to a nudge. Both directions are asserted: a
* one-directional test here would pass on a constant.
*/
class FleetdAssemblyRequireOperatorConfirmBehaviouralTest {
private static final class RecordingResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> {
throw new UnsupportedOperationException("no broker: block is configured");
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException("no coordinator: block is configured");
};
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
}
}
private static FleetConfig writeConfig(Path dir, boolean requireOperatorConfirm) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
leadHeartbeat:
idleAfterSeconds: 600
backoffMs: 15000
quietNudgeCap: 5
leadRollover:
handoverPath: handover.md
requireOperatorConfirm: %s
""".formatted(requireOperatorConfirm));
return FleetConfig.load(f);
}
/** Carries both the assembled loop under test AND its {@link RecordingResourcePorts}, so the
* caller can tear the assembly down (this test assembles the real daemon TWICE — see the class
* javadoc — and each assembly needs its own teardown, not just the last one). */
private record Assembled(LeadHeartbeatLoop heartbeat, RecordingResourcePorts ports) {
}
private static Assembled assembleHeartbeat(Path dir, boolean requireOperatorConfirm) throws Exception {
FleetConfig cfg = writeConfig(dir, requireOperatorConfirm);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
RecordingResourcePorts ports = new RecordingResourcePorts();
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
LeadHeartbeatLoop heartbeat = runtime.heartbeat();
assertNotNull(heartbeat, "control: leadHeartbeat: is configured, so FleetdAssembly.assembleAndStart "
+ "must have built a real LeadHeartbeatLoop");
return new Assembled(heartbeat, ports);
}
/** Pulls the REAL {@code requireOperatorConfirm} field off the REAL, assembled loop. */
private static boolean assembledRequireOperatorConfirm(LeadHeartbeatLoop heartbeat) throws Exception {
Field field = LeadHeartbeatLoop.class.getDeclaredField("requireOperatorConfirm");
field.setAccessible(true);
return field.getBoolean(heartbeat);
}
/** Calls the REAL, package-private {@code contextNotice(boolean, Reading, boolean, boolean)} via reflection. */
private static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified,
boolean requireOperatorConfirm) throws Exception {
Method method = LeadHeartbeatLoop.class.getDeclaredMethod("contextNotice", boolean.class,
LeadContextGauge.Reading.class, boolean.class, boolean.class);
method.setAccessible(true);
return (String) method.invoke(null, enabled, reading, alreadyNotified, requireOperatorConfirm);
}
@Test
@DisplayName("[BEHAVIOURAL] the real assembled LeadHeartbeatLoop's context-high notice tracks "
+ "leadRollover.requireOperatorConfirm — BOTH directions")
void assembledRequireOperatorConfirmControlsNoticeWording(@TempDir Path dir) throws Exception {
LeadContextGauge.Reading highReading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH,
250_000L, 2);
String noticeFalse;
String noticeTrue;
// --- direction 1: requireOperatorConfirm: false -----------------------------------------
Path falseDir = dir.resolve("false");
Files.createDirectories(falseDir);
Assembled assembledFalse = assembleHeartbeat(falseDir, false);
// Surefire runs the whole suite in one JVM fork, so each assembly's scheduler/loops must be
// torn down here, on the failure path too — hence try/finally per assembly (this test
// assembles TWICE, so both need their own teardown, not just the last one).
try {
boolean fieldFalse = assembledRequireOperatorConfirm(assembledFalse.heartbeat());
assertFalse(fieldFalse, "FleetdAssembly.java:402/:409 must thread leadRollover."
+ "requireOperatorConfirm: false into the assembled LeadHeartbeatLoop's own field — "
+ "dropping the 14th constructor argument selects the 13-argument overload, which "
+ "hardcodes true regardless of config (fleetd #621), and this would read true instead");
noticeFalse = contextNotice(true, highReading, false, fieldFalse);
assertTrue(noticeFalse.contains("Decide for yourself when to confirm"),
"with requireOperatorConfirm: false, the assembled loop's own notice must tell the "
+ "lead it can decide for itself — got: " + noticeFalse);
assertFalse(noticeFalse.contains("ask the operator") || noticeFalse.contains("Only the operator"),
"with requireOperatorConfirm: false, the assembled loop's own notice must NOT ask the "
+ "operator — got: " + noticeFalse);
} finally {
// Proof the teardown actually ran, not just an assurance that a finally was added: the
// captured shutdown hook's close order (FleetdAssemblyLifecycleTest) closes the herdr
// client last, so ports.herdr.closed flips to true only if this hook really executed.
assembledFalse.ports().shutdownHook.run();
assertTrue(assembledFalse.ports().herdr.closed, "the captured shutdown hook must have run "
+ "and closed herdr — proof this assembly's background loops/scheduler were torn down");
}
// --- direction 2: requireOperatorConfirm: true -------------------------------------------
Path trueDir = dir.resolve("true");
Files.createDirectories(trueDir);
Assembled assembledTrue = assembleHeartbeat(trueDir, true);
try {
boolean fieldTrue = assembledRequireOperatorConfirm(assembledTrue.heartbeat());
assertTrue(fieldTrue, "FleetdAssembly.java:402/:409 must thread leadRollover."
+ "requireOperatorConfirm: true into the assembled LeadHeartbeatLoop's own field");
noticeTrue = contextNotice(true, highReading, false, fieldTrue);
assertTrue(noticeTrue.contains("ask the operator") && noticeTrue.contains("Only the operator can approve the roll"),
"with requireOperatorConfirm: true, the assembled loop's own notice must ask the "
+ "operator — got: " + noticeTrue);
assertFalse(noticeTrue.contains("Decide for yourself when to confirm"),
"with requireOperatorConfirm: true, the assembled loop's own notice must NOT tell "
+ "the lead it can decide for itself — got: " + noticeTrue);
} finally {
assembledTrue.ports().shutdownHook.run();
assertTrue(assembledTrue.ports().herdr.closed, "the captured shutdown hook must have run "
+ "and closed herdr — proof this assembly's background loops/scheduler were torn down");
}
// --- the two directions must actually differ: a constant return would pass both assertion
// blocks above vacuously if they happened to share wording, so compare them directly too.
assertTrue(!noticeFalse.equals(noticeTrue),
"the two directions must produce genuinely different notice text — got the same "
+ "text for both: " + noticeFalse);
}
}
@@ -98,6 +98,14 @@ class FleetdAssemblyRoleFallbackBoundaryTest {
public void startHttp(Javalin app, String host, int port) {
// No real HTTP bind in a unit test.
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
}
private static final class SentinelReplyInbox implements ReplyInbox {
@@ -100,6 +100,14 @@ class FleetdBackendQuarantineAssemblyTest {
@Override
public void startHttp(Javalin app, String host, int port) {
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
}
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
@@ -126,6 +126,14 @@ class FleetdCompletionResolverAssemblyTest {
public void startHttp(Javalin app, String host, int port) {
// Deliberately never bind a real port.
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
}
private static FleetConfig writeConfig(Path dir, String profilesYaml, String extraGuardHost,
@@ -0,0 +1,218 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TurnToken;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.session.MemberSession;
import io.javalin.Javalin;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.OptionalLong;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
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.assertTrue;
/**
* fleetd #612 step 4, ranks 1 and 2 (publish side) — {@link FleetdAssembly} lines
* {@code liveExhaustedPatterns}/{@code exhaustedPatterns} (CB-578 stage A, the ticket's own "worst
* consequence in the whole sweep": a genuine usage-limit refusal handed back to a waiting caller
* AS REAL COMPLETED WORK) and {@code Fleetd.publishExhaustionSink(...)} (CB-578 stage B: the
* credential that hit the limit is never quarantined). None of these three lines is driven by an
* existing test through the real assembly: {@code FleetdExhaustedPatternLookupWiringTest} and
* {@code FleetdLiveExhaustedPatternsWiringTest} (fleetd #589) call {@code Fleetd.liveExhaustedPatterns}
* / {@code Fleetd.exhaustedPatternLookup} directly as factories, never through {@link
* FleetdAssembly#assembleAndStart} — they prove the FACTORY classifies correctly, never that THIS
* call site is the one that actually got wired into the running {@link CompletionResolver}. {@link
* FleetdBackendQuarantineAssemblyTest} drives {@code BackendQuarantine.withEscalation(...)}
* directly, a different call site from {@code publishExhaustionSink} here.
*
* <p>This test drives the REAL assembled {@link CompletionResolver} ({@link
* FleetdRuntime#completion()}) with a profile carrying a configured {@code exhaustedPattern},
* through a pane scrape that matches it, and asserts both halves of the production consequence:
* (1) the resolution is {@link Rendezvous.Kind#BACKEND_EXHAUSTED}, never a plain completion handed
* back as real work, and (2) the profile's credential is actually quarantined afterward, through
* the REAL {@link BackendQuarantine} the same assembly built ({@link
* FleetdRuntime#mcp()}{@code .quarantineSource().quarantine()}) — never a copy.
*
* <p>Same {@code ControllableResourcePorts} shape as {@code FleetdCompletionResolverAssemblyTest}:
* a fake, advanceable {@code nanoClock} so {@code CompletionResolver.MIN_TURN_NANOS} clears without
* a real sleep, and {@link FakeHerdr#readText} to drive the pane scrape.
*/
class FleetdExhaustedPatternAssemblyTest {
private static final class ControllableResourcePorts implements ResourcePorts {
final FakeHerdr herdr;
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); // arbitrary non-zero start
Runnable shutdownHook;
ControllableResourcePorts(FakeHerdr herdr) {
this.herdr = herdr;
}
void advanceSeconds(long seconds) {
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds));
}
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> {
throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block");
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException("leadMailboxOpener must not be called — no coordinator: block");
};
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
@Override
public LongSupplier nanoClock() {
return nowNanos::get;
}
@Override
public LongSupplier wallClockNanos() {
return nowNanos::get;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// Deliberately never bind a real port.
}
}
private static FleetConfig writeConfig(Path dir, int cooldownSeconds) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
quarantineCooldownSeconds: %d
profiles:
exhaustprofile:
baseUrl: http://exhausthost.local:8000
model: sonnet
exhaustedPattern: "usage limit reached"
guard:
offSubscriptionHosts:
- exhausthost.local
""".formatted(cooldownSeconds));
return FleetConfig.load(f);
}
/**
* fleetd #589's own description of this gap ({@code Fleetd#exhaustedPatternLookup}'s javadoc):
* "the worst consequence in the whole #589 sweep" — a genuine usage-limit refusal stops being
* classified as {@code BACKEND_EXHAUSTED} and is handed back to a waiting {@code fleet_send} as
* if it were real completed work. Pins {@code FleetdAssembly}'s {@code liveExhaustedPatterns}
* AND {@code exhaustedPatterns} lines (rank 1) together with {@code publishExhaustionSink}
* (rank 2, the non-OpenCode half) in one flow: classify, then quarantine.
*/
@Test
@DisplayName("[BEHAVIOURAL] a scrape matching the profile's exhaustedPattern resolves "
+ "BACKEND_EXHAUSTED (never a plain completion) and quarantines the credential")
void assembledResolverClassifiesExhaustionAndQuarantinesTheCredential(@TempDir Path dir) throws Exception {
int cooldownSeconds = 120;
FleetConfig cfg = writeConfig(dir, cooldownSeconds);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr());
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
try {
MemberSession session = runtime.sessions().acquire("exhaustprofile", null, dir.toString(), null);
String target = session.terminalId();
CompletionResolver completion = runtime.completion();
CompletableFuture<Rendezvous.Resolution> waiter = new CompletableFuture<>();
ports.herdr.readText("idle, nothing yet");
completion.onDelivered(target, new TurnToken(target, waiter, null));
// The matched text must START the pane line (CompletionResolver.startsWithExhaustion) —
// no preceding sentence — for the quarantine side-effect to fire, same as production.
ports.herdr.readText("usage limit reached: try again in a few hours");
ports.advanceSeconds(3); // clear CompletionResolver.MIN_TURN_NANOS (2s), no real sleep
completion.resolveBeforePostAction(target);
// CONTROL: the waiter must have resolved synchronously at all — if the assembled
// CompletionResolver were never actually driven (e.g. a wiring break upstream silently
// left the resolver unreachable), this fails loudly before the real assertions below
// ever run, rather than passing on an untouched waiter.
Rendezvous.Resolution resolution = waiter.getNow(null);
assertTrue(resolution != null, "CONTROL: the waiter must have resolved synchronously — "
+ "if this is null, the assembled resolver was never actually exercised");
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, resolution.kind(),
"a scrape matching the profile's configured exhaustedPattern must classify as "
+ "BACKEND_EXHAUSTED, not a plain completion handed back as real work — "
+ "replacing FleetdAssembly's liveExhaustedPatterns/exhaustedPatterns "
+ "lines with their inert forms (Map.of() / target -> null) must fail "
+ "this assertion; got: " + resolution);
assertTrue(resolution.text().contains("usage limit reached"), resolution.text());
BackendQuarantine quarantine = runtime.mcp().quarantineSource().quarantine();
assertTrue(quarantine.isQuarantined("exhaustprofile"),
"the real publishExhaustionSink-built sink must have quarantined the profile's "
+ "credential (effectiveCredentialId() == the profile name here, no "
+ "credentialId configured) — replacing FleetdAssembly's "
+ "publishExhaustionSink call site with a hardcoded ExhaustionSink.none() "
+ "must fail this assertion, since nothing would ever call "
+ "quarantine.quarantine(...)");
OptionalLong remaining = quarantine.remainingSeconds("exhaustprofile");
assertTrue(remaining.isPresent() && remaining.getAsLong() > 0
&& remaining.getAsLong() <= cooldownSeconds,
"a fresh quarantine must block for at most the configured base cooldown: " + remaining);
} finally {
if (ports.shutdownHook != null) ports.shutdownHook.run();
}
}
}
@@ -115,6 +115,14 @@ class FleetdLeadRolloverAssemblyTest {
@Override
public void startHttp(Javalin app, String host, int port) {
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
}
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
@@ -90,6 +90,14 @@ class FleetdLeadSeatAssemblyTest {
@Override
public void startHttp(Javalin app, String host, int port) {
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
}
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
@@ -0,0 +1,197 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.member.HerdrPeerLauncher;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.placement.BackendQuarantine;
import io.javalin.Javalin;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Field;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
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.assertTrue;
/**
* fleetd #612 step 4, rank 2 (OpenCode half) — {@link FleetdAssembly}'s {@code
* forwardingExhaustionSink} line ({@code Fleetd.forwardingExhaustionSink(exhaustionSinkRef)}),
* handed to {@link dev.ltms.fleet.member.OpenCodeLauncher} so its fleetd #175 model-mismatch check
* can quarantine a credential before {@code sessions} exists to build the real sink (the
* construction-order cycle documented at that call site). The ticket calls this independent from
* {@code publishExhaustionSink} (pinned by {@link FleetdExhaustedPatternAssemblyTest}): a credential
* that hits a usage limit through THIS path is never quarantined if {@code forwardingExhaustionSink}
* is swapped for a hardcoded {@link ExhaustionSink#none()} at that call site — the OpenCode
* launcher's own quarantine check keeps compiling and keeps "running", but it permanently talks to
* a sink that does nothing, independent of whatever {@code publishExhaustionSink} does later.
*
* <p>{@code FleetdExhaustionSinkForwardingWiringTest} (fleetd #589) already proves {@code
* Fleetd.forwardingExhaustionSink(ref)} forwards to whatever {@code ref} holds — as a bare factory
* call, never through {@link FleetdAssembly#assembleAndStart}. It proves nothing about whether
* THIS call site is the one FleetdAssembly actually wires into the real {@code OpenCodeLauncher}
* it builds, which is exactly the #602/#606-shaped gap this ticket exists to close.
*
* <p>No accessor on {@link FleetdRuntime} reaches the adapter instances (by design — see that
* class's own javadoc: only the final collaborators it owns directly are exposed), so this test
* reaches the REAL, assembled {@code OpenCodeLauncher}'s {@code exhaustionSink} field the same way
* {@code SessionManager}/{@code CompositePeerLauncher} wire it internally: a short, targeted
* reflective walk ({@code SessionManager.launcher} → {@code CompositePeerLauncher.byProfile} →
* {@code OpenCodeLauncher.exhaustionSink}) onto the exact object the assembly built — never a copy,
* and never a read of the source text. Reflection is used the same way elsewhere in this suite
* (e.g. {@code StatusPollerWatchdogTest}) to reach a private collaborator a production constructor
* intentionally does not expose a public accessor for.
*/
class FleetdOpenCodeExhaustionForwardingAssemblyTest {
private static final class ControllableResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L);
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> {
throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block");
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException("leadMailboxOpener must not be called — no coordinator: block");
};
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
@Override
public LongSupplier nanoClock() {
return nowNanos::get;
}
@Override
public LongSupplier wallClockNanos() {
return nowNanos::get;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// Deliberately never bind a real port.
}
}
private static FleetConfig writeConfig(Path dir, int cooldownSeconds) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
quarantineCooldownSeconds: %d
profiles:
gemini:
kind: opencode
model: google/gemini-2.5-pro
""".formatted(cooldownSeconds));
return FleetConfig.load(f);
}
/** Reach a declared field by name on {@code target}'s runtime class, bypassing the access check. */
private static Object readField(Object target, Class<?> declaringClass, String fieldName) throws Exception {
Field field = declaringClass.getDeclaredField(fieldName);
field.setAccessible(true);
return field.get(target);
}
@Test
@DisplayName("[BEHAVIOURAL] the real assembled OpenCodeLauncher's exhaustionSink field forwards "
+ "an onExhausted call into the real daemon's BackendQuarantine")
void assembledOpenCodeLauncherExhaustionSinkQuarantinesTheCredential(@TempDir Path dir) throws Exception {
int cooldownSeconds = 90;
FleetConfig cfg = writeConfig(dir, cooldownSeconds);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
ControllableResourcePorts ports = new ControllableResourcePorts();
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
try {
PeerLauncher launcherField = (PeerLauncher) readField(runtime.sessions(),
runtime.sessions().getClass(), "launcher");
// CONTROL: the composite launcher must actually be the real production type with a
// "gemini" -> OpenCodeLauncher entry — if this fails, nothing below exercised the real
// assembly at all, rather than silently passing on an empty/wrong object.
assertTrue(launcherField instanceof CompositePeerLauncher,
"CONTROL: SessionManager.launcher must be the real CompositePeerLauncher the "
+ "assembly built, got: " + launcherField);
@SuppressWarnings("unchecked")
Map<String, HerdrPeerLauncher> byProfile = (Map<String, HerdrPeerLauncher>)
readField(launcherField, CompositePeerLauncher.class, "byProfile");
HerdrPeerLauncher adapter = byProfile.get("gemini");
assertTrue(adapter != null && adapter.getClass().getSimpleName().equals("OpenCodeLauncher"),
"CONTROL: the 'gemini' profile must resolve to a real OpenCodeLauncher adapter, "
+ "got: " + adapter);
ExhaustionSink sink = (ExhaustionSink) readField(adapter, adapter.getClass(), "exhaustionSink");
assertTrue(sink != null, "CONTROL: OpenCodeLauncher.exhaustionSink must never be null");
// The exact call OpenCodeLauncher.SessionAwareHandle#checkModelMatch makes on a real
// model mismatch (fleetd #175): target, reason, and its own already-known profile name.
sink.onExhausted("term_gemini_1", "opencode model mismatch (test)", "gemini");
BackendQuarantine quarantine = runtime.mcp().quarantineSource().quarantine();
assertTrue(quarantine.isQuarantined("gemini"),
"the real forwardingExhaustionSink-wired field must have delegated into the "
+ "published production sink, which quarantines the profile's credential "
+ "('gemini' here — no credentialId configured) — replacing "
+ "FleetdAssembly's forwardingExhaustionSink call site with a hardcoded "
+ "ExhaustionSink.none() must fail this assertion, since the field read "
+ "above would then BE the inert no-op and nothing would ever reach "
+ "quarantine.quarantine(...)");
assertEquals(cooldownSeconds, quarantine.remainingSeconds("gemini").orElseThrow(
() -> new AssertionError("credential must report a remaining cooldown")));
} finally {
if (ports.shutdownHook != null) ports.shutdownHook.run();
}
}
}
@@ -0,0 +1,202 @@
package dev.ltms.fleet;
import dev.ltms.fleet.guard.GuardException;
import dev.ltms.fleet.herdr.HerdrClient;
import io.javalin.Javalin;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #625: pins {@link dev.ltms.fleet.guard.SubscriptionGuard#assertPrimaryClean}'s call site
* in {@link Fleetd#main(String[])} — the ONE place it runs at startup, and the check behind the
* bridge charter's invariant 1 (never let the primary carry {@code ANTHROPIC_BASE_URL}). Nothing
* pinned it before this ticket: deleting {@code guard.assertPrimaryClean(...)} from {@code main}
* left the full suite green, because the call site read the real process environment ({@code
* System.getenv()}), which a test cannot taint from inside the JVM.
*
* <p>{@link Fleetd#main(String[], ResourcePorts)} (added by this ticket) is the literal production
* sequence — not a copy of it — driven here with a {@link ResourcePorts} whose {@link
* ResourcePorts#environment()} is a plain {@code Map} a test controls. The guard itself was
* already pinned by {@code SubscriptionGuardTest}, directly, with a {@code Map} — that proves the
* method's behaviour, not that {@code main} still calls it at the right point. This class pins the
* call site and, separately, the ORDER: the guard must still run before {@code cfg.validateAll()}
* and before {@code FleetdAssembly.assembleAndStart} touches a socket, a broker, or HTTP — not just
* be present somewhere in {@code main}.
*
* <p>Presence alone is not enough (a fix that pins only presence trades an invisible deletion for
* an invisible reordering), so each test below is built so that EITHER deleting the guard call OR
* moving it later makes the <em>same</em> test fail — with a different exception type than the one
* asserted, not a vacuous pass. See each test's own javadoc for how.
*/
class FleetdSubscriptionGuardOrderingTest {
private static final Map<String, String> TAINTED_ENV =
Map.of("ANTHROPIC_BASE_URL", "http://tainted.example");
private static final Map<String, String> CLEAN_ENV = Map.of("PATH", "/usr/bin");
/**
* A {@link ResourcePorts} whose {@link #environment()} is fixed to whatever the test hands it,
* and whose every other method refuses to be called at all. That refusal is the ordering pin:
* if {@code main} ever reaches {@link FleetdAssembly#assembleAndStart} before the guard has had
* a chance to throw, the very first thing the assembly does with {@code ports} is {@link
* #connectHerdr} — so a test that expects {@link GuardException} and instead observes {@link
* UnsupportedOperationException} has just caught the guard running too late (or not at all).
*/
private static final class FixedEnvPorts implements ResourcePorts {
private final Map<String, String> env;
FixedEnvPorts(Map<String, String> env) {
this.env = env;
}
@Override
public Map<String, String> environment() {
return env;
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
throw new UnsupportedOperationException("connectHerdr must not be called before the guard runs");
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
throw new UnsupportedOperationException("replyInboxOpener must not be called before the guard runs");
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
throw new UnsupportedOperationException("leadMailboxOpener must not be called before the guard runs");
}
@Override
public LongSupplier nanoClock() {
throw new UnsupportedOperationException("nanoClock must not be called before the guard runs");
}
@Override
public LongSupplier wallClockNanos() {
throw new UnsupportedOperationException("wallClockNanos must not be called before the guard runs");
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
throw new UnsupportedOperationException("newScheduler must not be called before the guard runs");
}
@Override
public void addShutdownHook(Runnable hook) {
throw new UnsupportedOperationException("addShutdownHook must not be called before the guard runs");
}
@Override
public void startHttp(Javalin app, String host, int port) {
throw new UnsupportedOperationException("startHttp must not be called before the guard runs");
}
@Override
public Runnable herdrPollWait() {
throw new UnsupportedOperationException("herdrPollWait must not be called before the guard runs");
}
}
/**
* Otherwise-invalid: {@code bind.host: 0.0.0.0} with no {@code auth.mode: token} fails {@code
* cfg.validateAll()} (CB-501's auth-exposure check — the same fixture {@code
* FleetdStartupValidationTest#mainRefusesANonLoopbackBindWithoutTokenMode} uses), with an
* {@link IllegalStateException}. That is deliberate: it is what {@code main} would throw INSTEAD
* of {@link GuardException} if the guard call were deleted, or moved to run after {@code
* validateAll()} — a different, distinguishable exception type.
*/
private static Path invalidConfig(Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 0.0.0.0
port: 8765
""");
return f;
}
/** Passes {@code cfg.validateAll()} cleanly — nothing here trips any of its checks. */
private static Path validConfig(Path dir) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
""");
return f;
}
/**
* Pins the order against {@code cfg.validateAll()}. The environment is tainted and the config
* is otherwise invalid (see {@link #invalidConfig}). If the guard runs first (the required
* order), {@code main} throws {@link GuardException} before {@code validateAll()} is ever
* reached. If the guard were deleted, or reordered to run after {@code validateAll()}, {@code
* validateAll()} throws {@link IllegalStateException} instead and this assertion fails on the
* wrong exception type.
*/
@Test
void mainRefusesATaintedEnvironmentBeforeValidatingTheConfig(@TempDir Path dir) throws Exception {
Path config = invalidConfig(dir);
FixedEnvPorts ports = new FixedEnvPorts(TAINTED_ENV);
GuardException ex = assertThrows(GuardException.class,
() -> Fleetd.main(new String[]{config.toString()}, ports));
assertTrue(ex.getMessage().contains("tainted"),
"expected the primary-taint message, got: " + ex.getMessage());
}
/**
* Pins the order against {@code FleetdAssembly.assembleAndStart}. The environment is tainted
* and the config is otherwise VALID (see {@link #validConfig}), so {@code cfg.validateAll()}
* passes silently and the next thing that could possibly run is the assembly's first socket
* call. If the guard runs first (the required order), {@code main} throws {@link
* GuardException} before assembly starts. If the guard were deleted, or reordered to run after
* assembly begins touching {@code ports}, {@link FixedEnvPorts#connectHerdr} throws {@link
* UnsupportedOperationException} instead and this assertion fails on the wrong exception type.
*/
@Test
void mainRefusesATaintedEnvironmentBeforeAssemblyTouchesAnyPort(@TempDir Path dir) throws Exception {
Path config = validConfig(dir);
FixedEnvPorts ports = new FixedEnvPorts(TAINTED_ENV);
GuardException ex = assertThrows(GuardException.class,
() -> Fleetd.main(new String[]{config.toString()}, ports));
assertTrue(ex.getMessage().contains("tainted"),
"expected the primary-taint message, got: " + ex.getMessage());
}
/**
* The CONTROL for the two tests above. Same otherwise-valid config, same {@link FixedEnvPorts}
* whose every method but {@code environment()} refuses to be called — but a CLEAN environment.
* Without this, a guard that always threw {@link GuardException} regardless of input (the
* opposite bug — e.g. the check inverted) would make the two tests above pass for the wrong
* reason: not because they actually drove a real taint through a real guard, but because
* anything would have thrown {@code GuardException}. Here, with nothing to taint, the guard
* must let {@code main} proceed into {@code cfg.validateAll()} and on into the real assembly,
* which reaches {@code ports.connectHerdr} — and THAT throws. A loud, positive assertion: if
* the boot path never actually ran this far, there is no {@link UnsupportedOperationException}
* to catch, only a quiet, unexpected hang or an unrelated early failure.
*/
@Test
void mainProceedsPastTheGuardOnACleanEnvironment(@TempDir Path dir) throws Exception {
Path config = validConfig(dir);
FixedEnvPorts ports = new FixedEnvPorts(CLEAN_ENV);
UnsupportedOperationException ex = assertThrows(UnsupportedOperationException.class,
() -> Fleetd.main(new String[]{config.toString()}, ports));
assertTrue(ex.getMessage().contains("connectHerdr"),
"expected forward progress to reach the assembly's first port call, got: " + ex.getMessage());
}
}
@@ -476,6 +476,26 @@ class LeadHeartbeatLoopTest {
assertTrue(notice.contains("2 compactions"), notice);
}
// ── fleetd #621: the notice must track the effective requireOperatorConfirm value ─────────────
@Test
void contextNoticeKeepsAskingTheOperatorWhenRequireOperatorConfirmIsTrue() {
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
String notice = LeadHeartbeatLoop.contextNotice(true, reading, false, true);
assertTrue(notice.contains("ask the operator"), notice);
assertTrue(notice.contains("Only the operator can approve the roll"), notice);
}
@Test
void contextNoticeDropsTheOperatorAskWhenRequireOperatorConfirmIsFalse() {
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
String notice = LeadHeartbeatLoop.contextNotice(true, reading, false, false);
assertFalse(notice.contains("ask the operator"), notice);
assertFalse(notice.contains("Only the operator can approve the roll"), notice);
assertTrue(notice.contains("fleet_handover"), notice);
assertTrue(notice.contains("maxDocAgeSeconds"), notice);
}
// ── fleetd #609 review: the latch must mean "the notice reached the pane" ────────────────────
//
// These four drive LeadHeartbeatLoop.tick() directly (package-private, same reasoning as
@@ -1846,7 +1846,8 @@ class MessageServiceTest {
* but backed by {@link ManualScheduler} instead of a real timer (fleetd #608): its tick never
* fires on its own — a test drives it explicitly via {@link ManualScheduler#runDueTasks()}.
*/
private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler)
private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler,
ReplyPushLoop pushLoop)
implements AutoCloseable {
@Override
public void close() {
@@ -1861,14 +1862,20 @@ class MessageServiceTest {
* {@code anAlreadyCollectedTicketProducesNoNudge}, which sets it to 1 to prove exactly that.
*/
private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs) {
return wireWithManualScheduler(maxReminders, backoffMs, System::nanoTime);
}
/** As above, with an injectable clock for tests that exercise terminal-ticket pruning. */
private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs,
java.util.function.LongSupplier nowNanos) {
PrimaryRegistry registry = new PrimaryRegistry(null);
registry.recordDelegation(T, LEAD);
FakeHerdr leadHerdr = new FakeHerdr();
AgentControl leadAgents = new AgentControl(leadHerdr);
ManualScheduler scheduler = new ManualScheduler();
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, System::nanoTime);
return new ManualPushWiring(service, leadHerdr, scheduler);
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, nowNanos);
return new ManualPushWiring(service, leadHerdr, scheduler, pushLoop);
}
private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException {
@@ -1959,24 +1966,35 @@ class MessageServiceTest {
@Test
void severalAsyncTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception {
try (var wiring = wireWithPushLoop(1, 300)) { // wide backoff: both tickets land before the tick fires
String first = wiring.service().sendAsync(T, "first task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "first done"));
// Settle without polling: poll() itself marks a ticket collected (that's the point of
// anAlreadyCollectedTicketProducesNoNudge above) — using it here to detect completion
// would collect the ticket before the coalescing this test checks ever gets a chance.
Thread.sleep(100);
// A 1ms backoff is due immediately. ManualScheduler still cannot run it until this test
// explicitly calls runDueTasks(), so both terminal tickets join one scheduled tick.
try (var wiring = wireWithManualScheduler(1, 1)) {
String first;
String second;
java.util.concurrent.CountDownLatch firstTerminalReached = new java.util.concurrent.CountDownLatch(1);
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(firstTerminalReached::countDown);
try {
first = wiring.service().sendAsync(T, "first task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "first done"));
assertTrue(firstTerminalReached.await(5, TimeUnit.SECONDS),
"the first ticket never reached its terminal phase");
String second = wiring.service().sendAsync(T, "second task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "second done"));
Thread.sleep(100);
java.util.concurrent.CountDownLatch secondTerminalReached = new java.util.concurrent.CountDownLatch(1);
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(secondTerminalReached::countDown);
second = wiring.service().sendAsync(T, "second task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "second done"));
assertTrue(secondTerminalReached.await(5, TimeUnit.SECONDS),
"the second ticket never reached its terminal phase");
} finally {
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
}
awaitNudge(wiring.leadHerdr());
Thread.sleep(200); // settle — nothing more should arrive beyond the one coalesced nudge
assertEquals(1, wiring.scheduler().runDueTasks(),
"both terminal tickets must coalesce onto one scheduled tick");
long nudgeCount = wiring.leadHerdr().calls.stream()
.filter(c -> c.method().equals("agent.prompt")).count();
assertEquals(1, nudgeCount, "two tickets finishing together must produce ONE nudge, not two");
@@ -2055,7 +2073,8 @@ class MessageServiceTest {
@Test
void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception {
try (var wiring = wireWithPushLoop(5, 50)) {
// A 1ms backoff is due immediately, but ManualScheduler only ticks when this test asks it to.
try (var wiring = wireWithManualScheduler(5, 1)) {
String ticket = wiring.service().sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
@@ -2063,8 +2082,9 @@ class MessageServiceTest {
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
() -> wiring.service().ask(T, "which config file?", 5000));
MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING);
awaitQuestionPendingOn(wiring.pushLoop(), LEAD, asking.turnId());
awaitNudge(wiring.leadHerdr());
assertEquals(1, wiring.scheduler().runDueTasks(), "the open question must have one scheduled tick");
long callsBeforeAnswer = wiring.leadHerdr().calls.stream()
.filter(c -> c.method().equals("agent.prompt")).count();
@@ -2072,12 +2092,22 @@ class MessageServiceTest {
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1);
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(terminalReached::countDown);
assertTrue(rendezvous.resolve(T, "done"));
try {
assertTrue(terminalReached.await(5, TimeUnit.SECONDS),
"the answered ticket never reached its terminal phase");
} finally {
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
}
answer.get(5, TimeUnit.SECONDS);
// Let several more ticks (and the ticket's own now-legitimate terminal nudge) fire —
// none of them may still name the question's turnId, which is closed.
Thread.sleep(300);
// Run the ticket's legitimate terminal nudge and every later scheduled tick through its
// reminder cap. None may still name the closed question.
for (int tick = 0; tick < 6; tick++) {
assertEquals(1, wiring.scheduler().runDueTasks(), "expected one scheduled reminder tick");
}
boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.skip(callsBeforeAnswer)
@@ -2160,35 +2190,45 @@ class MessageServiceTest {
// decideTickets hits the cap and STOPs — activeLeads drops the lead, but (before the fix)
// pendingTickets never drops the ticket. That is the exact "cap already STOPped" branch of
// the bug report, reached deterministically rather than by timing it against a live tick.
try (var wiring = wireWithPushLoop(1, 50, clock::get)) {
String stale = wiring.service().sendAsync(T, "first task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "stale result"));
// A 1ms backoff is due immediately, but ManualScheduler runs only the ticks below.
try (var wiring = wireWithManualScheduler(1, 1, clock::get)) {
String stale;
String fresh;
java.util.concurrent.CountDownLatch staleTerminalReached = new java.util.concurrent.CountDownLatch(1);
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(staleTerminalReached::countDown);
try {
stale = wiring.service().sendAsync(T, "first task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "stale result"));
assertTrue(staleTerminalReached.await(5, TimeUnit.SECONDS),
"the stale ticket never reached its terminal phase");
// Let the reminder loop fire its one nudge and hit the cap (STOP removes it from
// activeLeads; pendingTickets is untouched either way — that asymmetry is the bug).
awaitNudge(wiring.leadHerdr());
Thread.sleep(300);
assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale),
"sanity: the stale ticket's own reminder must have fired first");
// Fire the stale ticket's one nudge, then its cap tick (STOP removes it from activeLeads;
// pendingTickets is untouched either way — that asymmetry is the bug).
assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one nudge tick");
assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one cap tick");
assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale),
"sanity: the stale ticket's own reminder must have fired first");
// Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly.
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
// Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly.
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
// A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes.
String fresh = wiring.service().sendAsync(T, "second task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "fresh result"));
// The fresh ticket restarts the (now-dormant) reminder loop with its own nudge.
long before = wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count();
long deadline = System.currentTimeMillis() + 3000;
while (wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count() <= before
&& System.currentTimeMillis() < deadline) {
Thread.sleep(10);
// A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes.
java.util.concurrent.CountDownLatch freshTerminalReached = new java.util.concurrent.CountDownLatch(1);
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(freshTerminalReached::countDown);
fresh = wiring.service().sendAsync(T, "second task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "fresh result"));
assertTrue(freshTerminalReached.await(5, TimeUnit.SECONDS),
"the fresh ticket never reached its terminal phase");
} finally {
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
}
// The fresh ticket restarts the now-dormant reminder loop with its own nudge.
assertEquals(1, wiring.scheduler().runDueTasks(), "the fresh ticket must have one nudge tick");
String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge);
assertFalse(latestNudge.contains(stale),