Merge branch 'worker/737-a263f3-3'
CI / shell-tests (push) Failing after 6s
CI / contract (push) Successful in 1m3s
CI / build (push) Failing after 2m20s

This commit is contained in:
Dai Ha
2026-10-05 06:01:49 +02:00
2 changed files with 207 additions and 32 deletions
@@ -149,9 +149,17 @@ public final class LeadRollover {
* calling lead's workspace, before storing it here; see this class's
* javadoc. This is the value the MCP layer hands back to the lead as
* "write your file here", so callers may rely on it always being absolute.
* @param rolloverKey the key {@link #confirm}'s single-flight claim is taken under. {@link
* #open} resolves this once, here, from {@code leadTerminal} while the
* calling lead is certainly still live: the lead's configured name, or
* {@code leadTerminal} itself when no name resolves. Carried rather than
* recomputed at release time, because by the time a roll's continuation
* releases its claim the OLD terminal may no longer resolve to any name at
* all — recomputing there would release a different key than the one the
* claim was taken under.
*/
public record PendingRollover(String token, String leadTerminal, String handoverPath,
long requestedAtMillis) {}
long requestedAtMillis, String rolloverKey) {}
/** Which check refused a {@link #confirm} call, named so a caller can act on it. */
public enum RefusalReason {
@@ -176,9 +184,9 @@ public final class LeadRollover {
*/
HANDOVER_STALE,
/**
* This lead terminal already has a roll running: an earlier {@link #confirm} call claimed
* it and that roll's continuation has not released it yet. {@code detail} names the lead
* terminal and the token that holds the claim.
* This lead already has a roll running: an earlier {@link #confirm} call claimed its
* single-flight key (see {@link PendingRollover#rolloverKey}) and that roll's continuation
* has not released it yet. {@code detail} names the key and the token that holds the claim.
*/
ROLL_ALREADY_RUNNING
}
@@ -340,9 +348,12 @@ public final class LeadRollover {
private final Function<String, String> leadWorkspace;
/**
* Terminal id → that lead's configured name under {@code fleet.leaders}, or {@code null} when
* the terminal names no currently-recognised lead. The deferred continuation calls this, on the
* OLD terminal, before tearing it down, so it knows which lead to pass to {@link
* LeadLauncher#relaunch}.
* the terminal names no currently-recognised lead. {@link #open} calls this on the calling
* lead's own terminal, while it is certainly still live, to resolve {@link
* PendingRollover#rolloverKey}. The deferred continuation also calls this, on the OLD terminal,
* before tearing it down, so it knows which lead to pass to {@link LeadLauncher#relaunch} — by
* that point the live roster may no longer contain the old terminal, so this lookup can return
* {@code null} here even though {@link #open}'s earlier call against the same terminal did not.
*/
private final Function<String, String> leadNameForTerminal;
/**
@@ -362,14 +373,14 @@ public final class LeadRollover {
private final Consumer<Runnable> continuationRunner;
private final Map<String, PendingRollover> pending = new ConcurrentHashMap<>();
/**
* Lead terminal → the token of the roll currently holding that terminal exclusive, for
* {@link #confirm}'s single-flight claim. {@link #confirm} claims an entry here with an
* atomic put-if-absent once every other gate has passed, refusing with {@link
* {@link PendingRollover#rolloverKey} → the token of the roll currently holding that key
* exclusive, for {@link #confirm}'s single-flight claim. {@link #confirm} claims an entry here
* with an atomic put-if-absent once every other gate has passed, refusing with {@link
* RefusalReason#ROLL_ALREADY_RUNNING} when a claim is already held; {@link #runRollover}
* releases it in a {@code finally}, on both the success and the thrown-exception path. A
* terminal absent from this map has no roll currently in flight for it.
* releases it in a {@code finally}, on both the success and the thrown-exception path. A key
* absent from this map has no roll currently in flight for it.
*/
private final Map<String, String> rollingByTerminal = new ConcurrentHashMap<>();
private final Map<String, String> rollingByLead = new ConcurrentHashMap<>();
/**
* Finished tokens → what actually happened, for {@link #status}. Bounded by {@link
* #OUTCOME_HISTORY_CAP}, oldest evicted first ({@code removeEldestEntry} on an insertion-order
@@ -438,8 +449,9 @@ public final class LeadRollover {
/**
* The lead says it is ready to be replaced. Generates a token and records the resolved
* handover path, this moment's wall-clock timestamp (the baseline {@link #confirm} checks the
* handover file's modified time against), and {@code leadTerminal} — only that exact terminal
* may later {@link #confirm} this token.
* handover file's modified time against), {@code leadTerminal} — only that exact terminal may
* later {@link #confirm} this token — and {@link PendingRollover#rolloverKey}, resolved here
* from {@code leadTerminal} while the calling lead is certainly still live.
*
* @param leadTerminal the calling lead's terminal id, resolved by the MCP layer from the
* connection (see this class's javadoc) — never a client-supplied value
@@ -459,7 +471,9 @@ public final class LeadRollover {
String token = UUID.randomUUID().toString();
long requestedAt = nowMillis.getAsLong();
String resolvedPath = resolveHandoverPath(cfg.handoverPath(), leadTerminal);
PendingRollover p = new PendingRollover(token, leadTerminal, resolvedPath, requestedAt);
String leadName = leadNameForTerminal.apply(leadTerminal);
String rolloverKey = (leadName == null || leadName.isBlank()) ? leadTerminal : leadName;
PendingRollover p = new PendingRollover(token, leadTerminal, resolvedPath, requestedAt, rolloverKey);
pending.put(token, p);
if (resolvedPath.equals(cfg.handoverPath())) {
log.info("lead-rollover: open token={} lead={} handoverPath={} reason={}",
@@ -565,12 +579,11 @@ public final class LeadRollover {
// Single-flight claim: atomic put-if-absent, taken only after every other gate has
// passed, so a refused confirm() never takes it. A non-null previous value means a
// different, still-running roll already holds this lead terminal.
String holder = rollingByTerminal.putIfAbsent(p.leadTerminal(), token);
// different, still-running roll already holds this lead's claim.
String holder = rollingByLead.putIfAbsent(p.rolloverKey(), token);
if (holder != null) {
return RollDecision.refused(RefusalReason.ROLL_ALREADY_RUNNING,
"lead terminal " + p.leadTerminal() + " already has a roll running under token "
+ holder);
"lead '" + p.rolloverKey() + "' already has a roll running under token " + holder);
}
// Record IN_PROGRESS BEFORE removing from `pending` — see RollState#IN_PROGRESS and
@@ -591,14 +604,14 @@ public final class LeadRollover {
} catch (RuntimeException e) {
// continuationRunner can reject the hand-off itself (e.g. a bounded executor's
// RejectedExecutionException) before runRollover ever starts, so runRollover's own
// finally — the only other place that releases rollingByTerminal — never runs either.
// finally — the only other place that releases rollingByLead — never runs either.
// Release the claim here and overwrite the IN_PROGRESS entry with a terminal outcome,
// or this lead terminal could never be rolled again and status() would report
// IN_PROGRESS forever for a roll that in fact never started.
// or this lead could never be rolled again and status() would report IN_PROGRESS
// forever for a roll that in fact never started.
log.warn("lead-rollover: continuationRunner rejected token={} lead={}: {} — the roll "
+ "never started; releasing its claim and reporting it as FAILED",
token, callerTerminal, e.toString(), e);
rollingByTerminal.remove(p.leadTerminal(), token);
rollingByLead.remove(p.rolloverKey(), token);
outcomes.put(token, new RollStatus(RollState.FAILED,
"continuationRunner rejected this roll before it ever started: " + e.toString()
+ " — the roll never ran; open() a fresh rollover request"));
@@ -641,10 +654,10 @@ public final class LeadRollover {
+ "a fresh rollover request"));
} finally {
// Release the single-flight claim on both the normal return and the thrown-exception
// path above — a release only on success would leave this lead terminal unrollable
// forever after one failure. The conditional two-argument remove only clears the
// entry this roll itself holds, never a different roll's claim on the same terminal.
rollingByTerminal.remove(p.leadTerminal(), p.token());
// path above — a release only on success would leave this lead unrollable forever
// after one failure. The conditional two-argument remove only clears the entry this
// roll itself holds, never a different roll's claim on the same key.
rollingByLead.remove(p.rolloverKey(), p.token());
}
}
@@ -16,6 +16,7 @@ import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -1127,8 +1128,13 @@ class LeadRolloverTest {
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
// Only LEAD resolves to a configured name — every other terminal (e.g. the distinct
// term_cap_N terminals evictionCountsInProgressEntriesTowardTheCap opens) falls back to
// keying its own single-flight claim on the terminal itself, exactly like a terminal
// the live roster does not recognise.
return new LeadRollover(agents, spaces, launcher, () -> config, _ -> null,
_ -> LEAD_NAME, () -> liveLeadTerminals, nowMillis, () -> { }, runner);
t -> LEAD.equals(t) ? LEAD_NAME : null, () -> liveLeadTerminals, nowMillis, () -> { },
runner);
}
@Test
@@ -1317,7 +1323,7 @@ class LeadRolloverTest {
+ "so an operator reading status() has something to act on: " + status.detail());
}
// ---- fleetd #726 unit 3: confirm() single-flights per lead terminal -----------------------
// ---- fleetd #726 unit 3: confirm() single-flights one roll at a time per lead ---------------
@Test
@DisplayName("[SINGLE-FLIGHT 1] a second confirm() for the SAME lead terminal is refused with "
@@ -1343,8 +1349,8 @@ class LeadRolloverTest {
assertFalse(secondDecision.accepted());
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, secondDecision.reason());
assertTrue(secondDecision.detail().contains(LEAD), "the refusal detail must name the lead "
+ "terminal: " + secondDecision.detail());
assertTrue(secondDecision.detail().contains(LEAD_NAME), "the refusal detail must name "
+ "the single-flight key (the lead's name): " + secondDecision.detail());
assertTrue(secondDecision.detail().contains(first.token()), "the refusal detail must name "
+ "the token holding the claim: " + secondDecision.detail());
assertEquals(1, runner.heldCount(), "the refused confirm() must never have reached "
@@ -1502,4 +1508,160 @@ class LeadRolloverTest {
+ "this lead terminal could never be rolled again: " + retryDecision.reason() + " / "
+ retryDecision.detail());
}
// ---- fleetd #737 unit 4: confirm() single-flights per LEAD NAME, not per lead terminal ------
@Test
@DisplayName("[NAME-KEYED 1] a second confirm() for the SAME lead is refused with "
+ "ROLL_ALREADY_RUNNING even when it is opened from a DIFFERENT terminal, while the "
+ "first roll's continuation is still in flight")
void secondConfirmForTheSameLeadIsRefusedEvenFromADifferentTerminal() throws IOException {
FakeHerdr herdr = herdrReadyForAFullRoll(); // the held roll WOULD complete once run
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
HoldingRunner runner = new HoldingRunner();
Map<String, String> liveLeadTerminals = Map.of(NEW_TERMINAL, LEAD_NAME);
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
// Both LEAD and OTHER_LEAD resolve to the SAME configured lead name — modelling a roll
// that has already replaced the lead's pane: the fresh pane (here, OTHER_LEAD) is still
// the SAME lead, just a different terminal id.
Function<String, String> leadNameForTerminal = _ -> LEAD_NAME;
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
_ -> null, leadNameForTerminal, () -> liveLeadTerminals, fixedClock(clock), () -> { },
runner);
LeadRollover.PendingRollover first = rollover.open(LEAD, "context is full");
LeadRollover.RollDecision firstDecision = rollover.confirm(LEAD, first.token(), true);
assertTrue(firstDecision.accepted(), "expected approval; got: " + firstDecision.reason()
+ " / " + firstDecision.detail());
assertEquals(1, runner.heldCount(), "sanity: the first roll is held, not run yet");
LeadRollover.PendingRollover second = rollover.open(OTHER_LEAD,
"a second request for the SAME lead, opened from a DIFFERENT terminal");
LeadRollover.RollDecision secondDecision = rollover.confirm(OTHER_LEAD, second.token(), true);
assertFalse(secondDecision.accepted(), "a different terminal resolving to the SAME lead "
+ "name must still be refused — the single-flight claim is keyed on the lead's "
+ "name, not its terminal");
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, secondDecision.reason());
assertTrue(secondDecision.detail().contains(LEAD_NAME), "the refusal detail must name "
+ "the single-flight key (the lead's name): " + secondDecision.detail());
assertTrue(secondDecision.detail().contains(first.token()), "the refusal detail must "
+ "name the token holding the claim: " + secondDecision.detail());
assertEquals(1, runner.heldCount(), "the refused confirm() must never have reached "
+ "continuationRunner — it ran exactly once, for the first roll only");
runner.runNext(); // let the first (and only) held roll finish
// The claim must have been released once the first roll's continuation finished — a
// fresh request for the SAME lead, even opened from yet another terminal, is now
// approved.
LeadRollover.PendingRollover third = rollover.open(OTHER_LEAD, "retry after the first roll finished");
LeadRollover.RollDecision thirdDecision = rollover.confirm(OTHER_LEAD, third.token(), true);
assertTrue(thirdDecision.accepted(), "the claim must have been released once the first "
+ "roll finished: " + thirdDecision.reason() + " / " + thirdDecision.detail());
}
@Test
@DisplayName("[NAME-KEYED 2] the claim is released under the key it was taken under, even "
+ "when the old terminal has already dropped out of the live-lead roster by release "
+ "time — success (non-throwing) path")
void claimReleasedUnderCarriedKeyEvenWhenOldTerminalIsNoLongerInTheRosterSuccessPath() throws IOException {
// herdrReadyForAFullRoll() lets the RETRY below run its continuation to a genuine ROLLED
// completion (pane gone on the first check, a pinned relaunch) — needed because, unlike
// the other NAME-KEYED tests, this one's retry resolves a REAL name ("opus") and must not
// spin forever against this test's fixed, non-advancing clock (see this class's javadoc).
FakeHerdr herdr = herdrReadyForAFullRoll();
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
AgentControl agents = new AgentControl(herdr);
WorkspaceControl spaces = new WorkspaceControl(herdr);
// A config that CAN relaunch 'opus' — unlike emptyFleetConfig(), its fleet() is non-null,
// so resolveLaunchable(null) below returns null cleanly instead of throwing a NullPointer
// out of cfg.fleet() itself; this test needs the NON-throwing relaunch-refused exit, not
// an incidental NPE.
LeadLauncher launcher = fakeLauncher(herdr, fleetConfigWithRelaunchableLead());
Map<String, String> roster = new HashMap<>();
roster.put(LEAD, LEAD_NAME);
Map<String, String> liveLeadTerminals = Map.of(NEW_TERMINAL, LEAD_NAME);
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
_ -> null, roster::get, () -> liveLeadTerminals, fixedClock(clock), () -> { }, Runnable::run);
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
// The live roster drops the OLD terminal before this roll's continuation runs — the
// hazard fleetd #737 names: leadNameForTerminal reads the LIVE roster, and (with the
// synchronous runner this test injects) the continuation runs INSIDE this confirm()
// call, strictly after open() already resolved and carried the key.
roster.remove(LEAD);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
LeadRollover.RollStatus status = rollover.status(pending.token());
assertEquals(LeadRollover.RollState.RELAUNCH_FAILED, status.state(), "sanity: leadName "
+ "resolved to null (the roster had already dropped the old terminal), so "
+ "relaunch(null) fails cleanly, through the non-throwing exit this test means "
+ "to cover: " + status.detail());
// A fresh terminal for the SAME lead name, confirmed live again — proves the claim this
// roll took under the lead's name was actually released, not left stuck under whatever
// a release that RECOMPUTES the key via leadNameForTerminal (now null for the old
// terminal) would have tried to remove instead.
roster.put(OTHER_LEAD, LEAD_NAME);
LeadRollover.PendingRollover retry = rollover.open(OTHER_LEAD, "a fresh pane for the same lead");
LeadRollover.RollDecision retryDecision = rollover.confirm(OTHER_LEAD, retry.token(), true);
assertTrue(retryDecision.accepted(), "the claim taken under the lead's name must have "
+ "been released even though the OLD terminal no longer resolved to any name at "
+ "release time: " + retryDecision.reason() + " / " + retryDecision.detail());
LeadRollover.RollStatus retryStatus = rollover.status(retry.token());
assertEquals(LeadRollover.RollState.ROLLED, retryStatus.state(), "sanity: this retry's "
+ "own claim (also 'opus') must not have been blocked by a leftover claim from "
+ "the first roll: " + retryStatus.detail());
}
@Test
@DisplayName("[NAME-KEYED 3] the claim is released under the key it was taken under on the "
+ "thrown-exception path too, even when the old terminal has already dropped out of "
+ "the live-lead roster")
void claimReleasedUnderCarriedKeyEvenWhenOldTerminalIsNoLongerInTheRosterThrowingPath() throws IOException {
FakeHerdr fake = new FakeHerdr(); // default idle — the turn-settle wait passes immediately
fake.paneCloseFailsWith("permission_denied"); // a real failure, not an already-gone code
Path handover = writeHandover("handover contents");
AtomicLong clock = new AtomicLong(1_000);
AgentControl agents = new AgentControl(fake);
WorkspaceControl spaces = new WorkspaceControl(fake);
LeadLauncher launcher = fakeLauncher(fake, emptyFleetConfig());
Map<String, String> roster = new HashMap<>();
roster.put(LEAD, LEAD_NAME);
LeadRollover rollover = new LeadRollover(agents, spaces, launcher, () -> cfg(handover.toString()),
_ -> null, roster::get, Map::of, fixedClock(clock), () -> { }, Runnable::run);
LeadRollover.PendingRollover first = rollover.open(LEAD, "context is full");
// Same hazard as [NAME-KEYED 2] above, now exercised on the thrown-exception exit.
roster.remove(LEAD);
LeadRollover.RollDecision decision = rollover.confirm(LEAD, first.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(first.token());
assertEquals(LeadRollover.RollState.FAILED, status.state(), "sanity: the continuation "
+ "threw and left a terminal FAILED outcome: " + status.detail());
// A fresh terminal for the SAME lead name, confirmed live again — proves the claim this
// roll took under the lead's name was actually released, not left stuck under whatever
// a release that RECOMPUTES the key via leadNameForTerminal (now null for the old
// terminal) would have tried to remove instead.
roster.put(OTHER_LEAD, LEAD_NAME);
LeadRollover.PendingRollover retry = rollover.open(OTHER_LEAD, "retry after the throw");
LeadRollover.RollDecision retryDecision = rollover.confirm(OTHER_LEAD, retry.token(), true);
assertTrue(retryDecision.accepted(), "the claim must have been released even though the "
+ "continuation threw AND the old terminal no longer resolved to any name at "
+ "release time: " + retryDecision.reason() + " / " + retryDecision.detail());
}
}