Compare commits

...

1 Commits

Author SHA1 Message Date
Dai Ha 8a1d73b39e fleetd #737 unit 4: key the rollover single-flight claim on the lead's name
CI / shell-tests (pull_request) Failing after 7s
CI / contract (pull_request) Successful in 52s
CI / build (pull_request) Failing after 1m59s
rollingByTerminal keyed the single-flight lock on p.leadTerminal(), the pane
address. A roll replaces the pane, so a second roll of the same lead opened
from the new terminal landed on a different map key and could run concurrent
with the first roll's still-in-flight continuation.

PendingRollover now carries rolloverKey, resolved once in open() from
leadNameForTerminal while the lead is certainly still live, falling back to
the terminal itself when the name resolves null or blank. confirm()'s claim
and both release sites (the continuationRunner-rejection catch and
runRollover's finally) use the carried key instead of recomputing it, since
leadNameForTerminal no longer resolves the old terminal by release time.
NOT_YOUR_ROLLOVER and the pane-teardown calls stay keyed on p.leadTerminal(),
unchanged.

rollingByTerminal is renamed rollingByLead to match.
2026-10-05 05:57:59 +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());
}
}