Compare commits

...

8 Commits

Author SHA1 Message Date
Dai Ha 428a12af62 fleetd #722: reconcile presence that arrives before a session's registry entry
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Failing after 1m48s
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 47s
CI / build (push) Failing after 1m46s
A member whose MCP contact lands between launcher.spawn() and registry.put()
had its presence marked, but the SPAWNING -> READY transition that markPresent
triggers found no registry entry yet and silently did nothing. The mark then
persisted while registration left the session in SPAWNING, with nothing to
retry the transition. That left the session undeliverable to reclaim/seat
accounting even though it was present and deliverable.

Add SessionManager.reconcilePresence, called right after registry.put in both
the plain-spawn and worktree-spawn paths, to retry the transition for a
terminal already marked present. One private helper serves both call sites.

Tests cover both orderings (contact-then-register and register-then-contact)
for both spawn paths, plus a terminal never marked present staying in
SPAWNING. The contact-then-register tests use a new PresenceRacingLauncher
test double that marks presence from inside spawn(), before acquire()'s own
registry.put runs.
2026-10-04 19:23:02 +02:00
Dai Ha a332dfdb2c Merge remote-tracking branch 'origin/worker/726-ea34a0-2'
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 54s
CI / build (push) Failing after 1m53s
2026-10-04 19:09:03 +02:00
Dai Ha d0f4ae057b fleetd #726 unit 3 fix: release the single-flight claim when continuationRunner rejects the hand-off
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Failing after 1m51s
confirm() takes the per-lead-terminal claim before handing the roll to
continuationRunner, and the only release path was runRollover's own
finally. If continuationRunner.accept itself throws, runRollover never
starts, so that finally never runs, and nothing else ever writes
rollingByTerminal — the claim is held forever and the terminal can never
be rolled again. This differs from fleetd #615, which covers a throw
INSIDE the continuation (runRollover already catches that and still
releases the claim) — this is a throw from the hand-off itself, which
is not reachable with today's virtual-thread runner but would be with
a bounded executor's RejectedExecutionException.

confirm() now catches that throw, releases the claim, and overwrites
the IN_PROGRESS outcome with a terminal FAILED one, matching how a
throw inside the continuation is already surfaced.
2026-10-04 19:06:57 +02:00
Dai Ha 7b3beaa209 Merge remote-tracking branch 'origin/worker/726-10cbf0-1'
CI / shell-tests (push) Failing after 10s
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m53s
2026-10-04 19:06:47 +02:00
Dai Ha b3b2bf3da6 fleetd #726 unit 1 review fixes: correct relaunch's javadoc and dedupe its resolve logic
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 57s
CI / build (pull_request) Failing after 1m49s
relaunch's javadoc said it reads the live config; it actually reads the
FleetConfig snapshot this launcher was constructed with (fleet.leaders
is the frozen half), so say that and note a profile/tab edit needs a
daemon restart.

relaunch's recognise-only refusal reused ensureLeads()'s log wording,
which claims the lead 'is not live' — true in ensureLeads()'s context
(reached only after a short live count), false in relaunch's (which
never counts liveness, by design). Dropped that clause.

Pulled the declared/creatable/profile-configured resolution shared by
ensureLeads() and relaunch() into one private resolveLaunchable(name)
helper (returns a new ResolvedLead(lead, profile) record, or null
having logged), so the three refusals and their wording live in one
place instead of two copies that can drift. Behaviour-preserving:
ensureLeads() keeps its own liveness-count logic around the shared
resolve, and the existing 37 LeadLauncherTest cases are unchanged and
still pass.
2026-10-04 18:59:56 +02:00
Dai Ha 8bb2aa0be4 Merge remote-tracking branch 'origin/worker/729-5961c6-3'
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 57s
CI / build (push) Failing after 2m0s
2026-10-04 18:59:54 +02:00
Dai Ha 4ce3149bfd fleetd #726 unit 3: make LeadRollover.confirm() single-flight per lead terminal
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Failing after 1m48s
Two open() calls for the same lead terminal minted two tokens that both
passed confirm()'s ownership check, so both could reach the deferred
continuation and roll the same lead twice. confirm() now claims a
per-lead-terminal slot (an atomic put-if-absent) once every other gate has
passed, refusing a concurrent confirm with the new ROLL_ALREADY_RUNNING
reason; runRollover releases the claim in a finally, on both the success
and the thrown-exception path.
2026-10-04 18:53:53 +02:00
Dai Ha 38544d467c fleetd #726 unit 1: give LeadLauncher a public single-lead relaunch seam
CI / shell-tests (pull_request) Failing after 8s
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Failing after 1m48s
Adds LeadLauncher.relaunch(name), which starts exactly the named lead
from the live config, outside of ensureLeads()'s instances bookkeeping.
It retries the whole launch attempt (not just the agent_name_taken/
agent_pane_busy cases ResilientAgentLaunch already retries inside one
agents.start call) up to RELAUNCH_ATTEMPTS times.

launch() now returns the started Agent (null on failure) instead of a
boolean, so relaunch() and ensureLeads() share the same primitive.
2026-10-04 18:52:38 +02:00
8 changed files with 660 additions and 21 deletions
@@ -83,6 +83,9 @@ public final class LeadLauncher {
private static final Logger log = LoggerFactory.getLogger(LeadLauncher.class);
/** Attempts {@link #relaunch(String)} makes before giving up and returning {@code null}. */
static final int RELAUNCH_ATTEMPTS = 3;
private final AgentControl agents;
private final WorkspaceControl spaces;
private final FleetConfig cfg;
@@ -201,24 +204,14 @@ public final class LeadLauncher {
log.info("lead '{}': {} live, {} wanted — nothing to start", name, running, wanted);
continue;
}
if (!lead.isCreatable()) {
// A lead with a `tab:` but no `profile:` is recognise-only by design: the operator
// opens it by hand. Say so once rather than looking like a silent failure.
log.info("lead '{}' is not live, and names no profile — it can be recognised but not "
+ "launched. Add `profile:` under fleet.leaders.{} to have fleetd start it.",
name, name);
continue;
}
FleetConfig.Profile profile = cfg.profiles().get(lead.profile());
if (profile == null) {
log.warn("lead '{}' names profile '{}', which is not configured — not launching",
name, lead.profile());
ResolvedLead resolved = resolveLaunchable(name);
if (resolved == null) {
continue;
}
for (int i = running; i < wanted; i++) {
if (launch(name, lead, profile)) {
if (launch(name, resolved.lead(), resolved.profile()) != null) {
started++;
}
}
@@ -226,6 +219,82 @@ public final class LeadLauncher {
return started;
}
/** A declared lead paired with the profile it launches on — {@link #resolveLaunchable}'s result. */
private record ResolvedLead(FleetConfig.Leader lead, FleetConfig.Profile profile) {
}
/**
* The declared {@code Leader} and its {@code Profile} for {@code name}, read from the config
* snapshot this launcher was constructed with.
*
* @return the resolved pair, or {@code null} (having logged) if {@code name} is not declared
* under {@code fleet.leaders}, that lead names no {@code profile:} (a {@code tab:}-only,
* recognise-only lead), or its {@code profile:} is not configured. Shared by
* {@link #ensureLeads()} and {@link #relaunch(String)} so the three refusals and their
* wording live in one place.
*/
private ResolvedLead resolveLaunchable(String name) {
FleetConfig.Leader lead = cfg.fleet().leaders().get(name);
if (lead == null) {
log.warn("lead '{}' is not declared under fleet.leaders — not launching", name);
return null;
}
if (!lead.isCreatable()) {
// A lead with a `tab:` but no `profile:` is recognise-only by design: the operator
// opens it by hand. Say so once rather than looking like a silent failure.
log.info("lead '{}' names no profile — it can be recognised but not launched. Add "
+ "`profile:` under fleet.leaders.{} to have fleetd start it.", name, name);
return null;
}
FleetConfig.Profile profile = cfg.profiles().get(lead.profile());
if (profile == null) {
log.warn("lead '{}' names profile '{}', which is not configured — not launching",
name, lead.profile());
return null;
}
return new ResolvedLead(lead, profile);
}
/**
* Start the named lead from the config snapshot this launcher was constructed with — not a
* live read, so a lead's {@code profile:} or {@code tab:} edited in config needs a daemon
* restart to take effect here — outside of {@link #ensureLeads()}'s {@code instances}
* bookkeeping.
*
* @return the started {@link Agent}, or {@code null} if {@code name} is not declared under
* {@code fleet.leaders}, that lead names no {@code profile:} (a {@code tab:}-only,
* recognise-only lead), its {@code profile:} is not configured, or every attempt up to
* {@link #RELAUNCH_ATTEMPTS} failed to start it. Never throws.
*
* <p>Does not count how many instances of this lead are already live. {@link #ensureLeads()}'s
* count exists to avoid starting a second orchestrator; the caller of this method has already
* decided to replace the lead and owns that decision.
*
* <p>Retries the whole launch attempt — not only the {@code agent_name_taken}/
* {@code agent_pane_busy} cases {@link ResilientAgentLaunch} already retries inside one
* {@code agents.start} call — up to {@link #RELAUNCH_ATTEMPTS} times, sleeping via the
* injected sleeper between attempts, and returns the agent from the first attempt that
* succeeds.
*/
public Agent relaunch(String name) {
ResolvedLead resolved = resolveLaunchable(name);
if (resolved == null) {
return null;
}
for (int attempt = 1; attempt <= RELAUNCH_ATTEMPTS; attempt++) {
Agent started = launch(name, resolved.lead(), resolved.profile());
if (started != null) {
return started;
}
if (attempt < RELAUNCH_ATTEMPTS) {
sleeper.run();
}
}
return null;
}
/**
* How many live leads exist per configured name, and which of that name's labelled tabs are
* <em>not</em> live: a running agent in a tab labelled with that lead's exact {@code tab}
@@ -335,7 +404,7 @@ public final class LeadLauncher {
}
/**
* Start one lead. Returns false (having logged) rather than throwing on any failure.
* Start one lead. Returns null (having logged) rather than throwing on any failure.
*
* <p>Goes through the same {@link ResilientAgentLaunch} seam every member spawn uses
* (fleetd #727): the assembled argv is refused outright if it cannot fit the pane line herdr
@@ -344,7 +413,7 @@ public final class LeadLauncher {
* relaunch, and a seed pane whose shell has not reached its prompt yet ({@code
* agent_pane_busy}) is retried rather than failing on the first miss.
*/
private boolean launch(String name, FleetConfig.Leader lead, FleetConfig.Profile profile) {
private Agent launch(String name, FleetConfig.Leader lead, FleetConfig.Profile profile) {
String label = lead.tabLabel();
String cwd = (lead.cwd() == null || lead.cwd().isBlank())
? System.getProperty("user.dir") : lead.cwd();
@@ -375,7 +444,7 @@ public final class LeadLauncher {
log.info("lead '{}' launched: profile={} tab={} pane={} terminal={} label='{}' cwd={}",
name, profile.profile(), tab.tab().tabId(), started.paneId(),
started.terminalId(), label, cwd);
return true;
return started;
} catch (RuntimeException e) {
log.warn("lead '{}' failed to launch on profile '{}': {}",
name, profile.profile(), e.getMessage());
@@ -387,7 +456,7 @@ public final class LeadLauncher {
tab.tab().tabId(), cleanup.getMessage());
}
}
return false;
return null;
}
}
@@ -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());
}
}
@@ -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
@@ -571,4 +571,137 @@ class LeadLauncherTest {
assertTrue(warn.contains("opus"), "names the profile: " + warn);
assertTrue(warn.contains("x".repeat(60)), "names the culprit argument: " + warn);
}
// ── fleetd #726 unit 1: the single-lead relaunch seam ─────────────────────────────────────
/**
* The returned agent's {@code terminalId()}/{@code paneId()} are the ones the fake
* {@code AgentControl} actually started — not a coincidental field left over from the caller.
* {@code paneId()} echoes the exact {@code pane_id} the launch's own {@code agent.start} call
* carried (protocol 19: the agent starts into the pane it is asked to), and {@code
* terminalId()} is herdr's own generated id, which the fake always shapes as {@code
* term_new_<n>}.
*/
@Test
void relaunchReturnsTheStartedAgent() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNotNull(started, "a launchable, configured lead must start");
Object startedPaneIdParam = ((Map<?, ?>) herdr.lastCall("agent.start").params()).get("pane_id");
assertEquals(startedPaneIdParam, started.paneId(),
"paneId() must be the pane the agent.start call actually targeted");
assertTrue(started.terminalId() != null && started.terminalId().startsWith("term_new_"),
"terminalId() must be herdr's own generated id: " + started.terminalId());
}
/** The new tab is labelled with the lead's configured {@code tab:}, and AFTER the start. */
@Test
void relaunchLabelsTheNewTabAfterStarting() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNotNull(started);
assertEquals("lead: opus", ((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"));
int startIndex = indexOfLastCall(herdr, "agent.start");
int renameIndex = indexOfLastCall(herdr, "tab.rename");
assertTrue(renameIndex > startIndex,
"the tab must be renamed AFTER the start succeeds, not before: start=" + startIndex
+ " rename=" + renameIndex);
}
private static int indexOfLastCall(FakeHerdr herdr, String method) {
int idx = -1;
List<FakeHerdr.Call> calls = herdr.calls;
for (int i = 0; i < calls.size(); i++) {
if (calls.get(i).method().equals(method)) {
idx = i;
}
}
return idx;
}
@Test
void relaunchOfAnUnknownLeadNameReturnsNullAndStartsNothing() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("not-declared");
assertNull(started);
assertFalse(herdr.called("agent.start"));
assertFalse(herdr.called("workspace.create"));
assertFalse(herdr.called("tab.create"));
}
@Test
void relaunchOfARecogniseOnlyLeadReturnsNullAndStartsNothing() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead(null, "lead: dead", 1))).relaunch("opus");
assertNull(started);
assertFalse(herdr.called("agent.start"));
}
@Test
void relaunchWithAnUnconfiguredProfileReturnsNullAndStartsNothing() {
FakeHerdr herdr = new FakeHerdr();
dev.ltms.fleet.herdr.Agent started =
launcher(herdr, configWith(lead("nope", "lead: opus", 1))).relaunch("opus");
assertNull(started);
assertFalse(herdr.called("agent.start"));
}
/**
* The outer retry {@link LeadLauncher#relaunch(String)} owns, separate from {@code
* ResilientAgentLaunch}'s internal {@code agent_name_taken} retry: a failed attempt must not
* be the end of the whole relaunch. Each of the first two attempts exhausts {@code
* ResilientAgentLaunch.NAME_RETRIES} name attempts (every one of them rejected), so each
* attempt's own tab is created and then closed; the third attempt's first name is free.
*/
@Test
void relaunchRetriesTheWholeAttemptAndSucceedsOnTheThird() {
FakeHerdr herdr = new FakeHerdr()
.agentNameTakenTimes(2 * ResilientAgentLaunch.NAME_RETRIES);
dev.ltms.fleet.herdr.Agent started =
fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNotNull(started, "the third attempt's first name is free — it must succeed");
assertEquals(3, herdr.calls.stream().filter(c -> c.method().equals("tab.create")).count(),
"one tab per attempt: three attempts");
assertEquals(2, herdr.calls.stream().filter(c -> c.method().equals("tab.close")).count(),
"the two failed attempts' tabs must be closed");
}
/**
* Every attempt fails outright (a herdr error {@code ResilientAgentLaunch} does not retry at
* all) — {@link LeadLauncher#relaunch(String)} must give up after exactly {@code
* RELAUNCH_ATTEMPTS} and must not leak any of the tabs it created along the way.
*/
@Test
void relaunchGivesUpAfterExactlyRelaunchAttemptsAndLeaksNoTab() {
FakeHerdr herdr = new FakeHerdr().agentStartFailsWith("some_other_error");
dev.ltms.fleet.herdr.Agent started =
fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).relaunch("opus");
assertNull(started, "every attempt failed — relaunch must give up, not hang or guess");
assertEquals(LeadLauncher.RELAUNCH_ATTEMPTS,
herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
"exactly RELAUNCH_ATTEMPTS attempts, no more, no fewer");
long tabsCreated = herdr.calls.stream().filter(c -> c.method().equals("tab.create")).count();
long tabsClosed = herdr.calls.stream().filter(c -> c.method().equals("tab.close")).count();
assertEquals(LeadLauncher.RELAUNCH_ATTEMPTS, tabsCreated);
assertEquals(tabsCreated, tabsClosed, "every tab this method created must be closed — no leaks");
}
}
@@ -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());
}
}
@@ -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();