Compare commits

..

1 Commits

Author SHA1 Message Date
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
5 changed files with 148 additions and 148 deletions
+11 -51
View File
@@ -167,59 +167,19 @@ 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.
- **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.
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.
- **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 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.
- **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.
## Writing style
@@ -577,13 +577,18 @@ public final class Fleetd {
// (keyed by configDir+sessionId), so this costs at most one extra bounded tail read per
// TTL window, never a shared-mutable-state hazard between the two callers.
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.
boolean requireOperatorConfirm = cfg.leadRollover() == null || cfg.leadRollover().requireOperatorConfirm();
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
pushLoop, heartbeatScheduler, System::nanoTime,
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
metrics,
leadContextSource(leadContextGauge, router.leadAgents(), leads,
leadConfigDirLookup(() -> config.get().profiles(), leaders)),
Boolean.TRUE.equals(hb.contextHighNudge()));
Boolean.TRUE.equals(hb.contextHighNudge()), requireOperatorConfirm);
heartbeat.start();
} else {
heartbeat = null;
@@ -76,6 +76,7 @@ public final class LeadHeartbeatLoop {
private final Metrics metrics; // CB-512 pattern: nullable — no registry in unit tests
private final LeadContextSource contextSource; // fleetd #609
private final boolean contextHighNudge; // fleetd #609: opt-in, like the loop itself
private final boolean requireOperatorConfirm; // fleetd #621: mirrors leadRollover.requireOperatorConfirm
/** When the current idle stretch began (nanos), or {@link #NOT_IDLE}. Single scheduler thread only. */
private long idleSinceNanos = NOT_IDLE;
@@ -106,12 +107,32 @@ public final class LeadHeartbeatLoop {
* fleetd #609: as above, plus the lead's own context source and whether a HIGH reading should
* append a hand-over notice to the loop's nudge. Pass {@link LeadContextSource#none()} and
* {@code false} to keep the pre-#609 behaviour exactly (both existing public constructors do).
*
* <p>fleetd #621: delegates to the full constructor with {@code requireOperatorConfirm=true} —
* the pre-#621 wording ("ask the operator ... only the operator can approve the roll") assumed
* the config default, so every caller of this overload keeps that text byte-identical.
*/
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
ScheduledExecutorService scheduler, LongSupplier clock,
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
LeadContextSource contextSource, boolean contextHighNudge) {
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock,
idleAfterNanos, backoffMs, quietNudgeCap, metrics, contextSource, contextHighNudge, true);
}
/**
* fleetd #621: as above, plus the daemon's effective {@code leadRollover.requireOperatorConfirm}
* value — threaded into {@link #contextNotice(boolean, LeadContextGauge.Reading, boolean, boolean)}
* so the notice's wording tracks the config the daemon actually enforces (see {@code
* LeadRollover.confirm}) instead of always asserting the operator gate is on.
*/
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
ScheduledExecutorService scheduler, LongSupplier clock,
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
LeadContextSource contextSource, boolean contextHighNudge,
boolean requireOperatorConfirm) {
this.primaryRegistry = primaryRegistry;
this.agents = agents;
this.inbox = inbox;
@@ -125,6 +146,7 @@ public final class LeadHeartbeatLoop {
this.metrics = metrics;
this.contextSource = contextSource;
this.contextHighNudge = contextHighNudge;
this.requireOperatorConfirm = requireOperatorConfirm;
}
/**
@@ -344,7 +366,7 @@ public final class LeadHeartbeatLoop {
// d.contextNotified() is the value to persist once delivery is confirmed, not the value the text
// itself should be built from. Otherwise a HIGH stretch that is still latched would never see the
// notice at all, defeating the very check this fixes.
String notice = contextNotice(contextHighNudge, reading, contextNotified);
String notice = contextNotice(contextHighNudge, reading, contextNotified, requireOperatorConfirm);
var lead = primaryRegistry.primaryTerminal();
boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
// The latch becomes true only when all three hold: decide() chose to notify, a notice was
@@ -381,8 +403,9 @@ public final class LeadHeartbeatLoop {
/**
* fleetd #609: the text appended to a nudge when the lead's own context is full — {@code ""}
* whenever the notice does not apply, so callers can unconditionally append this without an extra
* branch. Wording stays plain (CEFR B1) and honest that only the operator approves a roll — this
* loop only ever prints text, it never calls {@code fleet_handover} itself.
* branch. Wording stays plain (CEFR B1) and honest about who actually gates the roll — see the
* {@code requireOperatorConfirm} overload (fleetd #621) for which check that is. This loop only
* ever prints text, it never calls {@code fleet_handover} itself.
*
* @param enabled the {@code leadHeartbeat.contextHighNudge} config flag
* @param reading the lead's current {@link LeadContextGauge} reading
@@ -401,9 +424,33 @@ public final class LeadHeartbeatLoop {
* closing sentence ("You will not be told again until your context reads ok.") false. {@link
* #injectNudge} is the only caller that passes a non-default {@code alreadyNotified}.
*
* <p>fleetd #621: delegates with {@code requireOperatorConfirm=true} — the pre-#621 default and the
* value every existing caller of this overload (including every test written before #621) already
* assumed, so the text this overload returns stays byte-identical.
*
* @param alreadyNotified whether the lead has already been told about the current HIGH stretch
*/
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified) {
return contextNotice(enabled, reading, alreadyNotified, true);
}
/**
* fleetd #621: as {@link #contextNotice(boolean, LeadContextGauge.Reading, boolean)}, but the closing
* instructions also track the daemon's effective {@code leadRollover.requireOperatorConfirm} value,
* instead of always asserting that only the operator can approve the roll.
*
* <p>{@code LeadRollover.confirm(...)} already honours this flag: when it is {@code false}, the daemon
* itself gates the roll on the three handover-file checks alone (exists, modified after the {@code
* open()} request, and no older than {@code maxDocAgeSeconds}) and never consults {@code
* operatorConfirmed}. Before this parameter existed, this notice told the lead to ask the operator
* regardless — so a lead that followed its own instructions asked anyway, and setting the config knob
* to {@code false} stopped the daemon refusing the roll without stopping the operator being
* interrupted. This parameter is how the text is kept honest about which gate is actually live.
*
* @param requireOperatorConfirm the effective {@code leadRollover.requireOperatorConfirm} value
*/
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified,
boolean requireOperatorConfirm) {
if (!enabled || alreadyNotified || reading.state() != LeadContextGauge.State.HIGH) {
return "";
}
@@ -420,10 +467,18 @@ public final class LeadHeartbeatLoop {
sb.append(" (").append(reading.compactions()).append(' ').append(compactionWord)
.append(" so far).");
}
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
+ "again until your context reads ok.");
if (requireOperatorConfirm) {
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
+ "again until your context reads ok.");
} else {
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
+ "write the file it names, then call fleet_handover(action=\"confirm\", token). Decide for "
+ "yourself when to confirm: the roll goes through if the handover file exists, was "
+ "changed after you opened it, and is not older than maxDocAgeSeconds. You will not be "
+ "told again until your context reads ok.");
}
return sb.toString();
}
@@ -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,8 +1846,7 @@ 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,
ReplyPushLoop pushLoop)
private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler)
implements AutoCloseable {
@Override
public void close() {
@@ -1862,20 +1861,14 @@ 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, nowNanos);
return new ManualPushWiring(service, leadHerdr, scheduler, pushLoop);
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, System::nanoTime);
return new ManualPushWiring(service, leadHerdr, scheduler);
}
private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException {
@@ -1966,35 +1959,24 @@ class MessageServiceTest {
@Test
void severalAsyncTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception {
// 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");
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);
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);
}
String second = wiring.service().sendAsync(T, "second task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "second done"));
Thread.sleep(100);
assertEquals(1, wiring.scheduler().runDueTasks(),
"both terminal tickets must coalesce onto one scheduled tick");
awaitNudge(wiring.leadHerdr());
Thread.sleep(200); // settle — nothing more should arrive beyond the one coalesced nudge
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");
@@ -2073,8 +2055,7 @@ class MessageServiceTest {
@Test
void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception {
// A 1ms backoff is due immediately, but ManualScheduler only ticks when this test asks it to.
try (var wiring = wireWithManualScheduler(5, 1)) {
try (var wiring = wireWithPushLoop(5, 50)) {
String ticket = wiring.service().sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
@@ -2082,9 +2063,8 @@ 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());
assertEquals(1, wiring.scheduler().runDueTasks(), "the open question must have one scheduled tick");
awaitNudge(wiring.leadHerdr());
long callsBeforeAnswer = wiring.leadHerdr().calls.stream()
.filter(c -> c.method().equals("agent.prompt")).count();
@@ -2092,22 +2072,12 @@ 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);
// 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");
}
// 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);
boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream()
.filter(c -> c.method().equals("agent.prompt"))
.skip(callsBeforeAnswer)
@@ -2190,45 +2160,35 @@ 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.
// 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");
try (var wiring = wireWithPushLoop(1, 50, clock::get)) {
String stale = wiring.service().sendAsync(T, "first task");
awaitWaiting();
injectDelivery();
assertTrue(rendezvous.resolve(T, "stale result"));
// 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");
// 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");
// 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.
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);
// 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);
}
// 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),