Compare commits

...

14 Commits

Author SHA1 Message Date
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 640f4d5f23 Merge pull request 'fleetd #612 B3: behavioural replacements for lead-seat, quarantine, lead-rollover guards' (#628) from worker/612-b3-mcpwirings-da2b58-3 into worker/fleetd-612-unita-87807e-1
CI / shell-tests (pull_request) Failing after 12s
CI / contract (pull_request) Successful in 1m33s
CI / build (pull_request) Failing after 2m25s
2026-09-22 07:36:05 +02:00
Dai Ha 6edeb70bc4 fleetd #612 B3 correction: distinguish lead vs member herdr in FleetdLeadRolloverAssemblyTest
Ticket comment 17553 on fleetd #612 found that the test's single shared
FakeHerdr made router.leadAgents() and router.memberAgents() collapse to
the identical client (FleetdAssembly.java:140-142's no-distinct-socket
fallback), so a mutation swapping leadAgents() for memberAgents() at the
FleetdAssembly.java:408 call site was invisible to this test even though
the two are genuinely different daemons in production.

Configure two distinct herdr sockets and two distinct FakeHerdr instances
(the same TwoHerdrResourcePorts shape B2's FleetdAssemblyConnectionIdentityTest
uses) and assert the roll's /clear + bootstrap sends land on the LEAD fake
and never on the MEMBER one.

Proven red against the router.memberAgents() mutation, reverted, touched,
and re-run green — both outputs recorded in the PR.
2026-09-22 12:31:39 +07:00
Dai Ha bc49d87cb8 Merge remote-tracking branch 'origin/worker/fleetd-612-unita-87807e-1' into worker/612-b3-mcpwirings-da2b58-3 2026-09-22 12:27:47 +07:00
Dai Ha 2b52324d9a fleetd #612 step 2 unit B3: behavioural replacements for lead-seat, quarantine
and lead-rollover source-text guards

Replaces three FleetdAssembly.java source-text guards (each scraped
Fleetd.java for a call site that fleetd #612 Unit A moved into
FleetdAssembly.java) with tests that drive the real assembled objects
through FleetdAssembly.assembleAndStart(...) -> FleetdRuntime.mcp(),
per the step-2 B-unit split (issue #612 comment 17513).

- Deleted FleetdLeadSeatWiringTest (fleetd #176): pinned that
  FleetMcp's LeadSeatSource construction still wires
  Fleetd.leadSeatLookup(...) by scraping the constructor call's text.
  Replaced by FleetdLeadSeatAssemblyTest, which seeds one FakeHerdr
  tab labelled to match a configured fleet.leaders.opus.tab and
  asserts the REAL assembled LeadSeatSource (via
  runtime.mcp().leadSeatSource()) reports the live lead's seat against
  its own subscription profile -- 1, not the 0 LeadSeatSource.none()
  (the inert stand-in) could ever report.

- Deleted FleetdBackendQuarantineWiringTest (fleetd #466): pinned that
  the escalating BackendQuarantine.withEscalation(...) text was
  present and the flat two-argument constructor's text was absent.
  Replaced by FleetdBackendQuarantineAssemblyTest, which quarantines
  the same credential twice through the REAL assembled
  BackendQuarantine (via runtime.mcp().quarantineSource().quarantine())
  at controlled fake-clock offsets and asserts the second cooldown
  doubles (200s vs 100s) -- the one behavioural difference escalation
  and the flat constructor actually produce.

- Deleted FleetdLeadRolloverWiringTest (fleetd #480), all three
  methods: unrelatedAnchorStillPresent was a scaffold anchor with no
  independent claim, needing no replacement.
  mainStillCallsTheLeadRolloverFactory pinned the leadRollover
  assignment's call-site text. factoryGatesOnConfigPresence pinned
  that an absent leadRollover: config yields no LeadRollover.
  Replaced by FleetdLeadRolloverAssemblyTest's two tests:
  assembledLeadRolloverRunsTheRealClearAndBootstrapSequence drives the
  REAL assembled LeadRollover (via runtime.mcp().leadRollover())
  through open()/confirm() end to end and asserts /clear then
  bootstrapText were actually sent through the real herdr router,
  reaching ROLLED. absentLeadRolloverConfigMeansNoRolloverIsBuilt
  calls Fleetd.leadRollover(...) directly with no leadRollover: block
  and asserts null -- this claim was found uncovered elsewhere
  (LeadRolloverTest's only related assertion is vacuous, assertNull
  (null), and never calls the real factory).

Each of the three FleetdAssembly.java call sites (quarantine
line 179-180, leadRollover line 408, lead seats line 479) was mutated
to its named inert variant, run against ONLY its new test (RED),
reverted, touch'd (Maven mtime trap) and re-run (GREEN) -- six proven
runs, pasted in the PR body.

FleetMcp.java: adds three accessors (quarantineSource(),
leadSeatSource(), leadRollover()) alongside the existing
registeredTools() -- but public, not package-private, and this is a
deliberate deviation from that precedent, not an oversight: these new
assembly tests cannot live in package dev.ltms.fleet.mcp the way
registeredTools()'s callers do, because they also build the
ResourcePorts FleetdAssembly.assembleAndStart(...) needs, and
ResourcePorts' methods return Fleetd-nested types visible only from
package dev.ltms.fleet. Package-private would compile but be
unreachable from there.

Full mvn -o test in fleetd/: Tests run: 1879, Failures: 6 (down from
the branch baseline's 1880/9 by exactly the 3 guards this unit
deletes) -- the remaining 6 are FleetdCompletionResolverWiringTest (4)
and FleetdConnectionIdentityConstructionTest /
FleetdFleetAppConstructionTest (1 each), all out of this unit's scope
(B1/B2).
2026-09-22 12:18:48 +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
14 changed files with 945 additions and 279 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
@@ -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;
@@ -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) {
@@ -797,6 +797,46 @@ public final class FleetMcp {
return server.listTools();
}
/**
* fleetd #612 B3 — same reason as {@link #registeredTools()}: a test that must drive the REAL
* {@link QuarantineSource} (and the real {@link BackendQuarantine} it wraps) this daemon was
* assembled with, rather than scraping {@code FleetdAssembly.java}'s source text for the
* constructor call that built it. Unlike {@link #registeredTools()}'s callers, that test cannot
* live in this package: it also builds the {@code ResourcePorts} that drives
* {@code FleetdAssembly.assembleAndStart}, and {@code ResourcePorts}' methods return
* {@code Fleetd}-nested types that are only visible from package {@code dev.ltms.fleet} — so
* this accessor is {@code public}, not package-private, to stay reachable from there. {@code
* FleetdBackendQuarantineAssemblyTest} quarantines a credential twice through this exact
* instance and checks the second cooldown is longer than the first — the one behavioural
* difference {@link BackendQuarantine#withEscalation} and the flat two-argument constructor
* actually produce.
*/
public QuarantineSource quarantineSource() {
return quarantine;
}
/**
* fleetd #612 B3 — as {@link #quarantineSource()}, {@code public} for the same cross-package
* reason, for the real {@link LeadSeatSource} this daemon was assembled with. {@code
* FleetdLeadSeatAssemblyTest} calls {@code seatsFor} on this exact instance and checks it
* reports a live lead's seat, which {@link LeadSeatSource#none()} can never do (it is a
* constant-zero function regardless of input).
*/
public LeadSeatSource leadSeatSource() {
return leadSeats;
}
/**
* fleetd #612 B3 — as {@link #quarantineSource()}, {@code public} for the same cross-package
* reason, for the real {@link LeadRollover} (or {@code null}) this daemon was assembled with.
* {@code FleetdLeadRolloverAssemblyTest} drives {@code open}/{@code confirm} on this exact
* instance and waits for the real continuation to send {@code /clear} and {@code bootstrapText}
* through the real {@code router.leadAgents()}.
*/
public LeadRollover leadRollover() {
return leadRollover;
}
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
/**
@@ -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,199 @@
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.msg.ReplyInbox;
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.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.OptionalLong;
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 B3 — replaces {@code FleetdBackendQuarantineWiringTest} (fleetd #466), a source-text
* test that scraped {@code Fleetd.java} (now {@code FleetdAssembly.java}, moved there by fleetd #612
* Unit A) for the {@code BackendQuarantine.withEscalation(...)} call, and separately asserted the
* flat two-argument constructor's text was ABSENT. That proves the right method NAME appears in
* source; it proves nothing about what the constructed object actually DOES.
*
* <p>This test instead drives the REAL {@link BackendQuarantine} the real {@link
* FleetdAssembly#assembleAndStart} builds — reached through {@link
* dev.ltms.fleet.mcp.FleetMcp#quarantineSource()} on the real, live {@code FleetMcp} {@code
* FleetdRuntime} owns — and asserts the ONE behavioural difference {@code withEscalation} and the
* flat constructor actually produce (see {@link BackendQuarantine}'s own class doc, "Mechanism"):
* quarantining the same credential twice in a row, within one base cooldown of the first deadline,
* must escalate the second cooldown past the first. A flat instance reports the identical cooldown
* both times.
*/
class FleetdBackendQuarantineAssemblyTest {
/** Base cooldown used throughout — long enough that rounding never blurs the 2x escalation. */
private static final int COOLDOWN_SECONDS = 100;
private static final class RecordingResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
final AtomicLong clockNanos = new AtomicLong(0L);
@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) -> replyInbox;
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException(
"leadMailboxOpener must not be called — no coordinator: block is configured");
};
}
@Override
public LongSupplier nanoClock() {
// Controllable: the SAME LongSupplier instance BackendQuarantine.withEscalation(...) is
// built with, so advancing clockNanos after assembly moves the quarantine tracker's own
// clock, with no real sleep needed to observe escalation.
return clockNanos::get;
}
@Override
public LongSupplier wallClockNanos() {
return clockNanos::get;
}
@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 final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
@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;
}
@Override
public void close() {
}
}
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
broker:
uri: "amqp://fake-test-broker/vh"
quarantineCooldownSeconds: %d
""".formatted(COOLDOWN_SECONDS));
return FleetConfig.load(f);
}
@Test
@DisplayName("[BEHAVIOURAL] the real assembled BackendQuarantine escalates a repeated exhaustion, "
+ "which the flat two-argument constructor can never do")
void assembledQuarantineEscalatesOnARepeatedExhaustion(@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);
BackendQuarantine quarantine = runtime.mcp().quarantineSource().quarantine();
// First exhaustion, at clock=0: a fresh occurrence, blocked for exactly the base cooldown.
quarantine.quarantine("cred-x");
BackendQuarantine.Status first = quarantine.status("cred-x").orElseThrow(
() -> new AssertionError("credential must be quarantined immediately after quarantine()"));
assertEquals(1, first.repeatCount(), "the first call is repeat #1");
assertEquals(COOLDOWN_SECONDS, first.remainingSeconds(),
"a fresh quarantine blocks for exactly the base cooldown");
// Second exhaustion, arriving just after the first deadline — well within one base cooldown
// of it, so this is a CONTINUATION of the same streak (repeat #2), not a fresh occurrence.
long firstDeadlineNanos = COOLDOWN_SECONDS * 1_000_000_000L;
ports.clockNanos.set(firstDeadlineNanos + 1);
quarantine.quarantine("cred-x");
BackendQuarantine.Status second = quarantine.status("cred-x").orElseThrow(
() -> new AssertionError("credential must be quarantined immediately after the second "
+ "quarantine() call"));
assertEquals(2, second.repeatCount(), "the second call, arriving within one base cooldown of "
+ "the first deadline, continues the streak as repeat #2");
// The one behavioural difference: withEscalation doubles the cooldown on repeat #2 (capped
// well above this at 12x base), the flat two-argument constructor never grows past the base
// cooldown no matter how many times quarantine() is called in a row.
assertEquals(2 * COOLDOWN_SECONDS, second.remainingSeconds(),
"withEscalation's default backoff doubles the cooldown on the second consecutive "
+ "exhaustion — this is the exact call FleetdAssembly.java makes at the "
+ "BackendQuarantine.withEscalation(...) call site");
assertTrue(second.remainingSeconds() > first.remainingSeconds(),
"the flat two-argument BackendQuarantine constructor would report the SAME remaining "
+ "seconds both times — this inequality is what a mutation to the flat "
+ "constructor at that call site must fail");
// Also confirm isQuarantined/remainingSeconds agree, exercising the accessors a real caller
// (fleet_profiles / fleet_list, per BackendQuarantine's own class doc) actually reads.
assertTrue(quarantine.isQuarantined("cred-x"));
OptionalLong remaining = quarantine.remainingSeconds("cred-x");
assertTrue(remaining.isPresent());
assertEquals(2 * COOLDOWN_SECONDS, remaining.getAsLong());
}
}
@@ -1,65 +0,0 @@
package dev.ltms.fleet;
import java.nio.file.Files;
import java.nio.file.Path;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #466 follow-up: {@code Fleetd.main} builds the daemon's one {@code BackendQuarantine}
* from {@link dev.ltms.fleet.placement.BackendQuarantine#withEscalation(java.util.function.LongSupplier,
* long)} — the escalating factory — rather than the plain two-argument constructor, which is still a
* flat cooldown (kept for backward compatibility, see that class's doc). {@code
* BackendQuarantineTest} proves {@code withEscalation} itself escalates, is ceilinged, and resets;
* it says nothing about which one {@code main} actually calls.
*
* <p>Measured directly: reverting {@code main} to {@code new BackendQuarantine(System::nanoTime,
* TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()))} — the pre-#466 flat call — compiles
* with 0 errors and leaves the entire 1608-test suite (including every {@code BackendQuarantineTest}
* case) green, because no other test constructs its {@code BackendQuarantine} through {@code main};
* every one of them builds its own instance directly. That silent regression is exactly the shape
* {@link FleetdLeadSeatWiringTest} and {@link FleetdCompletionResolverWiringTest} already guard
* against for their own constructor arguments — this is the same class of gap for fleetd #466's
* factory choice, following their approach.
*
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
* BackendQuarantine} and never runs {@code main} — a green result here proves only that the exact
* text {@code main} calls {@code BackendQuarantine.withEscalation(...)} rather than the flat
* constructor. It does not prove that call actually executes at startup (no test here starts the
* daemon), and it does not prove the escalation reaches a real backend or credential — only
* {@code BackendQuarantineTest} proves the factory's own behaviour, and only a live daemon proves
* the wiring runs.
*/
class FleetdBackendQuarantineWiringTest {
private static String fleetdSource() throws Exception {
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
}
@Test
@DisplayName("[SOURCE TEXT] main's BackendQuarantine local is still built from BackendQuarantine.withEscalation(...)")
void mainStillWiresTheEscalatingQuarantineFactory() throws Exception {
String source = fleetdSource();
assertTrue(source.contains(
"BackendQuarantine quarantine = BackendQuarantine.withEscalation(System::nanoTime,\n"
+ " TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));"),
"Fleetd.main's BackendQuarantine local must still be built from "
+ "BackendQuarantine.withEscalation(System::nanoTime, "
+ "TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds())). Reverting to the flat "
+ "two-argument constructor (fleetd #466's measured regression) compiles with 0 errors "
+ "and leaves the whole suite green, including every BackendQuarantineTest case that "
+ "proves the escalation itself works — this source check is what must go red instead. "
+ "A reverted daemon would go back to retrying a weekly subscription limit on every "
+ "flat ~30-minute cooldown, about 336 times across the week.");
// Negative form of the same check: the pre-#466 flat call, if it ever reappears at this
// declaration, must not be mistaken for the escalating one by a looser positive-only check.
assertFalse(source.contains(
"BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,\n"
+ " TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));"),
"main's BackendQuarantine local must never regress to the flat two-argument constructor");
}
}
@@ -0,0 +1,285 @@
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.AgentControl;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.lead.LeadRollover;
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.nio.file.Files;
import java.nio.file.Path;
import java.util.LinkedHashMap;
import java.util.List;
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.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assertions.fail;
/**
* fleetd #612 B3 — replaces {@code FleetdLeadRolloverWiringTest} (fleetd #480). That class was a
* source-text test scraping {@code Fleetd.java} (now {@code FleetdAssembly.java}, moved there by
* fleetd #612 Unit A) with three methods: {@code unrelatedAnchorStillPresent} (a scaffold anchor,
* not an independent claim — needs no replacement of its own), {@code
* mainStillCallsTheLeadRolloverFactory} (the call-site pin replaced by {@link
* #assembledLeadRolloverRunsTheRealClearAndBootstrapSequence}), and {@code
* factoryGatesOnConfigPresence} (the absent-config claim replaced by {@link
* #absentLeadRolloverConfigMeansNoRolloverIsBuilt} — a claim this ticket found was NOT actually
* covered behaviourally anywhere else: {@code LeadRolloverTest}'s only related assertion is
* vacuous, {@code assertNull(null)}, and never calls the real factory).
*
* <p><strong>fleetd #612 B3 correction (ticket comment 17553):</strong> the first version of this
* test configured a single shared {@link FakeHerdr} for both the lead and member herdr sockets.
* {@code FleetdAssembly.java:140-142} falls back to {@code memberHerdr = herdr} whenever no
* distinct {@code memberHerdrSocket} is configured, so with one fake, {@code
* router.leadAgents()} and {@code router.memberAgents()} wrapped the identical client — a
* mutation swapping {@code Fleetd.leadRollover(cfg, router.leadAgents(), config, leads)} for
* {@code ..., router.memberAgents(), ...} at {@code FleetdAssembly.java:408} was therefore
* invisible to this test, even though the two are genuinely different daemons in production. This
* version configures two distinct sockets and two distinct {@link FakeHerdr} instances (the same
* pattern {@code FleetdAssemblyConnectionIdentityTest}, fleetd #612 B2, already uses to separate
* lead from member) and asserts the roll's {@code /clear}/bootstrap sends land on the LEAD fake
* and never on the MEMBER one.
*/
class FleetdLeadRolloverAssemblyTest {
private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock");
private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock");
/** Keys {@code connectHerdr} by socket path so the lead and member daemons can be two
* DIFFERENT {@link FakeHerdr}s — same shape as B2's {@code FleetdAssemblyConnectionIdentityTest
* .TwoHerdrResourcePorts}. */
private static final class RecordingResourcePorts implements ResourcePorts {
final Map<Path, HerdrClient> herdrsBySocket = new LinkedHashMap<>();
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
HerdrClient client = herdrsBySocket.get(socketPath);
if (client == null) {
throw new IllegalStateException("no fake herdr registered for socket " + socketPath);
}
return client;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> replyInbox;
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException(
"leadMailboxOpener must not be called — no coordinator: block is configured");
};
}
@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) {
}
@Override
public void startHttp(Javalin app, String host, int port) {
}
}
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
@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;
}
@Override
public void close() {
}
}
private static FleetConfig writeConfig(Path dir, Path leadCwd) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
herdrSocket: "%s"
memberHerdrSocket: "%s"
idleSleepGuard:
enabled: false
broker:
uri: "amqp://fake-test-broker/vh"
fleet:
leaders:
opus:
tab: "lead: opus"
cwd: "%s"
leadRollover:
handoverPath: handover.md
requireOperatorConfirm: false
""".formatted(LEAD_SOCKET, MEMBER_SOCKET, leadCwd.toString()));
return FleetConfig.load(f);
}
@SuppressWarnings("unchecked")
@Test
@DisplayName("[BEHAVIOURAL] the real assembled LeadRollover runs the full open/confirm/continuation "
+ "sequence — /clear, then bootstrapText — through the real herdr router")
void assembledLeadRolloverRunsTheRealClearAndBootstrapSequence(@TempDir Path dir) throws Exception {
Path leadCwd = dir.resolve("lead-workspace");
Files.createDirectories(leadCwd);
FleetConfig cfg = writeConfig(dir, leadCwd);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
RecordingResourcePorts ports = new RecordingResourcePorts();
// Two DISTINCT fakes — one per configured socket — so leadAgents()/memberAgents() wrap
// genuinely different clients, exactly like production when memberHerdrSocket is set.
FakeHerdr lead = new FakeHerdr();
lead.withTab("w2", "w2:t7", "lead: opus");
FakeHerdr member = new FakeHerdr();
ports.herdrsBySocket.put(LEAD_SOCKET, lead);
ports.herdrsBySocket.put(MEMBER_SOCKET, member);
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
LeadRollover rollover = runtime.mcp().leadRollover();
assertNotNull(rollover, "leadRollover: is present in this test's config, so "
+ "FleetdAssembly.assembleAndStart must have built a real LeadRollover through the "
+ "Fleetd.leadRollover(...) call site — a mutation to `LeadRollover leadRollover = "
+ "null;` at that call site can never pass this");
LeadRollover.PendingRollover pending = rollover.open("term_a", "fleetd #612 B3 test");
String expectedHandoverPath = leadCwd.resolve("handover.md").normalize().toString();
assertEquals(expectedHandoverPath, pending.handoverPath());
// Ensure the handover file's mtime lands strictly AFTER open()'s requestedAtMillis —
// LeadRollover.checkHandover refuses on mtime <= requestedAt (HANDOVER_STALE).
Thread.sleep(50);
Files.writeString(Path.of(pending.handoverPath()), "handover content for fleetd #612 B3");
LeadRollover.RollDecision decision = rollover.confirm("term_a", pending.token(), true);
assertTrue(decision.accepted(), "confirm() must approve: requireOperatorConfirm is false, "
+ "the caller terminal matches open()'s, and the handover file exists, is non-empty "
+ "and fresh — got: " + decision);
// The production LeadRollover constructor always runs the post-confirm continuation on a
// real virtual thread (see Fleetd.leadRollover, which never passes the package-private test
// constructor), so this polls the real FleetMcp.leadRollover() instance's status(token)
// until the real continuation finishes.
LeadRollover.RollStatus status = pollUntilTerminal(rollover, pending.token());
assertEquals(LeadRollover.RollState.ROLLED, status.state(),
"the full happy path must complete: FakeHerdr's default agent status is 'idle', so "
+ "the turn-boundary wait settles immediately and the post-/clear wait "
+ "releases via its pickup-grace path — detail: " + status.detail());
// Prove the real herdr router actually sent BOTH messages, in order, to the real LEAD
// pane — this is the one thing a source-text pin on the call site could never show.
List<FakeHerdr.Call> prompts = lead.calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.toList();
assertTrue(prompts.size() >= 2, "expected at least a /clear send and a bootstrapText send "
+ "on the LEAD daemon, got " + prompts.size() + " agent.prompt calls: " + prompts);
assertEquals("/clear", ((Map<String, Object>) prompts.get(0).params()).get("text"),
"the first send must be the literal /clear housekeeping command");
Object secondText = ((Map<String, Object>) prompts.get(1).params()).get("text");
assertTrue(secondText instanceof String && ((String) secondText).contains(expectedHandoverPath),
"the second send must be the default bootstrapText naming the resolved handover "
+ "path, got: " + secondText);
// fleetd #612 B3 correction: prove the roll never touches the MEMBER daemon. A mutation
// swapping router.leadAgents() for router.memberAgents() at the real call site would move
// both sends above onto `member` instead, which this assertion catches — the thing the
// single-fake version of this test could never see, because both wrapped the same client.
List<FakeHerdr.Call> memberPrompts = member.calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.toList();
assertTrue(memberPrompts.isEmpty(), "the roll must be wired to the LEAD daemon only — got "
+ memberPrompts.size() + " agent.prompt call(s) on the MEMBER daemon instead: "
+ memberPrompts);
}
private static LeadRollover.RollStatus pollUntilTerminal(LeadRollover rollover, String token)
throws InterruptedException {
long deadline = System.nanoTime() + java.util.concurrent.TimeUnit.SECONDS.toNanos(10);
while (System.nanoTime() < deadline) {
LeadRollover.RollStatus status = rollover.status(token);
if (status.state() != LeadRollover.RollState.PENDING
&& status.state() != LeadRollover.RollState.IN_PROGRESS) {
return status;
}
Thread.sleep(50);
}
fail("the real continuation did not reach a terminal state within 10s — last status: "
+ rollover.status(token));
throw new AssertionError("unreachable");
}
@Test
@DisplayName("[BEHAVIOURAL] Fleetd.leadRollover(...) returns null when leadRollover: is absent "
+ "from config — the opt-in gate FleetdLeadRolloverWiringTest's "
+ "factoryGatesOnConfigPresence pinned by source text alone")
void absentLeadRolloverConfigMeansNoRolloverIsBuilt(@TempDir Path dir) throws Exception {
Path yaml = dir.resolve("fleetd.yaml");
Files.writeString(yaml, """
bind:
host: 127.0.0.1
port: 8765
""");
ConfigRef config = new ConfigRef(yaml, FleetConfig.load(yaml));
AgentControl agents = new AgentControl(new FakeHerdr());
LeadRollover rollover = Fleetd.leadRollover(config.get(), agents, config, Map::of);
assertNull(rollover, "leadRollover: is absent from this config, so the factory's opt-in "
+ "gate (`if (cfg.leadRollover() == null) return null;`) must fire and no "
+ "LeadRollover must be constructed at all");
}
}
@@ -1,89 +0,0 @@
package dev.ltms.fleet;
import java.nio.file.Files;
import java.nio.file.Path;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #480 Unit A, hard requirement 6: pin {@code Fleetd.main}'s construction of {@link
* dev.ltms.fleet.lead.LeadRollover} with a source-text assertion, mirroring {@code
* FleetdCompletionResolverWiringTest}'s pattern — five log-only reporters in {@code Fleetd.main}
* already survived mutation batteries this exact way (fleetd #415's extraction antidote note).
*
* <p>What this class still covers, and what it never claimed to. {@code LeadRolloverTest}
* constructs its own {@code LeadRollover} directly (as every prior test of an extracted factory
* does) with a hand-built lookup, so a mutation that deletes the {@code leadRollover(...)} call
* from {@code main} — or replaces one of its arguments with something that still compiles, e.g.
* {@code router.leadAgents()} swapped for {@code null}, or the whole assignment swapped for a bare
* {@code null} literal — leaves every behavioural test green. This is a plain string read, guarded
* by an unrelated anchor assertion so a broken or empty file read cannot pass as a real change.
*
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
* LeadRollover} and never runs {@code main}. It pins the {@code leadRollover(...)} CALL SITE's
* argument list — that {@code main} still passes {@code leads} at all — never what the factory
* DOES with that argument once inside its own body.
*
* <p><b>Correction (fleetd #480 relative-handover-path follow-up): that gap used to be real, and
* now is not — but not here.</b> This class's javadoc previously claimed "no behavioural test can
* catch this wiring dropping out" for the whole factory, including the lambda {@code
* leadRollover(...)} builds internally (terminal → lead name → {@code Leader.cwd()}). That claim
* was proven true at the time — mutating that lambda's body to {@code String leadName = null;}
* (always "no lead found", which silently reintroduces the daemon-cwd bug this ticket fixes) left
* the full suite green, {@code Tests run: 1669, Failures: 0}. It is no longer true: {@code
* FleetdLeadRolloverWorkspaceLookupTest} now calls {@code Fleetd.leadRollover(...)} directly with a
* real {@link dev.ltms.fleet.config.ConfigRef} built from a temp {@code fleetd.yaml}, and fails
* against that exact one-line mutation. So: THIS class still covers only the call site's argument
* list; {@code FleetdLeadRolloverWorkspaceLookupTest} is what now covers the lambda's body. Neither
* one subsumes the other — keep both.
*/
class FleetdLeadRolloverWiringTest {
private static String fleetdSource() throws Exception {
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
}
@Test
@DisplayName("[SOURCE TEXT] unrelated anchor: Fleetd.java still declares the Fleetd class")
void unrelatedAnchorStillPresent() throws Exception {
// Guards the two assertions below: without this, a bad read (empty string, wrong file,
// truncated file) could vacuously fail to contain the leadRollover(...) call too, and a
// test that only asserts "contains X" would report a false pass for the wrong reason if X
// happened to match. Asserting an unrelated, structurally distant string first proves the
// read actually pulled real file content.
String source = fleetdSource();
assertTrue(source.contains("public final class Fleetd"),
"sanity anchor failed — the file read did not return real Fleetd.java source; the "
+ "leadRollover(...) wiring assertions below cannot be trusted until this passes");
}
@Test
@DisplayName("[SOURCE TEXT] main still constructs LeadRollover via the leadRollover(...) factory, exactly as heartbeat is constructed")
void mainStillCallsTheLeadRolloverFactory() throws Exception {
String source = fleetdSource();
assertTrue(source.contains(
"LeadRollover leadRollover = leadRollover(cfg, router.leadAgents(), config, leads);"),
"Fleetd.main must still assign `LeadRollover leadRollover = leadRollover(cfg, "
+ "router.leadAgents(), config, leads);`. Dropping this call, or swapping one of "
+ "its arguments for something that still compiles (e.g. null in place of "
+ "router.leadAgents()), leaves every behavioural test green — this source check is "
+ "what must go red instead. fleetd #480 correction 2 deliberately dropped "
+ "primaryRegistry from this call — see LeadRollover's class javadoc for why a "
+ "single-slot lookup was wrong here. The fleetd #480 relative-handover-path "
+ "follow-up added `leads` (terminal → lead name) so the factory can resolve a "
+ "relative handoverPath against the calling lead's own workspace.");
}
@Test
@DisplayName("[SOURCE TEXT] the leadRollover(...) factory itself gates construction on cfg.leadRollover() != null")
void factoryGatesOnConfigPresence() throws Exception {
String source = fleetdSource();
assertTrue(source.contains("if (cfg.leadRollover() == null) {"),
"Fleetd.leadRollover(...) must refuse to construct a LeadRollover when the "
+ "leadRollover: block is absent — an upgraded daemon must never silently acquire "
+ "the ability to clear the lead's own pane. See LeadHeartbeatLoop's construction "
+ "gate (cfg.leadHeartbeat() != null) for the pattern this mirrors.");
}
}
@@ -0,0 +1,172 @@
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.FleetMcp;
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.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.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
/**
* fleetd #612 B3 — replaces {@code FleetdLeadSeatWiringTest} (fleetd #176), a source-text test that
* scraped {@code Fleetd.java} (now {@code FleetdAssembly.java}, moved there by fleetd #612 Unit A)
* for the exact {@code new FleetMcp.LeadSeatSource(Fleetd.leadSeatLookup(...))} constructor-call
* text. That proves the right symbols appear in source; it proves nothing about what the daemon's
* live {@code fleet_list} actually reports.
*
* <p>This test instead drives the REAL {@link FleetMcp.LeadSeatSource} the real {@link
* FleetdAssembly#assembleAndStart} builds — including the REAL {@code LeadTabScanner} it wires
* {@code Fleetd.leadSeatLookup} through — reached via {@link FleetMcp#leadSeatSource()} on the
* live {@code FleetMcp} {@code FleetdRuntime} owns. It seeds one FakeHerdr tab labelled to match a
* configured {@code fleet.leaders.opus.tab}, with a live agent already in it (FakeHerdr's own
* default {@code agent.list}/{@code pane.list} entries for {@code term_a}/{@code w2:p7}/{@code
* w2:t7} — no FakeHerdr change needed), and asserts the assembled seat source reports exactly the
* seat {@link FleetMcp.LeadSeatSource#none()} (the inert stand-in) could never produce: 1, not 0.
*/
class FleetdLeadSeatAssemblyTest {
private static final class RecordingResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
@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) -> replyInbox;
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException(
"leadMailboxOpener must not be called — no coordinator: block is configured");
};
}
@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) {
}
@Override
public void startHttp(Javalin app, String host, int port) {
}
}
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
@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;
}
@Override
public void close() {
}
}
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
broker:
uri: "amqp://fake-test-broker/vh"
fleet:
leaders:
opus:
tab: "lead: opus"
profile: sonnet
profiles:
sonnet:
subscription: true
argv: ["ccs", "sonnet"]
""");
return FleetConfig.load(f);
}
@Test
@DisplayName("[BEHAVIOURAL] the real assembled LeadSeatSource, backed by the real LeadTabScanner, "
+ "reports a live lead's seat against its own subscription profile")
void assembledLeadSeatSourceReportsALiveLeadsSeat(@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();
// Label FakeHerdr's own default pane's tab (term_a / w2:p7 / w2:t7, already carrying a live
// agent) to match fleet.leaders.opus.tab exactly — no FakeHerdr change needed at all.
ports.herdr.withTab("w2", "w2:t7", "lead: opus");
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
FleetMcp.LeadSeatSource seatSource = runtime.mcp().leadSeatSource();
assertEquals(1, seatSource.seatsFor().apply("sonnet"),
"the real LeadTabScanner recognises the labelled tab as a live 'opus' lead on "
+ "profile 'sonnet' (same credential, matched by Fleetd.leadSeatLookup), so "
+ "subscription profile 'sonnet' must be charged one seat — "
+ "FleetMcp.LeadSeatSource.none() (the inert stand-in this test's mutation "
+ "swaps the call site for) always reports 0, whatever the input");
// A profile no lead is running on gets no seat charged — the same seat source, applied to
// an input that must stay at the inert answer even on the real, non-inert instance.
assertEquals(0, seatSource.seatsFor().apply("no-such-profile"));
}
}
@@ -1,43 +0,0 @@
package dev.ltms.fleet;
import java.nio.file.Files;
import java.nio.file.Path;
import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #176: {@code Fleetd.main} builds its {@code FleetMcp} from a 14-argument constructor whose
* last argument is a {@code FleetMcp.LeadSeatSource} wrapping {@link Fleetd#leadSeatLookup}. That
* argument is exactly the kind of wiring fleetd #248 warned about: dropping it (or swapping it for
* the inert {@code FleetMcp.LeadSeatSource.none()}) compiles with 0 errors and leaves every test
* that builds its own {@code FleetMcp}/{@code CapacitySource} directly — every test that predates
* this ticket — green, because none of them go through {@code main} at all.
*
* <p>{@link FleetdLeadSeatLookupTest} proves the factory's own matching logic; this class is the
* plain source-text assertion that proves {@code main} still passes its result in, mirroring
* {@code FleetdCompletionResolverWiringTest}'s approach for the same class of gap.
*
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a
* {@code FleetMcp} and never runs {@code main}.
*/
class FleetdLeadSeatWiringTest {
private static String fleetdSource() throws Exception {
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
}
@Test
@DisplayName("[SOURCE TEXT] FleetMcp's construction call still passes a LeadSeatSource built from leadSeatLookup(...)")
void fleetMcpConstructionStillWiresLeadSeatLookup() throws Exception {
String source = fleetdSource();
assertTrue(source.contains("new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), "
+ "leaders, leads))"),
"FleetMcp's construction call must still pass a LeadSeatSource built from "
+ "Fleetd.leadSeatLookup(...). Dropping it or swapping in "
+ "FleetMcp.LeadSeatSource.none() (fleetd #176's would-be silent regression, the same "
+ "shape as fleetd #248's measured mutations) compiles with 0 errors and leaves every "
+ "existing behavioural test green — this source check is what must go red instead.");
}
}
@@ -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),