diff --git a/.claude/skills/handover/SKILL.md b/.claude/skills/handover/SKILL.md index 1b63437..8422d0e 100644 --- a/.claude/skills/handover/SKILL.md +++ b/.claude/skills/handover/SKILL.md @@ -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 diff --git a/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java b/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java index 03ebcdc..c51c3b4 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java +++ b/fleetd/src/main/java/dev/ltms/fleet/FleetdAssembly.java @@ -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; diff --git a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java index 94f4916..62da9a8 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java +++ b/fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java @@ -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. * - *

Which of the two inputs Claude Code actually follows when they disagree is intentionally - * not 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}. + *

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. * *

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)); } diff --git a/fleetd/src/main/java/dev/ltms/fleet/launch/ClaudeCodeArguments.java b/fleetd/src/main/java/dev/ltms/fleet/launch/ClaudeCodeArguments.java index ce2e214..263c2e4 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/launch/ClaudeCodeArguments.java +++ b/fleetd/src/main/java/dev/ltms/fleet/launch/ClaudeCodeArguments.java @@ -15,13 +15,15 @@ public final class ClaudeCodeArguments { * Append the configured Claude Code auto-compaction window when the profile opts in. * *

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 withAutoCompactWindow(List argv, FleetConfig.Profile profile) { if (profile.autoCompactWindow() == null) { diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java index c40f31d..37c6caa 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/LeadHeartbeatLoop.java @@ -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). + * + *

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> 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> 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}. * + *

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. + * + *

{@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(); } diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java index 9f48af9..4cee349 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/LeadHeartbeatLoopTest.java @@ -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 diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index 99f8f80..fb5b7f5 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -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 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),