Compare commits
13 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 428a12af62 | |||
| a332dfdb2c | |||
| d0f4ae057b | |||
| 7b3beaa209 | |||
| b3b2bf3da6 | |||
| 8bb2aa0be4 | |||
| 4ce3149bfd | |||
| 1fc9e85bf1 | |||
| 38544d467c | |||
| 809b7d9b20 | |||
| 8cf7215d56 | |||
| cf0c9b9316 | |||
| 337dbd491e |
@@ -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());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -355,6 +355,14 @@ public final class MessageService {
|
||||
*/
|
||||
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
|
||||
private final AtomicLong ticketSeq = new AtomicLong();
|
||||
/**
|
||||
* Minted once per {@code MessageService} instance and folded into every ticket id (see
|
||||
* {@link #sendAsync(String, String, Runnable, String)}). {@link #ticketSeq} alone restarts at
|
||||
* zero for every instance, so without this a ticket id can be reused across instances and
|
||||
* resolve to an unrelated {@link Task} 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 ticketBootNonce = UUID.randomUUID().toString().substring(0, 6);
|
||||
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
|
||||
Thread.ofVirtual().name("bridge-async-", 0).factory());
|
||||
|
||||
@@ -1321,7 +1329,7 @@ public final class MessageService {
|
||||
* @return the ticket to poll for the eventual result
|
||||
*/
|
||||
public String sendAsync(String target, String content, Runnable onAccepted, String creatorTerminal) {
|
||||
String ticket = "task-" + ticketSeq.incrementAndGet();
|
||||
String ticket = "task-" + ticketBootNonce + "-" + ticketSeq.incrementAndGet();
|
||||
Task task = new Task(ticket, target, nowNanos, creatorTerminal);
|
||||
tasks.put(ticket, task);
|
||||
if (pushLoop != null) {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -888,6 +888,39 @@ class MessageServiceTest {
|
||||
assertNull(messages.poll("task-999999"), "a ticket that was never minted is unknown");
|
||||
}
|
||||
|
||||
// --- fleetd #719: a per-boot nonce keeps one instance's ticket ids out of another's space ---
|
||||
|
||||
/** A second, fully independent instance — its own agents/injector/rendezvous/inbox, not shared. */
|
||||
private MessageService newIndependentInstance() {
|
||||
FakeHerdr otherHerdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
|
||||
AgentControl otherAgents = new AgentControl(otherHerdr);
|
||||
Injector otherInjector = new Injector(otherAgents);
|
||||
return new MessageService(otherAgents, otherInjector, new Rendezvous(), new InMemoryReplyInbox());
|
||||
}
|
||||
|
||||
@Test
|
||||
void twoInstancesMintDisjointTicketIds() {
|
||||
MessageService other = newIndependentInstance();
|
||||
String ticketFromThis = messages.sendAsync(T, "task on first instance", null, null);
|
||||
String ticketFromOther = other.sendAsync(T, "task on second instance", null, null);
|
||||
assertNotEquals(ticketFromThis, ticketFromOther,
|
||||
"each instance mints its own id space, so even a first ticket from each must differ");
|
||||
}
|
||||
|
||||
@Test
|
||||
void foreignInstanceTicketDoesNotResolve() {
|
||||
MessageService other = newIndependentInstance();
|
||||
String ticket = messages.sendAsync(T, "task on first instance", null, null);
|
||||
// `other` must reach the same sequence number, or this test passes against an empty map
|
||||
// instead of against a colliding id.
|
||||
other.sendAsync(T, "task on second instance", null, null);
|
||||
|
||||
// control: the id resolves in the instance that minted it, so a null below cannot be
|
||||
// explained by broken plumbing — only by the ticket being foreign to `other`.
|
||||
assertNotNull(messages.poll(ticket), "the minting instance must still resolve its own ticket");
|
||||
assertNull(other.poll(ticket), "a ticket minted by a different instance must not resolve here");
|
||||
}
|
||||
|
||||
// --- fleetd #705: a ticket's creator terminal gates who may poll it -------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -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