Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 8a1d73b39e |
@@ -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());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user