Resolves the one conflict in Fleetd.java. main's side is the inline boot block that Unit A had already moved into FleetdAssembly.assembleAndStart, so the resolution keeps Unit A's single assembly call. That resolution is not purely mechanical. #622 (fleetd #621) landed on main AFTER Unit A forked, and it added boolean requireOperatorConfirm = cfg.leadRollover() == null || cfg.leadRollover().requireOperatorConfirm(); plus a 14th argument to the LeadHeartbeatLoop constructor, inside the very block Unit A moved. Taking Unit A's side alone would have dropped both and silently reverted the operator's #621 fix: the 13-argument overload still exists and delegates with `true`, so the daemon would go back to telling every lead to ask the operator before a context roll. Both are carried into FleetdAssembly here. Measured: with the carried line removed, the full suite is `Tests run: 1883, Failures: 0, Errors: 0` — nothing pins it. That is the fleetd #612 defect shape applied to #621's own wiring, and it is filed separately rather than fixed here, because this commit is a merge resolution and must not also introduce new tests. Full suite on this resolved tree: Tests run: 1883, Failures: 0, Errors: 0.
This commit is contained in:
@@ -167,19 +167,59 @@ fails.
|
|||||||
- **There is no terminal or session parameter, on purpose.** The pane is always your own, resolved
|
- **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.
|
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
|
- **`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`
|
are confident. Ask, wait for the answer, then pass what they said.
|
||||||
defaults to `true` and this is the only thing standing between a judgement call and a wiped
|
|
||||||
session.
|
- **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.
|
- **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.
|
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
|
- **The bootstrap prompt works end to end. Measured 2026-09-22.** This used to say the fix was
|
||||||
real rollover, on 2026-09-12, joined `/clear` and the bootstrap text into one line and Claude Code
|
unproven (fleetd #489) and told you to expect a failure. That is no longer true. The daemon log
|
||||||
refused it as `Unknown command: /clearFresh`. The pane was never cleared and no context was lost,
|
now holds four `lead-rollover: rolled` lines, and three of them ran on 2026-09-22 at 10:01:43,
|
||||||
so the failure was safe — the roll simply did nothing. PR #490 fixed the cause and is deployed,
|
10:38:28 and 11:15:47. Each one cleared the old lead and started a fresh session against the
|
||||||
but no roll has bootstrapped a fresh session end to end yet. **Assume it may still fail, and tell
|
handover file, with the configured `bootstrapText` arriving as its first message. No context was
|
||||||
the operator so before you confirm.** The recovery is the same either way: the file is already
|
lost. The old `Unknown command: /clearFresh` failure from 2026-09-12 does not appear in the log
|
||||||
written, so the operator starts a session and points it at the file. That is why you write the
|
at all. Re-measure both numbers with:
|
||||||
file before you confirm, and never the other way round.
|
|
||||||
|
```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
|
## Writing style
|
||||||
|
|
||||||
|
|||||||
@@ -392,13 +392,21 @@ final class FleetdAssembly {
|
|||||||
if (cfg.leadHeartbeat() != null) {
|
if (cfg.leadHeartbeat() != null) {
|
||||||
var hb = cfg.leadHeartbeat();
|
var hb = cfg.leadHeartbeat();
|
||||||
var leadContextGauge = new LeadContextGauge();
|
var leadContextGauge = new LeadContextGauge();
|
||||||
|
// fleetd #621: the context-high notice's own wording must track this same effective
|
||||||
|
// value — LeadRollover.confirm(...) already gates the roll on it (LeadRollover.java:480),
|
||||||
|
// and absent `leadRollover:` entirely the roll is unusable regardless (NOT_CONFIGURED),
|
||||||
|
// so `true` (the FleetConfig.LeadRollover default) is the safe, byte-identical fallback.
|
||||||
|
// Carried in from Fleetd.main when #612 Unit A merged main: #622 added this line to the
|
||||||
|
// block Unit A had already moved here, so the merge would otherwise have silently
|
||||||
|
// dropped it — with a fully green suite, because nothing pins it (see the follow-up issue).
|
||||||
|
boolean requireOperatorConfirm = cfg.leadRollover() == null || cfg.leadRollover().requireOperatorConfirm();
|
||||||
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
|
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
|
||||||
pushLoop, heartbeatScheduler, ports.nanoClock(),
|
pushLoop, heartbeatScheduler, ports.nanoClock(),
|
||||||
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
|
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
|
||||||
metrics,
|
metrics,
|
||||||
Fleetd.leadContextSource(leadContextGauge, router.leadAgents(), leads,
|
Fleetd.leadContextSource(leadContextGauge, router.leadAgents(), leads,
|
||||||
Fleetd.leadConfigDirLookup(() -> config.get().profiles(), leaders)),
|
Fleetd.leadConfigDirLookup(() -> config.get().profiles(), leaders)),
|
||||||
Boolean.TRUE.equals(hb.contextHighNudge()));
|
Boolean.TRUE.equals(hb.contextHighNudge()), requireOperatorConfirm);
|
||||||
heartbeat.start();
|
heartbeat.start();
|
||||||
} else {
|
} else {
|
||||||
heartbeat = null;
|
heartbeat = null;
|
||||||
|
|||||||
@@ -2240,11 +2240,11 @@ public record FleetConfig(
|
|||||||
* information — which profiles, and now both values, so they can fix it without reading the
|
* information — which profiles, and now both values, so they can fix it without reading the
|
||||||
* source — without ever taking the fleet down.
|
* source — without ever taking the fleet down.
|
||||||
*
|
*
|
||||||
* <p>Which of the two inputs Claude Code actually follows when they disagree is intentionally
|
* <p>fleetd #618 measured which of the two inputs Claude Code actually follows when they
|
||||||
* <em>not</em> asserted here. {@code ClaudeCodeArguments}'s javadoc used to state the
|
* disagree: the environment variable wins, so {@code autoCompactWindow} is inert on a profile
|
||||||
* environment variable always wins; nobody had measured that, and this host's own
|
* that also sets the env var. This method only detects and reports the disagreement — it does
|
||||||
* {@code fleetd.yaml} asserts the opposite in a comment. This method only detects and reports
|
* not correct it — see {@link dev.ltms.fleet.launch.ClaudeCodeArguments} for the full measured
|
||||||
* the disagreement — see {@link dev.ltms.fleet.launch.ClaudeCodeArguments}.
|
* precedence.
|
||||||
*
|
*
|
||||||
* <p>Equal values never warn: either input then produces the same session window, so there is
|
* <p>Equal values never warn: either input then produces the same session window, so there is
|
||||||
* nothing to reconcile.
|
* nothing to reconcile.
|
||||||
@@ -2285,9 +2285,11 @@ public record FleetConfig(
|
|||||||
names.sort(String::compareTo);
|
names.sort(String::compareTo);
|
||||||
detail.sort(String::compareTo);
|
detail.sort(String::compareTo);
|
||||||
log.warn("Claude Code profile(s) {} set disagreeing autoCompactWindow and env."
|
log.warn("Claude Code profile(s) {} set disagreeing autoCompactWindow and env."
|
||||||
+ "CLAUDE_CODE_AUTO_COMPACT_WINDOW — the daemon starts anyway. Fix by "
|
+ "CLAUDE_CODE_AUTO_COMPACT_WINDOW — the daemon starts anyway: {}. fleetd "
|
||||||
+ "removing one key or setting equal values on each: {}. Which input Claude "
|
+ "#618 measured that CLAUDE_CODE_AUTO_COMPACT_WINDOW wins, so "
|
||||||
+ "Code actually follows when they disagree is not verified here.",
|
+ "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));
|
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.
|
* 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
|
* <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
|
* disagree, and fleetd #618 measured which one Claude Code actually follows: the environment
|
||||||
* javadoc used to claim the environment variable always wins, but nobody had measured that, and
|
* variable wins, ahead of this {@code --autocompact} flag, ahead of the settings file, ahead of
|
||||||
* this host's own {@code fleetd.yaml} asserts the opposite in a comment. So this javadoc no
|
* clientdata, the experiment, and the model default. So when a profile sets both, the flag this
|
||||||
* longer picks a side. {@link FleetConfig#load(java.nio.file.Path)} only WARNS when a Claude
|
* method appends has NO effect — Claude Code reads {@code CLAUDE_CODE_AUTO_COMPACT_WINDOW}
|
||||||
* Code profile sets both to different values (see {@code
|
* 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,
|
* 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) {
|
public static List<String> withAutoCompactWindow(List<String> argv, FleetConfig.Profile profile) {
|
||||||
if (profile.autoCompactWindow() == null) {
|
if (profile.autoCompactWindow() == null) {
|
||||||
|
|||||||
@@ -76,6 +76,7 @@ public final class LeadHeartbeatLoop {
|
|||||||
private final Metrics metrics; // CB-512 pattern: nullable — no registry in unit tests
|
private final Metrics metrics; // CB-512 pattern: nullable — no registry in unit tests
|
||||||
private final LeadContextSource contextSource; // fleetd #609
|
private final LeadContextSource contextSource; // fleetd #609
|
||||||
private final boolean contextHighNudge; // fleetd #609: opt-in, like the loop itself
|
private final boolean contextHighNudge; // fleetd #609: opt-in, like the loop itself
|
||||||
|
private final boolean requireOperatorConfirm; // fleetd #621: mirrors leadRollover.requireOperatorConfirm
|
||||||
|
|
||||||
/** When the current idle stretch began (nanos), or {@link #NOT_IDLE}. Single scheduler thread only. */
|
/** When the current idle stretch began (nanos), or {@link #NOT_IDLE}. Single scheduler thread only. */
|
||||||
private long idleSinceNanos = NOT_IDLE;
|
private long idleSinceNanos = NOT_IDLE;
|
||||||
@@ -106,12 +107,32 @@ public final class LeadHeartbeatLoop {
|
|||||||
* fleetd #609: as above, plus the lead's own context source and whether a HIGH reading should
|
* fleetd #609: as above, plus the lead's own context source and whether a HIGH reading should
|
||||||
* append a hand-over notice to the loop's nudge. Pass {@link LeadContextSource#none()} and
|
* append a hand-over notice to the loop's nudge. Pass {@link LeadContextSource#none()} and
|
||||||
* {@code false} to keep the pre-#609 behaviour exactly (both existing public constructors do).
|
* {@code false} to keep the pre-#609 behaviour exactly (both existing public constructors do).
|
||||||
|
*
|
||||||
|
* <p>fleetd #621: delegates to the full constructor with {@code requireOperatorConfirm=true} —
|
||||||
|
* the pre-#621 wording ("ask the operator ... only the operator can approve the roll") assumed
|
||||||
|
* the config default, so every caller of this overload keeps that text byte-identical.
|
||||||
*/
|
*/
|
||||||
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||||
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
|
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
|
||||||
ScheduledExecutorService scheduler, LongSupplier clock,
|
ScheduledExecutorService scheduler, LongSupplier clock,
|
||||||
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
|
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
|
||||||
LeadContextSource contextSource, boolean contextHighNudge) {
|
LeadContextSource contextSource, boolean contextHighNudge) {
|
||||||
|
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock,
|
||||||
|
idleAfterNanos, backoffMs, quietNudgeCap, metrics, contextSource, contextHighNudge, true);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* fleetd #621: as above, plus the daemon's effective {@code leadRollover.requireOperatorConfirm}
|
||||||
|
* value — threaded into {@link #contextNotice(boolean, LeadContextGauge.Reading, boolean, boolean)}
|
||||||
|
* so the notice's wording tracks the config the daemon actually enforces (see {@code
|
||||||
|
* LeadRollover.confirm}) instead of always asserting the operator gate is on.
|
||||||
|
*/
|
||||||
|
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||||
|
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
|
||||||
|
ScheduledExecutorService scheduler, LongSupplier clock,
|
||||||
|
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
|
||||||
|
LeadContextSource contextSource, boolean contextHighNudge,
|
||||||
|
boolean requireOperatorConfirm) {
|
||||||
this.primaryRegistry = primaryRegistry;
|
this.primaryRegistry = primaryRegistry;
|
||||||
this.agents = agents;
|
this.agents = agents;
|
||||||
this.inbox = inbox;
|
this.inbox = inbox;
|
||||||
@@ -125,6 +146,7 @@ public final class LeadHeartbeatLoop {
|
|||||||
this.metrics = metrics;
|
this.metrics = metrics;
|
||||||
this.contextSource = contextSource;
|
this.contextSource = contextSource;
|
||||||
this.contextHighNudge = contextHighNudge;
|
this.contextHighNudge = contextHighNudge;
|
||||||
|
this.requireOperatorConfirm = requireOperatorConfirm;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -344,7 +366,7 @@ public final class LeadHeartbeatLoop {
|
|||||||
// d.contextNotified() is the value to persist once delivery is confirmed, not the value the text
|
// d.contextNotified() is the value to persist once delivery is confirmed, not the value the text
|
||||||
// itself should be built from. Otherwise a HIGH stretch that is still latched would never see the
|
// itself should be built from. Otherwise a HIGH stretch that is still latched would never see the
|
||||||
// notice at all, defeating the very check this fixes.
|
// notice at all, defeating the very check this fixes.
|
||||||
String notice = contextNotice(contextHighNudge, reading, contextNotified);
|
String notice = contextNotice(contextHighNudge, reading, contextNotified, requireOperatorConfirm);
|
||||||
var lead = primaryRegistry.primaryTerminal();
|
var lead = primaryRegistry.primaryTerminal();
|
||||||
boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
|
boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
|
||||||
// The latch becomes true only when all three hold: decide() chose to notify, a notice was
|
// The latch becomes true only when all three hold: decide() chose to notify, a notice was
|
||||||
@@ -381,8 +403,9 @@ public final class LeadHeartbeatLoop {
|
|||||||
/**
|
/**
|
||||||
* fleetd #609: the text appended to a nudge when the lead's own context is full — {@code ""}
|
* fleetd #609: the text appended to a nudge when the lead's own context is full — {@code ""}
|
||||||
* whenever the notice does not apply, so callers can unconditionally append this without an extra
|
* whenever the notice does not apply, so callers can unconditionally append this without an extra
|
||||||
* branch. Wording stays plain (CEFR B1) and honest that only the operator approves a roll — this
|
* branch. Wording stays plain (CEFR B1) and honest about who actually gates the roll — see the
|
||||||
* loop only ever prints text, it never calls {@code fleet_handover} itself.
|
* {@code requireOperatorConfirm} overload (fleetd #621) for which check that is. This loop only
|
||||||
|
* ever prints text, it never calls {@code fleet_handover} itself.
|
||||||
*
|
*
|
||||||
* @param enabled the {@code leadHeartbeat.contextHighNudge} config flag
|
* @param enabled the {@code leadHeartbeat.contextHighNudge} config flag
|
||||||
* @param reading the lead's current {@link LeadContextGauge} reading
|
* @param reading the lead's current {@link LeadContextGauge} reading
|
||||||
@@ -401,9 +424,33 @@ public final class LeadHeartbeatLoop {
|
|||||||
* closing sentence ("You will not be told again until your context reads ok.") false. {@link
|
* closing sentence ("You will not be told again until your context reads ok.") false. {@link
|
||||||
* #injectNudge} is the only caller that passes a non-default {@code alreadyNotified}.
|
* #injectNudge} is the only caller that passes a non-default {@code alreadyNotified}.
|
||||||
*
|
*
|
||||||
|
* <p>fleetd #621: delegates with {@code requireOperatorConfirm=true} — the pre-#621 default and the
|
||||||
|
* value every existing caller of this overload (including every test written before #621) already
|
||||||
|
* assumed, so the text this overload returns stays byte-identical.
|
||||||
|
*
|
||||||
* @param alreadyNotified whether the lead has already been told about the current HIGH stretch
|
* @param alreadyNotified whether the lead has already been told about the current HIGH stretch
|
||||||
*/
|
*/
|
||||||
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified) {
|
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified) {
|
||||||
|
return contextNotice(enabled, reading, alreadyNotified, true);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* fleetd #621: as {@link #contextNotice(boolean, LeadContextGauge.Reading, boolean)}, but the closing
|
||||||
|
* instructions also track the daemon's effective {@code leadRollover.requireOperatorConfirm} value,
|
||||||
|
* instead of always asserting that only the operator can approve the roll.
|
||||||
|
*
|
||||||
|
* <p>{@code LeadRollover.confirm(...)} already honours this flag: when it is {@code false}, the daemon
|
||||||
|
* itself gates the roll on the three handover-file checks alone (exists, modified after the {@code
|
||||||
|
* open()} request, and no older than {@code maxDocAgeSeconds}) and never consults {@code
|
||||||
|
* operatorConfirmed}. Before this parameter existed, this notice told the lead to ask the operator
|
||||||
|
* regardless — so a lead that followed its own instructions asked anyway, and setting the config knob
|
||||||
|
* to {@code false} stopped the daemon refusing the roll without stopping the operator being
|
||||||
|
* interrupted. This parameter is how the text is kept honest about which gate is actually live.
|
||||||
|
*
|
||||||
|
* @param requireOperatorConfirm the effective {@code leadRollover.requireOperatorConfirm} value
|
||||||
|
*/
|
||||||
|
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified,
|
||||||
|
boolean requireOperatorConfirm) {
|
||||||
if (!enabled || alreadyNotified || reading.state() != LeadContextGauge.State.HIGH) {
|
if (!enabled || alreadyNotified || reading.state() != LeadContextGauge.State.HIGH) {
|
||||||
return "";
|
return "";
|
||||||
}
|
}
|
||||||
@@ -420,10 +467,18 @@ public final class LeadHeartbeatLoop {
|
|||||||
sb.append(" (").append(reading.compactions()).append(' ').append(compactionWord)
|
sb.append(" (").append(reading.compactions()).append(' ').append(compactionWord)
|
||||||
.append(" so far).");
|
.append(" so far).");
|
||||||
}
|
}
|
||||||
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
if (requireOperatorConfirm) {
|
||||||
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
|
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
||||||
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
|
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
|
||||||
+ "again until your context reads ok.");
|
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
|
||||||
|
+ "again until your context reads ok.");
|
||||||
|
} else {
|
||||||
|
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
||||||
|
+ "write the file it names, then call fleet_handover(action=\"confirm\", token). Decide for "
|
||||||
|
+ "yourself when to confirm: the roll goes through if the handover file exists, was "
|
||||||
|
+ "changed after you opened it, and is not older than maxDocAgeSeconds. You will not be "
|
||||||
|
+ "told again until your context reads ok.");
|
||||||
|
}
|
||||||
return sb.toString();
|
return sb.toString();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -476,6 +476,26 @@ class LeadHeartbeatLoopTest {
|
|||||||
assertTrue(notice.contains("2 compactions"), notice);
|
assertTrue(notice.contains("2 compactions"), notice);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ── fleetd #621: the notice must track the effective requireOperatorConfirm value ─────────────
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void contextNoticeKeepsAskingTheOperatorWhenRequireOperatorConfirmIsTrue() {
|
||||||
|
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||||
|
String notice = LeadHeartbeatLoop.contextNotice(true, reading, false, true);
|
||||||
|
assertTrue(notice.contains("ask the operator"), notice);
|
||||||
|
assertTrue(notice.contains("Only the operator can approve the roll"), notice);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void contextNoticeDropsTheOperatorAskWhenRequireOperatorConfirmIsFalse() {
|
||||||
|
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||||
|
String notice = LeadHeartbeatLoop.contextNotice(true, reading, false, false);
|
||||||
|
assertFalse(notice.contains("ask the operator"), notice);
|
||||||
|
assertFalse(notice.contains("Only the operator can approve the roll"), notice);
|
||||||
|
assertTrue(notice.contains("fleet_handover"), notice);
|
||||||
|
assertTrue(notice.contains("maxDocAgeSeconds"), notice);
|
||||||
|
}
|
||||||
|
|
||||||
// ── fleetd #609 review: the latch must mean "the notice reached the pane" ────────────────────
|
// ── fleetd #609 review: the latch must mean "the notice reached the pane" ────────────────────
|
||||||
//
|
//
|
||||||
// These four drive LeadHeartbeatLoop.tick() directly (package-private, same reasoning as
|
// These four drive LeadHeartbeatLoop.tick() directly (package-private, same reasoning as
|
||||||
|
|||||||
@@ -1846,7 +1846,8 @@ class MessageServiceTest {
|
|||||||
* but backed by {@link ManualScheduler} instead of a real timer (fleetd #608): its tick never
|
* 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()}.
|
* 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 {
|
implements AutoCloseable {
|
||||||
@Override
|
@Override
|
||||||
public void close() {
|
public void close() {
|
||||||
@@ -1861,14 +1862,20 @@ class MessageServiceTest {
|
|||||||
* {@code anAlreadyCollectedTicketProducesNoNudge}, which sets it to 1 to prove exactly that.
|
* {@code anAlreadyCollectedTicketProducesNoNudge}, which sets it to 1 to prove exactly that.
|
||||||
*/
|
*/
|
||||||
private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs) {
|
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);
|
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||||
registry.recordDelegation(T, LEAD);
|
registry.recordDelegation(T, LEAD);
|
||||||
FakeHerdr leadHerdr = new FakeHerdr();
|
FakeHerdr leadHerdr = new FakeHerdr();
|
||||||
AgentControl leadAgents = new AgentControl(leadHerdr);
|
AgentControl leadAgents = new AgentControl(leadHerdr);
|
||||||
ManualScheduler scheduler = new ManualScheduler();
|
ManualScheduler scheduler = new ManualScheduler();
|
||||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
|
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
|
||||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, System::nanoTime);
|
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, nowNanos);
|
||||||
return new ManualPushWiring(service, leadHerdr, scheduler);
|
return new ManualPushWiring(service, leadHerdr, scheduler, pushLoop);
|
||||||
}
|
}
|
||||||
|
|
||||||
private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException {
|
private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException {
|
||||||
@@ -1959,24 +1966,35 @@ class MessageServiceTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
void severalAsyncTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception {
|
void severalAsyncTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception {
|
||||||
try (var wiring = wireWithPushLoop(1, 300)) { // wide backoff: both tickets land before the tick fires
|
// A 1ms backoff is due immediately. ManualScheduler still cannot run it until this test
|
||||||
String first = wiring.service().sendAsync(T, "first task");
|
// explicitly calls runDueTasks(), so both terminal tickets join one scheduled tick.
|
||||||
awaitWaiting();
|
try (var wiring = wireWithManualScheduler(1, 1)) {
|
||||||
injectDelivery();
|
String first;
|
||||||
assertTrue(rendezvous.resolve(T, "first done"));
|
String second;
|
||||||
// Settle without polling: poll() itself marks a ticket collected (that's the point of
|
java.util.concurrent.CountDownLatch firstTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||||
// anAlreadyCollectedTicketProducesNoNudge above) — using it here to detect completion
|
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(firstTerminalReached::countDown);
|
||||||
// would collect the ticket before the coalescing this test checks ever gets a chance.
|
try {
|
||||||
Thread.sleep(100);
|
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");
|
java.util.concurrent.CountDownLatch secondTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||||
awaitWaiting();
|
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(secondTerminalReached::countDown);
|
||||||
injectDelivery();
|
second = wiring.service().sendAsync(T, "second task");
|
||||||
assertTrue(rendezvous.resolve(T, "second done"));
|
awaitWaiting();
|
||||||
Thread.sleep(100);
|
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());
|
assertEquals(1, wiring.scheduler().runDueTasks(),
|
||||||
Thread.sleep(200); // settle — nothing more should arrive beyond the one coalesced nudge
|
"both terminal tickets must coalesce onto one scheduled tick");
|
||||||
long nudgeCount = wiring.leadHerdr().calls.stream()
|
long nudgeCount = wiring.leadHerdr().calls.stream()
|
||||||
.filter(c -> c.method().equals("agent.prompt")).count();
|
.filter(c -> c.method().equals("agent.prompt")).count();
|
||||||
assertEquals(1, nudgeCount, "two tickets finishing together must produce ONE nudge, not two");
|
assertEquals(1, nudgeCount, "two tickets finishing together must produce ONE nudge, not two");
|
||||||
@@ -2055,7 +2073,8 @@ class MessageServiceTest {
|
|||||||
|
|
||||||
@Test
|
@Test
|
||||||
void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception {
|
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");
|
String ticket = wiring.service().sendAsync(T, "task that asks");
|
||||||
awaitWaiting();
|
awaitWaiting();
|
||||||
injectDelivery();
|
injectDelivery();
|
||||||
@@ -2063,8 +2082,9 @@ class MessageServiceTest {
|
|||||||
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
|
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
|
||||||
() -> wiring.service().ask(T, "which config file?", 5000));
|
() -> wiring.service().ask(T, "which config file?", 5000));
|
||||||
MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING);
|
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()
|
long callsBeforeAnswer = wiring.leadHerdr().calls.stream()
|
||||||
.filter(c -> c.method().equals("agent.prompt")).count();
|
.filter(c -> c.method().equals("agent.prompt")).count();
|
||||||
|
|
||||||
@@ -2072,12 +2092,22 @@ class MessageServiceTest {
|
|||||||
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
|
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
|
||||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||||
awaitWaiting();
|
awaitWaiting();
|
||||||
|
java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||||
|
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(terminalReached::countDown);
|
||||||
assertTrue(rendezvous.resolve(T, "done"));
|
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);
|
answer.get(5, TimeUnit.SECONDS);
|
||||||
|
|
||||||
// Let several more ticks (and the ticket's own now-legitimate terminal nudge) fire —
|
// Run the ticket's legitimate terminal nudge and every later scheduled tick through its
|
||||||
// none of them may still name the question's turnId, which is closed.
|
// reminder cap. None may still name the closed question.
|
||||||
Thread.sleep(300);
|
for (int tick = 0; tick < 6; tick++) {
|
||||||
|
assertEquals(1, wiring.scheduler().runDueTasks(), "expected one scheduled reminder tick");
|
||||||
|
}
|
||||||
boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream()
|
boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream()
|
||||||
.filter(c -> c.method().equals("agent.prompt"))
|
.filter(c -> c.method().equals("agent.prompt"))
|
||||||
.skip(callsBeforeAnswer)
|
.skip(callsBeforeAnswer)
|
||||||
@@ -2160,35 +2190,45 @@ class MessageServiceTest {
|
|||||||
// decideTickets hits the cap and STOPs — activeLeads drops the lead, but (before the fix)
|
// 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
|
// 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.
|
// the bug report, reached deterministically rather than by timing it against a live tick.
|
||||||
try (var wiring = wireWithPushLoop(1, 50, clock::get)) {
|
// A 1ms backoff is due immediately, but ManualScheduler runs only the ticks below.
|
||||||
String stale = wiring.service().sendAsync(T, "first task");
|
try (var wiring = wireWithManualScheduler(1, 1, clock::get)) {
|
||||||
awaitWaiting();
|
String stale;
|
||||||
injectDelivery();
|
String fresh;
|
||||||
assertTrue(rendezvous.resolve(T, "stale result"));
|
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
|
// Fire the stale ticket's one nudge, then its cap tick (STOP removes it from activeLeads;
|
||||||
// activeLeads; pendingTickets is untouched either way — that asymmetry is the bug).
|
// pendingTickets is untouched either way — that asymmetry is the bug).
|
||||||
awaitNudge(wiring.leadHerdr());
|
assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one nudge tick");
|
||||||
Thread.sleep(300);
|
assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one cap tick");
|
||||||
assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale),
|
assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale),
|
||||||
"sanity: the stale ticket's own reminder must have fired first");
|
"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.
|
// 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));
|
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
|
||||||
|
|
||||||
// A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes.
|
// A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes.
|
||||||
String fresh = wiring.service().sendAsync(T, "second task");
|
java.util.concurrent.CountDownLatch freshTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||||
awaitWaiting();
|
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(freshTerminalReached::countDown);
|
||||||
injectDelivery();
|
fresh = wiring.service().sendAsync(T, "second task");
|
||||||
assertTrue(rendezvous.resolve(T, "fresh result"));
|
awaitWaiting();
|
||||||
|
injectDelivery();
|
||||||
// The fresh ticket restarts the (now-dormant) reminder loop with its own nudge.
|
assertTrue(rendezvous.resolve(T, "fresh result"));
|
||||||
long before = wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count();
|
assertTrue(freshTerminalReached.await(5, TimeUnit.SECONDS),
|
||||||
long deadline = System.currentTimeMillis() + 3000;
|
"the fresh ticket never reached its terminal phase");
|
||||||
while (wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count() <= before
|
} finally {
|
||||||
&& System.currentTimeMillis() < deadline) {
|
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
|
||||||
Thread.sleep(10);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// The fresh ticket restarts the now-dormant reminder loop with its own nudge.
|
||||||
|
assertEquals(1, wiring.scheduler().runDueTasks(), "the fresh ticket must have one nudge tick");
|
||||||
String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
|
String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
|
||||||
assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge);
|
assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge);
|
||||||
assertFalse(latestNudge.contains(stale),
|
assertFalse(latestNudge.contains(stale),
|
||||||
|
|||||||
Reference in New Issue
Block a user