Compare commits

...

10 Commits

Author SHA1 Message Date
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
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
ltms 17127efb88 Merge #601: pass auto-compact window to leads; warn instead of refusing on a conflict (CB-617)
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m16s
CI / build (push) Failing after 2m33s
2026-09-22 05:22:56 +02:00
ltms 203f034528 Merge #617: write FAILED instead of leaving a dead roll stuck at IN_PROGRESS (fleetd #615)
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 1m14s
CI / build (push) Failing after 1m47s
2026-09-22 05:21:21 +02:00
ltms 9ee16f5b85 Merge #616: report role-fallback gaps at boot, name contextHighNudge (fleetd #613)
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m15s
CI / build (push) Failing after 1m43s
2026-09-22 05:17:52 +02:00
Dai Ha 388ef5a3c3 fleetd #615: write FAILED instead of leaving status(token) stuck at IN_PROGRESS
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 2m25s
LeadRollover.runRollover made two unwrapped agents.send calls. HerdrException
is unchecked, and the production continuationRunner is a bare virtual thread
with no uncaught-exception handler, so a throw from either send call killed
the continuation silently — confirm() had already written IN_PROGRESS into
outcomes before scheduling it, and nothing ever overwrote that entry with a
terminal state.

Wrap the whole continuation body in one try/catch(RuntimeException), matching
the local convention already used around agents.status in
waitUntilAtTurnBoundary. On a throw, write a new terminal RollState.FAILED
entry naming the exception, in the same diagnostic style as
TURN_NEVER_SETTLED and CLEAR_NEVER_SETTLED.

Two new tests make send() throw on the /clear call and on the bootstrap-text
call respectively, each asserting status(token) reports FAILED, not
IN_PROGRESS. Reverting only the production catch (keeping the tests) turns
both red; restoring it turns them green again.
2026-09-22 10:17:12 +07:00
Dai Ha 987ccef4c7 fleetd #613: log role-fallback gaps at boot, name contextHighNudge in the heartbeat line
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 1m17s
CI / build (pull_request) Failing after 1m48s
- reportRoleFallbackGaps(cfg), called right after cfg.validateAll() in Fleetd.main, logs every
  MemberRole with no fleet.<role>s: pool (naming the profile count and the resolved
  defaultProfileFor(role) first choice) and, separately, every role with no
  fleet.charters.<role>: entry. Log only — the deliberate 'unconstrained' fallback in
  FleetConfig#candidateProfiles / CompositePeerLauncher#poolFor is unchanged, and a config with
  profiles: and no fleet: block still starts and still spawns.
- LeadHeartbeatLoop#start()'s boot line now also names contextHighNudge (fleetd #609), alongside
  the three settings it already logged.
- RoleFallbackGapReportTest (new) and two new LeadHeartbeatLoopTest cases pin both lines' content
  via a ListAppender, raising the dev.ltms.fleet logger past logback-test.xml's WARN override for
  the INFO-level lines.
2026-09-22 10:12:09 +07:00
10 changed files with 671 additions and 77 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
@@ -53,6 +53,7 @@ import dev.ltms.fleet.rest.FleetApp;
import dev.ltms.fleet.session.GitWorktrees;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.session.SessionReaper;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
@@ -176,6 +177,12 @@ public final class Fleetd {
// and the Fleetd-startup tests actually pin — see FleetConfig#validateAll's javadoc for
// why a name-by-name list here would have the same defect it replaces.
cfg.validateAll();
// fleetd #613: validateAll() (validateMembers() inside it) only refuses a slot that names a
// bad role or profile — it says nothing about a role that has NO pool or NO charter at all,
// because both are legitimate ("unconstrained") states, not errors. Report them here, right
// after validation passes, so an operator sees the gap once per restart instead of finding
// it later in a roster row (see reportRoleFallbackGaps' javadoc for the measured cause).
reportRoleFallbackGaps(cfg);
// fleetd #469, follow-up to #464: validateAll() (and validateCharters() inside it) only
// checks that a charter's KEY is a role wire name and its text is non-blank — it never
// looks at what the text actually names. This is the separate check that does: it asks
@@ -2175,6 +2182,58 @@ public final class Fleetd {
}
}
/**
* fleetd #613: {@code FleetConfig.candidateProfiles(MemberRole)} (FleetConfig.java:1827) and
* {@code CompositePeerLauncher.poolFor} (CompositePeerLauncher.java:598-603) both fall back to
* <em>every</em> configured profile when a role has no {@code fleet.<role>s:} pool — a
* deliberate "unconstrained" behaviour, kept unchanged here, that lets a config with only
* {@code profiles:} and no {@code fleet:} block still spawn. That fallback is silent, and on
* the host that opened this ticket it widened an unqualified {@code hunter} spawn to all 8
* configured profiles and picked {@code local} as the resolved first choice — a profile every
* other pool on that same config gives weight 0 to. Report it once at boot instead, naming both
* how many profiles the gap opens onto and the exact first choice, since the first choice (not
* the pool size) is what actually surprised the operator.
*
* <p>A missing {@code fleet.charters.<role>:} entry is reported separately: the member still
* runs, but with only the launcher's own reply charter and no role contract. Unlike the pool
* gap this already has a per-spawn instrument ({@code HerdrPeerLauncher.logCharterReceipt},
* {@code SessionManager}'s {@code charterSource} roster field) — this boot line is the same
* information surfaced once, up front, rather than discovered per member later.
*
* <p>Never refuses to start over either gap — both are legitimate configurations, and this is a
* report, not a validation. Package-private so a test can call it directly and capture the log
* via a {@link ch.qos.logback.core.read.ListAppender}, the same pattern {@link
* #reportExhaustedPatternGap} and {@link #reportMemberCredentialsGap} already use.
*/
static void reportRoleFallbackGaps(FleetConfig cfg) {
List<String> poolGaps = new ArrayList<>();
List<String> charterGaps = new ArrayList<>();
int profileCount = cfg.profiles().size();
for (MemberRole role : MemberRole.values()) {
boolean hasPool = cfg.fleet() != null && !cfg.fleet().profilesFor(role).isEmpty();
if (!hasPool) {
String firstChoice = cfg.defaultProfileFor(role);
poolGaps.add(role.wireName() + " (may land on any of " + profileCount
+ " profile(s), first choice "
+ (firstChoice == null ? "none — no profiles configured" : "'" + firstChoice + "'")
+ ")");
}
String charter = cfg.fleet() == null ? null : cfg.fleet().charterFor(role);
if (charter == null || charter.isBlank()) {
charterGaps.add(role.wireName());
}
}
if (!poolGaps.isEmpty()) {
log.info("role fallback: no fleet.<role>s: pool for {} — an unqualified spawn of that "
+ "role falls back to every configured profile (deliberate; see "
+ "FleetConfig#candidateProfiles)", poolGaps);
}
if (!charterGaps.isEmpty()) {
log.info("role fallback: no fleet.charters: entry for {} — that role runs with only "
+ "the launcher's reply charter, no role contract", charterGaps);
}
}
/**
* fleetd #474: the one place both the startup call (right after {@code cfg.validateAll()} in
* {@link #main}) and the reload call (wired into {@code config}'s {@code extraValidation} above,
@@ -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) {
@@ -231,7 +231,9 @@ public final class LeadRollover {
* #status} could wrongly answer {@link #UNKNOWN} ("nothing was ever requested") for a roll
* that is, in fact, actively running. This is not sticky: the deferred continuation
* overwrites this same entry with a terminal state ({@link #ROLLED}, {@link
* #TURN_NEVER_SETTLED}, or {@link #CLEAR_NEVER_SETTLED}) once it finishes.
* #TURN_NEVER_SETTLED}, {@link #CLEAR_NEVER_SETTLED}, or {@link #FAILED}) once it finishes
* — including by throwing, which fleetd #615's catch in {@link #runRollover} now turns into
* {@link #FAILED} instead of leaving this entry stuck forever.
*/
IN_PROGRESS,
/**
@@ -253,6 +255,19 @@ public final class LeadRollover {
* within {@code clearSettleSeconds} — {@code bootstrapText} was never sent.
*/
CLEAR_NEVER_SETTLED,
/**
* fleetd #615: the deferred continuation threw a {@link RuntimeException} — most likely a
* {@link dev.ltms.fleet.herdr.HerdrException} out of one of the two unwrapped {@code
* agents.send} calls in {@link #runRollover} — and the continuation thread died with it.
* Before this state existed, that throw left {@link #outcomes} holding {@link #IN_PROGRESS}
* forever, because the production {@code continuationRunner} is a bare virtual thread with
* no uncaught-exception handler and nothing downstream of the throw ever ran to write a
* terminal outcome. {@code detail} names the exception, so a reader has something to act on
* — the same diagnostic style as {@link #TURN_NEVER_SETTLED} and {@link
* #CLEAR_NEVER_SETTLED}. The roll is dead at this point and does not retry itself; a stuck
* lead must {@link #open} a fresh request.
*/
FAILED,
/**
* {@code token} names nothing this instance currently knows about: never issued by {@link
* #open}, dropped by {@link #cancel}, or aged out of {@link #outcomes}'s bounded history.
@@ -493,8 +508,42 @@ public final class LeadRollover {
* entirely after {@link #confirm} has returned to its caller — see this class's javadoc for the
* four-step order. There is no result to return to by this point, so every outcome is logged
* only.
*
* <p><strong>fleetd #615 — the whole body is wrapped in one {@code try}.</strong> The two {@code
* agents.send} calls below are not wrapped individually: {@code send} → {@code agentCall} →
* {@code herdr.call} can throw an unchecked {@link dev.ltms.fleet.herdr.HerdrException} (see
* {@code AgentControl.java}), and the production {@code continuationRunner} is a bare virtual
* thread with no uncaught-exception handler (see this class's public constructor). Before this
* fix, either throw killed the continuation thread silently, leaving the {@link
* RollState#IN_PROGRESS} entry {@link #confirm} wrote at hand-off stuck forever — {@link
* #status} had no way to tell a dead roll from one still genuinely running. The {@code catch}
* below is scoped to the method body rather than to each {@code send} call individually, so it
* also covers anything else added to this continuation later, not just today's two call sites —
* the same reasoning that put the write-a-terminal-outcome step at each of this method's other
* exits (see the {@link RollState#TURN_NEVER_SETTLED} and {@link RollState#CLEAR_NEVER_SETTLED}
* branches below) rather than inside the helpers that detect them.</p>
*
* <p>Only {@link RuntimeException} is caught, matching the local convention {@link
* #waitUntilAtTurnBoundary} already set around its own {@code agents.status} call — not the
* broader {@link Exception} or {@link Throwable}, which would also swallow something like an
* {@link OutOfMemoryError} this continuation has no business handling.</p>
*/
private void runRollover(PendingRollover p, FleetConfig.LeadRollover cfg) {
try {
runRolloverUnguarded(p, cfg);
} catch (RuntimeException e) {
log.warn("lead-rollover: continuation for token={} lead={} threw {} — the roll is dead; "
+ "no further step in this continuation will run",
p.token(), p.leadTerminal(), e.toString(), e);
outcomes.put(p.token(), new RollStatus(RollState.FAILED,
"the roll's continuation threw " + e.toString() + " — the roll is dead and will "
+ "not retry itself; check the daemon log for the stack trace, then open() "
+ "a fresh rollover request"));
}
}
/** The actual body of {@link #runRollover}, unwrapped — see that method's javadoc for the catch. */
private void runRolloverUnguarded(PendingRollover p, FleetConfig.LeadRollover cfg) {
String lead = p.leadTerminal();
long rollStartMillis = nowMillis.getAsLong();
TurnSettleResult turnResult = waitUntilAtTurnBoundary(lead, cfg.turnSettleSeconds());
@@ -271,8 +271,13 @@ public final class LeadHeartbeatLoop {
* does not evaluate the lead's idle state before the fleet has settled.
*/
public void start() {
log.info("idle-lead heartbeat: on — nudge lead after {}s idle (recheck {}ms, quiet cap {})",
TimeUnit.NANOSECONDS.toSeconds(idleAfterNanos), backoffMs, quietNudgeCap);
// fleetd #613: contextHighNudge added alongside the three settings already here — an
// operator otherwise cannot tell from the boot log whether the #609 handover notice is
// armed, and had to load the deployed jar's config to confirm it.
log.info("idle-lead heartbeat: on — nudge lead after {}s idle (recheck {}ms, quiet cap {}, "
+ "context-high nudge {})",
TimeUnit.NANOSECONDS.toSeconds(idleAfterNanos), backoffMs, quietNudgeCap,
contextHighNudge);
scheduler.schedule(this::tick, backoffMs, TimeUnit.MILLISECONDS);
}
@@ -0,0 +1,252 @@
package dev.ltms.fleet;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.fleet.config.FleetConfig;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #613: a {@code MemberRole} with no {@code fleet.<role>s:} pool falls back to
* <em>every</em> configured profile ({@code FleetConfig#candidateProfiles}), and one with no
* {@code fleet.charters.<role>:} entry runs with only the launcher's reply charter. Both are
* deliberate, legitimate states — neither is refused by {@code validateMembers()} — but both were
* silent at boot. On the host that opened this ticket, an unqualified {@code hunter} spawn silently
* widened to all 8 configured profiles and its resolved first choice was {@code local}, a profile
* every other pool on that same config gave weight 0 to.
*
* <p>{@link Fleetd#reportRoleFallbackGaps} must name every gapped role, and for a pool gap, the
* exact resolved first-choice profile — that number, not the pool size, is what actually surprised
* the operator. Mirrors {@link ExhaustedPatternGapReportTest}'s pattern: capture the real log via a
* {@link ListAppender} rather than asserting on the call site's source text.
*/
class RoleFallbackGapReportTest {
private static FleetConfig load(Path dir, String yaml) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, yaml);
return FleetConfig.load(f);
}
/**
* The level this logger had before {@link #attach()} raised it, so {@link #detach} can put it
* back. {@code null} means "inherit from the parent" — the state this logger starts in.
*/
private static Level originalLevel;
/**
* {@code reportRoleFallbackGaps} logs at INFO, and {@code logback-test.xml} sets
* {@code dev.ltms.fleet} to WARN — so INFO events are dropped by the level check before any
* appender sees them. Raise the level for the duration of the test, exactly like {@code
* GitHostShapeReportTest#attach}.
*/
private static ListAppender<ILoggingEvent> attach() {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
originalLevel = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(ListAppender<ILoggingEvent> appender) {
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
logger.detachAppender(appender);
logger.setLevel(originalLevel);
}
private static List<String> infoMessages(ListAppender<ILoggingEvent> appender) {
return appender.list.stream()
.filter(e -> e.getLevel() == Level.INFO)
.map(ILoggingEvent::getFormattedMessage)
.toList();
}
/**
* Reproduces the shape measured in the ticket: {@code dev}, {@code reviewer} and
* {@code architect} each have a pool and a charter; {@code hunter} has neither. The pool-gap
* line must name {@code hunter}, the profile count (3), and the resolved first choice
* ({@code local}, the first profile in definition order) — and must not name the three healthy
* roles. The charter-gap line must separately name only {@code hunter}.
*/
@Test
void hunterWithNoPoolOrCharterIsNamedWithItsResolvedFirstChoice(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
sonnet:
baseUrl: https://llm.ltms.dev/v1
terra:
baseUrl: https://llm.ltms.dev/v1
fleet:
developers:
a:
profile: sonnet
reviewers:
b:
profile: terra
architects:
c:
profile: sonnet
charters:
dev: "dev charter text"
reviewer: "reviewer charter text"
architect: "architect charter text"
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportRoleFallbackGaps(cfg);
} finally {
detach(appender);
}
List<String> infos = infoMessages(appender);
String poolLine = infos.stream()
.filter(m -> m.contains("no fleet.<role>s: pool"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a pool-gap INFO line: " + infos));
assertTrue(poolLine.contains("hunter"), poolLine);
assertTrue(poolLine.contains("3"), "must name the profile count the gap opens onto: " + poolLine);
assertTrue(poolLine.contains("'local'"),
"must name the resolved first-choice profile, the number that actually surprised "
+ "the operator: " + poolLine);
for (String healthy : List.of("dev", "reviewer", "architect")) {
assertFalse(poolLine.contains(healthy),
"pool-gap line must not name a role that has a pool: " + poolLine);
}
String charterLine = infos.stream()
.filter(m -> m.contains("no fleet.charters: entry"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a charter-gap INFO line: " + infos));
assertTrue(charterLine.contains("hunter"), charterLine);
for (String healthy : List.of("dev", "reviewer", "architect")) {
assertFalse(charterLine.contains(healthy),
"charter-gap line must not name a role that has a charter: " + charterLine);
}
}
/**
* The resolved first choice must be genuinely computed from definition order, not hardcoded —
* reordering {@code profiles:} so a different entry comes first changes the reported choice.
*/
@Test
void theResolvedFirstChoiceFollowsProfileDefinitionOrder(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
sonnet:
baseUrl: https://llm.ltms.dev/v1
local:
baseUrl: http://gx00.gw:8000
fleet:
developers:
a:
profile: sonnet
charters:
dev: "dev charter text"
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportRoleFallbackGaps(cfg);
} finally {
detach(appender);
}
String poolLine = infoMessages(appender).stream()
.filter(m -> m.contains("no fleet.<role>s: pool"))
.findFirst()
.orElseThrow();
// hunter, reviewer and architect all lack a pool here; each falls back to the full 2-profile
// set and the first-choice is 'sonnet' because it is first in profiles: definition order.
assertTrue(poolLine.contains("'sonnet'"), poolLine);
assertFalse(poolLine.contains("'local'"), poolLine);
}
/** A config with a pool and a charter for every role produces no role-fallback log at all. */
@Test
void everyRoleWithAPoolAndACharterProducesNoLogAtAll(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
sonnet:
baseUrl: https://llm.ltms.dev/v1
fleet:
developers:
a:
profile: sonnet
reviewers:
b:
profile: sonnet
hunters:
c:
profile: sonnet
architects:
d:
profile: sonnet
charters:
dev: "dev charter text"
reviewer: "reviewer charter text"
hunter: "hunter charter text"
architect: "architect charter text"
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportRoleFallbackGaps(cfg);
} finally {
detach(appender);
}
assertTrue(appender.list.isEmpty(),
"a config with no gaps must not print a per-role block: " + infoMessages(appender));
}
/**
* A config with only {@code profiles:} and no {@code fleet:} block at all must still be
* reported (every role is gapped, both pool and charter) rather than throwing — this is the
* exact shape {@code candidateProfiles}' fallback exists to keep starting.
*/
@Test
void aConfigWithNoFleetBlockAtAllReportsEveryRoleGapped(@TempDir Path dir) throws Exception {
FleetConfig cfg = load(dir, """
profiles:
sonnet:
baseUrl: https://llm.ltms.dev/v1
""");
ListAppender<ILoggingEvent> appender = attach();
try {
Fleetd.reportRoleFallbackGaps(cfg);
} finally {
detach(appender);
}
List<String> infos = infoMessages(appender);
String poolLine = infos.stream()
.filter(m -> m.contains("no fleet.<role>s: pool"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a pool-gap INFO line: " + infos));
String charterLine = infos.stream()
.filter(m -> m.contains("no fleet.charters: entry"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected a charter-gap INFO line: " + infos));
for (String role : List.of("dev", "hunter", "reviewer", "architect")) {
assertTrue(poolLine.contains(role), poolLine);
assertTrue(charterLine.contains(role), charterLine);
}
}
}
@@ -1359,4 +1359,92 @@ class LeadRolloverTest {
+ "it were the measured wait duration: " + message);
}
}
// ---- fleetd #615: a HerdrException out of either unwrapped agents.send call must leave a ----
// ---- TERMINAL FAILED outcome, never a stuck IN_PROGRESS ---------------------------------------
@Test
@DisplayName("[fleetd #615 — 1] send() throwing on the /clear call leaves status(token) "
+ "reporting FAILED, not stuck at IN_PROGRESS")
void sendThrowingOnClearLeavesStatusReportingFailed() throws IOException {
FakeHerdr fake = new FakeHerdr(); // default idle — the turn-settle wait passes immediately
HerdrClient throwsOnClear = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) throws HerdrException {
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
throw new HerdrException("simulated herdr transport failure sending /clear");
}
return fake.call(method, params);
}
@Override
public void close() {
fake.close();
}
};
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(throwsOnClear, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "every synchronous gate passes; the throw happens only "
+ "inside the deferred continuation, which this test's synchronous runner has "
+ "already run to completion by the time confirm() returns");
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.FAILED, status.state(),
"a HerdrException out of the /clear send must leave a TERMINAL FAILED outcome — "
+ "before fleetd #615's fix, the continuation thread died silently and "
+ "status() was stuck reporting the IN_PROGRESS confirm() wrote at hand-off, "
+ "forever: got " + status.state() + " / " + status.detail());
assertNotEquals(LeadRollover.RollState.IN_PROGRESS, status.state());
assertTrue(status.detail().contains("HerdrException"), "the detail must name the exception "
+ "so an operator reading status() has something to act on: " + status.detail());
}
@Test
@DisplayName("[fleetd #615 — 2] send() throwing on the bootstrap-text call (after /clear "
+ "succeeded and the pane settled) also leaves status(token) reporting FAILED — a "
+ "DIFFERENT exit from the /clear-throw case above")
void sendThrowingOnBootstrapTextLeavesStatusReportingFailed() throws IOException {
FakeHerdr fake = new FakeHerdr(); // default idle throughout — both settle waits pass promptly
HerdrClient throwsOnBootstrapText = new HerdrClient() {
@Override
public JsonNode call(String method, Object params) throws HerdrException {
if ("agent.prompt".equals(method) && String.valueOf(params).contains("read the handover file")) {
throw new HerdrException("simulated herdr transport failure sending bootstrapText");
}
return fake.call(method, params);
}
@Override
public void close() {
fake.close();
}
};
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
LeadRollover rollover = newRollover(throwsOnBootstrapText, cfg(handover.toString()), fixedClock(clock));
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "every synchronous gate passes; the throw happens only "
+ "inside the deferred continuation, which this test's synchronous runner has "
+ "already run to completion by the time confirm() returns");
assertEquals(1, promptCallCount(fake), "sanity: /clear was sent and settled — only the "
+ "SECOND agent.prompt call (bootstrapText) threw");
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.FAILED, status.state(),
"a HerdrException out of the bootstrapText send — a DIFFERENT exit from the /clear "
+ "throw, reached only after /clear already succeeded and the pane already "
+ "settled — must also leave a TERMINAL FAILED outcome, not a stuck "
+ "IN_PROGRESS: got " + status.state() + " / " + status.detail());
assertNotEquals(LeadRollover.RollState.IN_PROGRESS, status.state());
assertTrue(status.detail().contains("HerdrException"), "the detail must name the exception "
+ "so an operator reading status() has something to act on: " + status.detail());
}
}
@@ -1,5 +1,9 @@
package dev.ltms.fleet.msg;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.fleet.herdr.AgentControl;
@@ -12,6 +16,7 @@ import dev.ltms.fleet.session.MemberSession;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.List;
@@ -107,6 +112,58 @@ class LeadHeartbeatLoopTest {
"never inject into a state the loop cannot read");
}
// ── fleetd #613: the boot line names all four heartbeat settings ──────────────────────────
/**
* fleetd #613: {@link LeadHeartbeatLoop#start()}'s boot line named only 3 of the 4 constructor
* settings — {@code contextHighNudge} (fleetd #609) was missing, so an operator could not tell
* from the log whether the handover notice was armed. Captures the real log via a
* {@link ListAppender}, raising the logger's level past {@code logback-test.xml}'s
* {@code dev.ltms.fleet -> WARN} override for the duration of the call — the same seam {@code
* GitHostShapeReportTest#attach} uses for its own INFO-level boot line.
*/
private static String heartbeatBootLine(boolean contextHighNudge, ScheduledExecutorService scheduler) {
Logger logger = (Logger) LoggerFactory.getLogger(LeadHeartbeatLoop.class);
Level originalLevel = logger.getLevel();
logger.setLevel(Level.INFO);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
LeadHeartbeatLoop l = new LeadHeartbeatLoop(
new PrimaryRegistry("term_lead"), null /*agents*/, null /*inbox*/, List::of,
null /*pushLoop*/, scheduler, () -> 0L, IDLE_AFTER_NANOS, 1_000L, 3, null,
LeadHeartbeatLoop.LeadContextSource.none(), contextHighNudge);
l.start();
} finally {
logger.detachAppender(appender);
logger.setLevel(originalLevel);
}
return appender.list.stream()
.filter(e -> e.getLevel() == Level.INFO)
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.startsWith("idle-lead heartbeat:"))
.findFirst()
.orElseThrow(() -> new AssertionError("expected the heartbeat boot line to be logged"));
}
@Test
void theBootLineNamesContextHighNudgeWhenArmed() {
String line = heartbeatBootLine(true, scheduler);
assertTrue(line.contains("300s idle"), line);
assertTrue(line.contains("recheck 1000ms"), line);
assertTrue(line.contains("quiet cap 3"), line);
assertTrue(line.contains("context-high nudge true"),
"the boot line must name the 4th setting, contextHighNudge, alongside the other "
+ "three: " + line);
}
@Test
void theBootLineNamesContextHighNudgeWhenOff() {
String line = heartbeatBootLine(false, scheduler);
assertTrue(line.contains("context-high nudge false"), line);
}
// ── (b) an idle lead within the quiet period is not yet injected ───────────────────────────
@Test
@@ -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),