Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 428a12af62 | |||
| a332dfdb2c | |||
| d0f4ae057b | |||
| 7b3beaa209 | |||
| 8bb2aa0be4 | |||
| 4ce3149bfd | |||
| 1fc9e85bf1 |
@@ -164,7 +164,13 @@ public final class LeadRollover {
|
||||
* The handover file's modified time is not after {@link #open}'s request timestamp, or is
|
||||
* older than {@code maxDocAgeSeconds}.
|
||||
*/
|
||||
HANDOVER_STALE
|
||||
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.
|
||||
*/
|
||||
ROLL_ALREADY_RUNNING
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -305,6 +311,15 @@ 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
|
||||
* 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.
|
||||
*/
|
||||
private final Map<String, String> rollingByTerminal = 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
|
||||
@@ -487,6 +502,16 @@ public final class LeadRollover {
|
||||
return docCheck;
|
||||
}
|
||||
|
||||
// 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);
|
||||
if (holder != null) {
|
||||
return RollDecision.refused(RefusalReason.ROLL_ALREADY_RUNNING,
|
||||
"lead terminal " + p.leadTerminal() + " already has a roll running under token "
|
||||
+ holder);
|
||||
}
|
||||
|
||||
// Record IN_PROGRESS BEFORE removing from `pending` — see RollState#IN_PROGRESS and
|
||||
// OUTCOME_HISTORY_CAP's javadoc. This ordering means `token` is written into `outcomes`
|
||||
// while it is STILL present in `pending`; status() checks `outcomes` first (see that
|
||||
@@ -499,7 +524,23 @@ public final class LeadRollover {
|
||||
pending.remove(token);
|
||||
log.info("lead-rollover: confirmed token={} lead={} — roll scheduled once the calling turn ends",
|
||||
token, callerTerminal);
|
||||
continuationRunner.accept(() -> runRollover(p, cfg));
|
||||
try {
|
||||
continuationRunner.accept(() -> runRollover(p, cfg));
|
||||
} 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.
|
||||
// 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.
|
||||
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);
|
||||
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"));
|
||||
}
|
||||
return RollDecision.approved();
|
||||
}
|
||||
|
||||
@@ -539,6 +580,12 @@ public final class LeadRollover {
|
||||
"the roll's continuation threw " + e.toString() + " — the roll is dead and will "
|
||||
+ "not retry itself; check the daemon log for the stack trace, then open() "
|
||||
+ "a fresh rollover request"));
|
||||
} 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());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
@@ -101,6 +102,14 @@ public final class Rendezvous {
|
||||
/** Reverse rendezvous (CB-205): worker questions awaiting the primary's answer, keyed by {@code turnId}. */
|
||||
private final ConcurrentHashMap<String, AskWaiter> asks = new ConcurrentHashMap<>();
|
||||
private final AtomicLong askSeq = new AtomicLong();
|
||||
/**
|
||||
* Minted once per {@code Rendezvous} instance and folded into every {@code turnId} (see
|
||||
* {@link #openAsk(String)}). {@link #askSeq} alone restarts at zero for every instance, so
|
||||
* without this a {@code turnId} minted by one instance could be minted again by another and
|
||||
* resolve to an unrelated ask with no error; this nonce makes that impossible, because an id
|
||||
* minted by one instance can never match the id space of another.
|
||||
*/
|
||||
private final String askBootNonce = UUID.randomUUID().toString().substring(0, 6);
|
||||
/** Per-session index of the currently-open ask, so duplicate fleet_ask calls coalesce onto one turn. */
|
||||
private final ConcurrentHashMap<String, String> openAsksBySession = new ConcurrentHashMap<>();
|
||||
|
||||
@@ -187,7 +196,7 @@ public final class Rendezvous {
|
||||
while (true) {
|
||||
AskWaiter[] minted = { null };
|
||||
String turnId = openAsksBySession.computeIfAbsent(session, _ -> {
|
||||
String newTurnId = session + "#" + askSeq.incrementAndGet();
|
||||
String newTurnId = session + "#" + askBootNonce + "-" + askSeq.incrementAndGet();
|
||||
CompletableFuture<String> answer = new CompletableFuture<>();
|
||||
AskWaiter waiter = new AskWaiter(session, answer, ownerOf(session));
|
||||
asks.put(newTurnId, waiter);
|
||||
|
||||
@@ -255,6 +255,10 @@ public final class SessionManager implements TurnListener {
|
||||
handle.id(), handle.terminalId(), resolvedProfile, actualRole, cwd, ownerTerminal, now, now, 0,
|
||||
MemberSession.State.SPAWNING, null, null, handle.charterReceipt(), handle.agentSessionId());
|
||||
registry.put(handle.id(), session);
|
||||
// A presence contact that already arrived for this terminal found no registry
|
||||
// entry to transition and gave up silently. Retry it now that one exists; remove
|
||||
// this call and such a session stays in SPAWNING even though it is present.
|
||||
reconcilePresence(handle.terminalId());
|
||||
handles.put(handle.id(), handle);
|
||||
log.debug("acquired session id={} terminal={} profile={} owner={}",
|
||||
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
|
||||
@@ -809,6 +813,10 @@ public final class SessionManager implements TurnListener {
|
||||
handle.charterReceipt(),
|
||||
handle.agentSessionId());
|
||||
registry.put(handle.id(), session);
|
||||
// A presence contact that already arrived for this terminal found no registry entry to
|
||||
// transition and gave up silently. Retry it now that one exists; remove this call and
|
||||
// such a session stays in SPAWNING even though it is present.
|
||||
reconcilePresence(handle.terminalId());
|
||||
handles.put(handle.id(), handle);
|
||||
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
|
||||
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
|
||||
@@ -983,6 +991,20 @@ public final class SessionManager implements TurnListener {
|
||||
transitionByTerminal(terminalId, MemberSession.State.SPAWNING, MemberSession.State.READY);
|
||||
}
|
||||
|
||||
/**
|
||||
* Completes a newly registered session's {@code SPAWNING -> READY} transition when {@code
|
||||
* terminalId} was already marked present before this ran. A terminal never marked present is
|
||||
* left in {@code SPAWNING}; it reaches {@code READY} normally through {@link #onReady} once
|
||||
* its own contact arrives. Callers must run this only once the session's registry entry is
|
||||
* already visible — {@link #onReady}'s transition matches against that entry, and reconciling
|
||||
* before the entry exists finds nothing to transition.
|
||||
*/
|
||||
private void reconcilePresence(String terminalId) {
|
||||
if (terminalId != null && !terminalId.isBlank() && presence.isPresent(terminalId)) {
|
||||
onReady(terminalId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Lifecycle hook: a message was delivered into the worker — it is now busy on a turn.
|
||||
* The turn count is bumped and the activity timestamp is refreshed. A {@code DONE} session
|
||||
|
||||
@@ -1305,8 +1305,12 @@ class LeadRolloverTest {
|
||||
int rolls = LeadRollover.OUTCOME_HISTORY_CAP + 50;
|
||||
String[] tokens = new String[rolls];
|
||||
for (int i = 0; i < rolls; i++) {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full #" + i);
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
// a distinct terminal per iteration: confirm() single-flights per lead terminal, so
|
||||
// reusing LEAD here would refuse every roll after the first instead of building up
|
||||
// OUTCOME_HISTORY_CAP+50 simultaneously in-flight entries.
|
||||
String terminal = "term_cap_" + i;
|
||||
LeadRollover.PendingRollover pending = rollover.open(terminal, "context is full #" + i);
|
||||
LeadRollover.RollDecision decision = rollover.confirm(terminal, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "roll #" + i + " should have been approved: " + decision.detail());
|
||||
tokens[i] = pending.token();
|
||||
}
|
||||
@@ -1447,4 +1451,195 @@ class LeadRolloverTest {
|
||||
assertTrue(status.detail().contains("HerdrException"), "the detail must name the exception "
|
||||
+ "so an operator reading status() has something to act on: " + status.detail());
|
||||
}
|
||||
|
||||
// ---- fleetd #726 unit 3: confirm() single-flights per lead terminal -----------------------
|
||||
|
||||
@Test
|
||||
@DisplayName("[SINGLE-FLIGHT 1] a second confirm() for the SAME lead terminal is refused with "
|
||||
+ "ROLL_ALREADY_RUNNING while the first roll is still in flight, continuationRunner ran "
|
||||
+ "exactly once, and the claim is released once that first roll finishes")
|
||||
void secondConfirmForSameTerminalIsRefusedWhileFirstRollIsStillInFlight() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr(); // default idle — the held roll WOULD complete once run
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
HoldingRunner runner = new HoldingRunner();
|
||||
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
|
||||
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(LEAD, "a second request for the same terminal");
|
||||
LeadRollover.RollDecision secondDecision = rollover.confirm(LEAD, second.token(), true);
|
||||
|
||||
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(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
|
||||
assertEquals(2, promptCallCount(herdr), "continuationRunner ran exactly once: /clear then "
|
||||
+ "bootstrapText, for the first roll only");
|
||||
|
||||
// The claim must have been released once the first roll's continuation finished — a
|
||||
// fresh request for the SAME terminal is now approved.
|
||||
LeadRollover.PendingRollover third = rollover.open(LEAD, "retry after the first roll finished");
|
||||
LeadRollover.RollDecision thirdDecision = rollover.confirm(LEAD, third.token(), true);
|
||||
assertTrue(thirdDecision.accepted(), "the claim must have been released once the first "
|
||||
+ "roll finished: " + thirdDecision.reason() + " / " + thirdDecision.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SINGLE-FLIGHT 2] two DIFFERENT lead terminals confirm without either being "
|
||||
+ "refused, and the continuation runs once per terminal")
|
||||
void differentLeadTerminalsConfirmIndependently() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr(); // default idle — both rolls complete
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, cfg(handover.toString()), fixedClock(clock));
|
||||
|
||||
LeadRollover.PendingRollover a = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.PendingRollover b = rollover.open(OTHER_LEAD, "context is full too");
|
||||
|
||||
LeadRollover.RollDecision aDecision = rollover.confirm(LEAD, a.token(), true);
|
||||
LeadRollover.RollDecision bDecision = rollover.confirm(OTHER_LEAD, b.token(), true);
|
||||
|
||||
assertTrue(aDecision.accepted(), "expected approval; got: " + aDecision.reason() + " / " + aDecision.detail());
|
||||
assertTrue(bDecision.accepted(), "expected approval; got: " + bDecision.reason() + " / " + bDecision.detail());
|
||||
assertEquals(4, promptCallCount(herdr), "both rolls ran to completion — /clear and "
|
||||
+ "bootstrapText, once per terminal");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SINGLE-FLIGHT 3] a confirm() refused for OPERATOR_NOT_CONFIRMED does not take "
|
||||
+ "the single-flight claim — a later valid confirm on the same terminal is still approved")
|
||||
void refusedConfirmDoesNotTakeTheClaim() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, cfg(handover.toString(), true), fixedClock(clock));
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision refused = rollover.confirm(LEAD, pending.token(), false);
|
||||
assertFalse(refused.accepted());
|
||||
assertEquals(LeadRollover.RefusalReason.OPERATOR_NOT_CONFIRMED, refused.reason());
|
||||
|
||||
// the same token stays pending after a gate refusal — retrying it with operator
|
||||
// confirmation must succeed, proving the earlier refusal never took the claim.
|
||||
LeadRollover.RollDecision approved = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(approved.accepted(), "a confirm() refused for an existing gate must never have "
|
||||
+ "taken the single-flight claim: " + approved.reason() + " / " + approved.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SINGLE-FLIGHT 4] a roll whose continuation throws a RuntimeException still "
|
||||
+ "releases the claim — a fresh confirm() on that terminal is approved afterwards")
|
||||
void throwingContinuationStillReleasesTheClaim() throws IOException {
|
||||
FakeHerdr fake = new FakeHerdr(); // default idle — the turn-settle wait passes immediately
|
||||
HerdrClient throwsOnClear = new HerdrClient() {
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) throws HerdrException {
|
||||
if ("agent.prompt".equals(method) && String.valueOf(params).contains("/clear")) {
|
||||
throw new HerdrException("simulated herdr transport failure sending /clear");
|
||||
}
|
||||
return fake.call(method, params);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
fake.close();
|
||||
}
|
||||
};
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(throwsOnClear, cfg(handover.toString()), fixedClock(clock));
|
||||
|
||||
LeadRollover.PendingRollover first = rollover.open(LEAD, "context is full");
|
||||
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());
|
||||
|
||||
LeadRollover.PendingRollover second = rollover.open(LEAD, "retry after the failure");
|
||||
LeadRollover.RollDecision retryDecision = rollover.confirm(LEAD, second.token(), true);
|
||||
assertTrue(retryDecision.accepted(), "the claim must have been released even though the "
|
||||
+ "continuation threw — a release only on the success path would leave this lead "
|
||||
+ "terminal unrollable forever after one failure: " + retryDecision.reason() + " / "
|
||||
+ retryDecision.detail());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SINGLE-FLIGHT 5] cancel(token) after a successful confirm(token) returns false "
|
||||
+ "and does not release the claim held by that roll")
|
||||
void cancelAfterConfirmDoesNotReleaseTheClaim() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
HoldingRunner runner = new HoldingRunner();
|
||||
LeadRollover rollover = newRolloverWithHoldingRunner(herdr, cfg(handover.toString()),
|
||||
fixedClock(clock), runner);
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason() + " / " + decision.detail());
|
||||
|
||||
assertFalse(rollover.cancel(pending.token()), "confirm() already removed the token from "
|
||||
+ "pending, so cancel() must find nothing for it");
|
||||
|
||||
LeadRollover.PendingRollover another = rollover.open(LEAD,
|
||||
"a second request while the first roll is still held");
|
||||
LeadRollover.RollDecision anotherDecision = rollover.confirm(LEAD, another.token(), true);
|
||||
assertFalse(anotherDecision.accepted(), "cancel() on the already-confirmed token must not "
|
||||
+ "have released the claim held by that roll");
|
||||
assertEquals(LeadRollover.RefusalReason.ROLL_ALREADY_RUNNING, anotherDecision.reason());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SINGLE-FLIGHT 6] a continuationRunner that rejects the hand-off still releases "
|
||||
+ "the claim and leaves a terminal FAILED outcome, not a stuck IN_PROGRESS")
|
||||
void continuationRunnerThatRejectsTheHandOffStillReleasesTheClaim() throws IOException {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path handover = writeHandover("handover contents");
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
// A continuationRunner that rejects every hand-off, standing in for e.g. a bounded
|
||||
// executor's RejectedExecutionException — the continuation never runs, so runRollover's
|
||||
// own finally never gets a chance to release the claim either.
|
||||
LeadRollover rollover = new LeadRollover(agents, () -> cfg(handover.toString()), _ -> null,
|
||||
fixedClock(clock), () -> { },
|
||||
r -> { throw new RuntimeException("simulated continuationRunner rejection"); });
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
assertTrue(decision.accepted(), "every gate passed before the hand-off itself rejected — "
|
||||
+ "confirm() treats that rejection the same way it treats a throw INSIDE the "
|
||||
+ "continuation: logged and surfaced only through status(), never as a refusal here: "
|
||||
+ decision.reason() + " / " + decision.detail());
|
||||
|
||||
LeadRollover.RollStatus status = rollover.status(pending.token());
|
||||
assertEquals(LeadRollover.RollState.FAILED, status.state(), "a rejected hand-off must leave "
|
||||
+ "a TERMINAL outcome, never a stuck IN_PROGRESS — status() would otherwise have no "
|
||||
+ "way to tell a roll that never started from one still genuinely running: got "
|
||||
+ status.state() + " / " + status.detail());
|
||||
assertNotEquals(LeadRollover.RollState.IN_PROGRESS, status.state());
|
||||
|
||||
LeadRollover.PendingRollover retry = rollover.open(LEAD, "retry after the rejected hand-off");
|
||||
LeadRollover.RollDecision retryDecision = rollover.confirm(LEAD, retry.token(), true);
|
||||
assertTrue(retryDecision.accepted(), "the claim must have been released even though "
|
||||
+ "continuationRunner itself threw before the continuation ever ran — otherwise "
|
||||
+ "this lead terminal could never be rolled again: " + retryDecision.reason() + " / "
|
||||
+ retryDecision.detail());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -191,6 +191,53 @@ class RendezvousTest {
|
||||
assertNull(rendezvous.askOwner(t.turnId()), "a closed ask no longer reports an owner");
|
||||
}
|
||||
|
||||
// ── fleetd #729: per-boot nonce guards turnId against cross-instance reuse ────────────────
|
||||
|
||||
@Test
|
||||
void twoInstancesMintDisjointTurnIds() {
|
||||
Rendezvous other = new Rendezvous();
|
||||
Rendezvous.AskTicket fromThis = rendezvous.openAsk(W);
|
||||
Rendezvous.AskTicket fromOther = other.openAsk(W);
|
||||
assertNotEquals(fromThis.turnId(), fromOther.turnId(),
|
||||
"each instance mints its own id space, so even a first ask from each must differ");
|
||||
}
|
||||
|
||||
@Test
|
||||
void foreignInstanceTurnIdDoesNotResolve() {
|
||||
Rendezvous other = new Rendezvous();
|
||||
|
||||
// `other` must reach the same sequence number as `rendezvous` (two asks each, the first
|
||||
// closed so the second mints fresh), or this test passes against an empty map instead of
|
||||
// against a colliding id.
|
||||
Rendezvous.AskTicket firstFromThis = rendezvous.openAsk(W);
|
||||
rendezvous.closeAsk(firstFromThis.turnId());
|
||||
Rendezvous.AskTicket secondFromThis = rendezvous.openAsk(W);
|
||||
|
||||
Rendezvous.AskTicket firstFromOther = other.openAsk(W);
|
||||
other.closeAsk(firstFromOther.turnId());
|
||||
other.openAsk(W);
|
||||
|
||||
// control: the id resolves in the instance that minted it, so a false below cannot be
|
||||
// explained by broken plumbing — only by the turnId being foreign to `other`.
|
||||
assertTrue(rendezvous.answerAsk(secondFromThis.turnId(), "answer from this instance"),
|
||||
"the minting instance must still resolve its own turnId");
|
||||
assertFalse(other.answerAsk(secondFromThis.turnId(), "answer from other instance"),
|
||||
"a turnId minted by a different instance must not resolve here");
|
||||
}
|
||||
|
||||
@Test
|
||||
void openAskStillCoalescesDuplicatesAndStillMintsDistinctIdsPerAsk() {
|
||||
Rendezvous.AskTicket t1 = rendezvous.openAsk(W);
|
||||
Rendezvous.AskTicket t2 = rendezvous.openAsk(W);
|
||||
assertEquals(t1.turnId(), t2.turnId(),
|
||||
"a second openAsk while one is open still coalesces onto the same turn");
|
||||
assertFalse(t2.fresh(), "the coalesced ask is still reported as not fresh");
|
||||
|
||||
rendezvous.closeAsk(t1.turnId());
|
||||
Rendezvous.AskTicket t3 = rendezvous.openAsk(W);
|
||||
assertNotEquals(t1.turnId(), t3.turnId(), "two asks from the same session still get different turnIds");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ownerPermitsIsFailClosedOnARecordAndThreeStatesAreDistinct() {
|
||||
assertFalse(Rendezvous.Owner.permits(null, null),
|
||||
|
||||
@@ -0,0 +1,92 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
import dev.ltms.fleet.inject.MemberPresence;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.PeerHandle;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
|
||||
/**
|
||||
* {@link PeerLauncher} decorator that marks presence for a spawned terminal before returning its
|
||||
* handle to the caller — the contact-then-register ordering fleetd #722 covers, where the
|
||||
* terminal's MCP contact lands before {@link SessionManager#acquire} runs its own
|
||||
* {@code registry.put}. The presence view is set after construction, via {@link #presence},
|
||||
* because it is owned by the {@link SessionManager} this launcher is passed into.
|
||||
*/
|
||||
final class PresenceRacingLauncher implements PeerLauncher {
|
||||
|
||||
private final PeerLauncher delegate;
|
||||
volatile MemberPresence presence;
|
||||
|
||||
PresenceRacingLauncher(PeerLauncher delegate) {
|
||||
this.delegate = delegate;
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
PeerHandle handle = delegate.spawn(req);
|
||||
presence.markPresent(handle.terminalId());
|
||||
return handle;
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
PeerHandle handle = delegate.spawn(req, decision);
|
||||
presence.markPresent(handle.terminalId());
|
||||
return handle;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
return delegate.capabilities();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilitiesFor(String profileName) {
|
||||
return delegate.capabilitiesFor(profileName);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return delegate.profiles();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String defaultProfile() {
|
||||
return delegate.defaultProfile();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
return delegate.effectiveCwd(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> parityOverlay(String profileName) {
|
||||
return delegate.parityOverlay(profileName);
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<?> list() {
|
||||
return delegate.list();
|
||||
}
|
||||
|
||||
@Override
|
||||
public int reapOrphanWorkers() {
|
||||
return delegate.reapOrphanWorkers();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
delegate.stop(id);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean clearContext(String id) {
|
||||
return delegate.clearContext(id);
|
||||
}
|
||||
}
|
||||
@@ -422,6 +422,53 @@ class SessionManagerTest {
|
||||
"turn completion moves BUSY → DONE");
|
||||
}
|
||||
|
||||
// --- fleetd #722: registration and presence must reach READY whichever lands first --------
|
||||
|
||||
@Test
|
||||
void registerThenContactReachesReadyForPlainSpawn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
|
||||
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"a presence contact that arrives after registration reaches READY");
|
||||
}
|
||||
|
||||
@Test
|
||||
void contactThenRegisterStillReachesReadyForPlainSpawn() {
|
||||
// The racing launcher marks presence for the spawned terminal from inside spawn() —
|
||||
// before SessionManager.acquire's own registry.put runs — modeling an MCP contact that
|
||||
// lands in that window.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
PresenceRacingLauncher race = new PresenceRacingLauncher(workers);
|
||||
SessionManager sessions = new SessionManager(race);
|
||||
race.presence = sessions.asPresence();
|
||||
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
|
||||
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"a presence contact that lands before registry.put must still reach READY");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aTerminalNeverMarkedPresentStaysSpawningAfterRegistration() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
|
||||
assertEquals(MemberSession.State.SPAWNING, sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"registration alone must not advance a terminal that was never marked present");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseTearsDownWorkerAndRemovesFromRosterAndIsIdempotent() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -159,6 +159,40 @@ class WorktreeSessionManagerTest {
|
||||
assertEquals(expectedPath, s.cwd(), "session cwd is the worktree path");
|
||||
}
|
||||
|
||||
// --- fleetd #722: registration and presence must reach READY whichever lands first --------
|
||||
|
||||
@Test
|
||||
void registerThenContactReachesReadyForWorktreeSpawn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
|
||||
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
|
||||
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
|
||||
new WorktreeRequest("cb-722", null));
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
|
||||
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"a presence contact that arrives after worktree registration reaches READY");
|
||||
}
|
||||
|
||||
@Test
|
||||
void contactThenRegisterStillReachesReadyForWorktreeSpawn() {
|
||||
// The racing launcher marks presence for the spawned terminal from inside spawn() —
|
||||
// before SessionManager.acquireWithWorktree's own registry.put runs — modeling an MCP
|
||||
// contact that lands in that window.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
|
||||
PresenceRacingLauncher race = new PresenceRacingLauncher(workerService(herdr));
|
||||
SessionManager sessions = new SessionManager(race, worktrees);
|
||||
race.presence = sessions.asPresence();
|
||||
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
|
||||
new WorktreeRequest("cb-722", null));
|
||||
|
||||
assertEquals(MemberSession.State.READY, sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"a presence contact that lands before worktree registration must still reach READY");
|
||||
}
|
||||
|
||||
@Test
|
||||
void worktreeArchitectAcquireAlsoBindsItsSlot() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
Reference in New Issue
Block a user