Compare commits
10 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| b6b006c651 | |||
| 7d9a807243 | |||
| 8915e40c7d | |||
| 6cb31a10e4 | |||
| 8368a274a0 | |||
| 17127efb88 | |||
| 203f034528 | |||
| 9ee16f5b85 | |||
| 388ef5a3c3 | |||
| 987ccef4c7 |
@@ -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),
|
||||
|
||||
Reference in New Issue
Block a user