Compare commits
14 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 1fc9e85bf1 | |||
| 809b7d9b20 | |||
| 8cf7215d56 | |||
| 1ef93e57cc | |||
| cf0c9b9316 | |||
| 337dbd491e | |||
| 7cf6075b79 | |||
| fbdcd709c9 | |||
| 8d3f10d291 | |||
| 38f4fd64ee | |||
| 3fc39b981d | |||
| 70328ca0f8 | |||
| efab9b8c49 | |||
| efeffb4ab7 |
@@ -102,6 +102,8 @@ below are the procedure — run them in order, every task, not only the big ones
|
||||
`fleet_status`, never by reading its terminal; it also reports an open question and the `turnId`
|
||||
that answers it — but **only to the caller that created that delegation**, so a question raised
|
||||
under an architect's brief is invisible to you, and seeing none does not mean there is none.
|
||||
**Only that same creator can answer it.** A `turnId` you came by any other way is refused, so an
|
||||
architect's worker waits for that architect and not for you.
|
||||
**A worker's ask waits ~55 seconds, and no nudge makes that longer** — so never
|
||||
brief a worker to "ask me". Decide before you delegate, or give it an explicit default.
|
||||
**A correction cannot reach a busy member.** A `fleet_send` to a working member is *accepted* and
|
||||
@@ -167,7 +169,7 @@ you decide.
|
||||
| See the fleet | `fleet_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) + `loopHealth` (`RUNNING`, `STALLED`, or `STOPPED` for `statusPoller` and `sessionReaper`) · one peer's state: `fleet_status{sessionId}` |
|
||||
| Delegate (blocking) | `fleet_send{sessionId, content}` |
|
||||
| Delegate (long task) | `fleet_send{sessionId, content, wait:false}` → ticket → `fleet_poll{ticket}` |
|
||||
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` |
|
||||
| Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId`, and only the caller that created that delegation |
|
||||
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
|
||||
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task |
|
||||
| Message a **collaborator** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` reports a `collaborators` array, and each row carries that peer's `name` and the `sessionId` you send to. It is visible to you, to an architect and to another collaborator, never to a worker. Coordination only, **never** a task |
|
||||
@@ -332,10 +334,22 @@ must obey belongs in the charter, not here.
|
||||
(#362). **Read `plugin/` before designing anything about onboarding a project.** Two limits are
|
||||
structural, not bugs: a plugin cannot carry the role agent files, because
|
||||
`ClaudeCodeLauncher.java:371` requires `<cwd>/.claude/agents/<role>.md` in the member's own
|
||||
worktree; and a plugin cannot deliver anything to members at all, because
|
||||
`ClaudeCodeLauncher.java:285` exports `CLAUDE_CONFIG_DIR` and every Claude profile here sets it,
|
||||
so a member never reads the operator's plugin store. **The plugin is the lead-side surface;
|
||||
member-facing assets travel in the worktree.**
|
||||
worktree; and a plugin reaches a member only through `CLAUDE_CONFIG_DIR`, which
|
||||
`ClaudeCodeLauncher.java:286` exports with `putIfPresent` — so only for a profile that sets
|
||||
`configDir`. Every `claude-code` profile does set one (the four without are `opencode`, which
|
||||
never reads that variable). **But measured 2026-10-04: two of them point at
|
||||
`~/.ccs/instances/ltms`, which is the operator's own `CLAUDE_CONFIG_DIR` on this host.** So for an
|
||||
`opus` or `sonnet` member, "a member never reads the operator's plugin store" is false — it reads
|
||||
the same store, because that store is the one its `configDir` names. It stays true for `local` and
|
||||
`local-direct`, which point at `~/.ccs/instances/gx10`. `ClaudeCodeLauncher`'s own javadoc names
|
||||
the related hazard: that file is rewritten on every spawn, so for those two profiles fleetd and the
|
||||
operator's live session write the same `.claude.json`, and its compare-and-swap "narrows the
|
||||
lost-update window, it does not close it". Re-measure which profiles share the operator's dir with
|
||||
`awk '/^profiles:/{i=1;next} /^[a-z]/{i=0} i&&/^ [a-z-]+:$/{p=$1} i&&/configDir:/{print p,$2}'
|
||||
fleetd/fleetd.yaml` against `echo $CLAUDE_CONFIG_DIR`; delete this note once no profile names the
|
||||
operator's dir. **Treat the plugin as the lead-side surface and put member-facing assets in the
|
||||
worktree** — that conclusion holds either way, because a worktree asset does not depend on which
|
||||
config dir a member reads.
|
||||
- **Never commit** `.mcp.json` (the primary's local copy, flagged `--skip-worktree`) or `wiki/`
|
||||
(a submodule with its own remote).
|
||||
- **A provisioned worktree neutralizes `.mcp.json`, `opencode.json` and `.autoenv`** — the repo's
|
||||
|
||||
@@ -51,7 +51,6 @@ import dev.ltms.fleet.power.CaffeinateSleepAssertionMechanism;
|
||||
import dev.ltms.fleet.power.IdleSleepGuard;
|
||||
import dev.ltms.fleet.rest.FleetApp;
|
||||
import dev.ltms.fleet.session.GitWorktrees;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.SessionReaper;
|
||||
import io.javalin.Javalin;
|
||||
@@ -480,13 +479,9 @@ final class FleetdAssembly {
|
||||
new PaneLocator(herdr, memberHerdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
||||
|
||||
// fleetd #669 Unit D: a live spawned member resolves as its own role, whatever a tab map
|
||||
// says about the same terminal — read from the roster meant for a hot path (SessionManager
|
||||
// javadoc), never rosterResolved(), since resolve() runs on every request.
|
||||
Function<String, MemberRole> spawnedMemberRole = terminal -> sessions.roster().stream()
|
||||
.filter(s -> terminal.equals(s.terminalId()))
|
||||
.map(MemberSession::role)
|
||||
.findFirst()
|
||||
.orElse(null);
|
||||
// says about the same terminal. fleetd #702: SessionManager.spawnedMemberRole also answers
|
||||
// for a pane mid-teardown, not only one still in the registry — see its javadoc.
|
||||
Function<String, MemberRole> spawnedMemberRole = sessions::spawnedMemberRole;
|
||||
|
||||
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
|
||||
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
|
||||
|
||||
@@ -0,0 +1,146 @@
|
||||
package dev.ltms.fleet.herdr;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.List;
|
||||
import java.util.function.IntFunction;
|
||||
|
||||
/**
|
||||
* The three protections every {@code agent.start} caller needs against herdr's pane-typed launch
|
||||
* surface (fleetd #220, #727): a byte-limit check on the assembled command line, a bounded retry
|
||||
* on {@code agent_pane_busy} (the target pane's shell has not reached its prompt yet), and a
|
||||
* bounded retry on {@code agent_name_taken} with a fresh name each attempt. One implementation —
|
||||
* every caller of {@code agent.start}, lead or member, goes through this seam rather than carrying
|
||||
* its own copy.
|
||||
*/
|
||||
public final class ResilientAgentLaunch {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(ResilientAgentLaunch.class);
|
||||
|
||||
private ResilientAgentLaunch() {
|
||||
}
|
||||
|
||||
/**
|
||||
* The pty line buffer herdr types a launch command into: BSD/macOS {@code MAX_CANON}. Not a
|
||||
* fleetd choice and not configurable — see {@link #checkFits}.
|
||||
*/
|
||||
public static final int PANE_COMMAND_BYTE_LIMIT = 1024;
|
||||
|
||||
/** Per-argument allowance for the separating space and a shell quote pair fleetd cannot see. */
|
||||
private static final int QUOTING_OVERHEAD_PER_ARG = 3;
|
||||
|
||||
/** herdr rejects a duplicate agent {@code name}; a caller retries a bumped name this many times. */
|
||||
public static final int NAME_RETRIES = 8;
|
||||
|
||||
/**
|
||||
* Retries for {@code agent.start} against a seed pane whose shell has not reached its prompt
|
||||
* yet — {@code tab.create}/{@code pane.split} return as soon as the pane exists, and herdr
|
||||
* refuses to start an agent in a pane that is not "an available shell" ({@code agent_pane_busy}).
|
||||
*/
|
||||
public static final int SHELL_READY_RETRIES = 20;
|
||||
|
||||
/** Raised by {@link #checkFits} when the assembled command cannot fit the pane line. */
|
||||
public static final class TooLargeException extends RuntimeException {
|
||||
public TooLargeException(String message) {
|
||||
super(message);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Verify the assembled launch line fits the pane herdr types it into. herdr does not exec the
|
||||
* launch command — it TYPES it into the pane as one line, and a pty line buffer holds only
|
||||
* {@value #PANE_COMMAND_BYTE_LIMIT} bytes. Everything past that byte is dropped with no error
|
||||
* anywhere: herdr answers "agent started", the backend exits on the mangled argument it was
|
||||
* handed, the pane closes, and the only symptom is a readiness timeout with no reason. That is
|
||||
* how fleetd #214 broke every claude-code spawn — one 50-byte flag pushed a 978-byte command to
|
||||
* 1028, and the tail that got cut was {@code --autocompact 250000}.
|
||||
*
|
||||
* <p>So measure it here and refuse, loudly and immediately, rather than start something that
|
||||
* cannot work. The estimate is deliberately conservative: fleetd cannot see herdr's quoting, so
|
||||
* every argument is charged its own bytes plus a separator and a quote pair. An over-estimate
|
||||
* costs a clear error at a length that was already unsafe; an under-estimate would let the
|
||||
* silent truncation back in.
|
||||
*
|
||||
* @param label names the launch in the refusal message (a profile name)
|
||||
* @param argv the full argv, including the executable at index 0
|
||||
* @throws TooLargeException naming the limit, the estimate, and the longest argument
|
||||
*/
|
||||
public static void checkFits(String label, List<String> argv) {
|
||||
int bytes = 0;
|
||||
String longest = null;
|
||||
int longestBytes = 0;
|
||||
for (String arg : argv) {
|
||||
int argBytes = arg == null ? 0 : arg.getBytes(StandardCharsets.UTF_8).length;
|
||||
bytes += argBytes + QUOTING_OVERHEAD_PER_ARG;
|
||||
if (argBytes > longestBytes) {
|
||||
longestBytes = argBytes;
|
||||
longest = arg;
|
||||
}
|
||||
}
|
||||
if (bytes <= PANE_COMMAND_BYTE_LIMIT) {
|
||||
return;
|
||||
}
|
||||
String culprit = longest == null ? "<none>"
|
||||
: longest.substring(0, Math.min(longest.length(), 60)) + (longest.length() > 60 ? "…" : "");
|
||||
throw new TooLargeException(
|
||||
"launch command for " + label + " is about " + bytes + " bytes, over the "
|
||||
+ PANE_COMMAND_BYTE_LIMIT + "-byte limit of the pane line herdr types it into. "
|
||||
+ "The pty would drop the tail silently and the backend would exit on a mangled "
|
||||
+ "argument. Longest argument is " + longestBytes + " bytes: " + culprit
|
||||
+ " — move it off the command line (a file flag) or shorten it.");
|
||||
}
|
||||
|
||||
/**
|
||||
* Start an agent into {@code paneId}, retrying {@code agent_pane_busy} up to {@code retries}
|
||||
* times with {@code sleeper} run between attempts.
|
||||
*
|
||||
* @throws HerdrException the last {@code agent_pane_busy} failure once {@code retries} is
|
||||
* spent, or immediately for any other herdr failure
|
||||
*/
|
||||
public static Agent startAwaitingShellPrompt(AgentControl agents, String name, String kind,
|
||||
List<String> args, String paneId,
|
||||
int retries, Runnable sleeper) {
|
||||
HerdrException busy = null;
|
||||
for (int attempt = 0; attempt < retries; attempt++) {
|
||||
try {
|
||||
return agents.start(name, kind, args, paneId);
|
||||
} catch (HerdrException e) {
|
||||
if (!"agent_pane_busy".equals(e.code())) throw e;
|
||||
log.debug("pane {} not at its shell prompt yet, retrying agent.start", paneId);
|
||||
busy = e;
|
||||
sleeper.run();
|
||||
}
|
||||
}
|
||||
throw busy;
|
||||
}
|
||||
|
||||
/**
|
||||
* Start an agent under a freshly generated name each attempt, retrying {@code agent_name_taken}
|
||||
* up to {@code nameRetries} times — herdr refuses a duplicate {@code name} outright, so a stale
|
||||
* registry entry (a crashed session, a name the registry has not yet released) must not block a
|
||||
* legitimate relaunch. Each attempt also carries its own {@link #startAwaitingShellPrompt} retry.
|
||||
*
|
||||
* @param nameForAttempt called once per attempt (0-based) to produce that attempt's name
|
||||
* @throws HerdrException the last {@code agent_name_taken} failure once {@code nameRetries} is
|
||||
* spent, or immediately for any other herdr failure
|
||||
*/
|
||||
public static Agent startUniquelyNamed(AgentControl agents, String kind, List<String> args,
|
||||
String paneId, IntFunction<String> nameForAttempt,
|
||||
int nameRetries, int shellReadyRetries, Runnable sleeper) {
|
||||
HerdrException last = null;
|
||||
for (int attempt = 0; attempt < nameRetries; attempt++) {
|
||||
String name = nameForAttempt.apply(attempt);
|
||||
try {
|
||||
return startAwaitingShellPrompt(agents, name, kind, args, paneId,
|
||||
shellReadyRetries, sleeper);
|
||||
} catch (HerdrException e) {
|
||||
if (!"agent_name_taken".equals(e.code())) throw e;
|
||||
log.debug("agent name '{}' taken, retrying", name);
|
||||
last = e;
|
||||
}
|
||||
}
|
||||
throw last;
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,7 @@ import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.PendingCloseMarker;
|
||||
import dev.ltms.fleet.herdr.ResilientAgentLaunch;
|
||||
import dev.ltms.fleet.herdr.Tab;
|
||||
import dev.ltms.fleet.herdr.Workspace;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
@@ -12,6 +13,7 @@ import dev.ltms.fleet.launch.ClaudeCodeArguments;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.security.SecureRandom;
|
||||
import java.util.ArrayList;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.LinkedHashSet;
|
||||
@@ -19,6 +21,7 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
|
||||
/**
|
||||
@@ -83,6 +86,13 @@ public final class LeadLauncher {
|
||||
private final AgentControl agents;
|
||||
private final WorkspaceControl spaces;
|
||||
private final FleetConfig cfg;
|
||||
private final Runnable sleeper;
|
||||
|
||||
// Per-process token mixed into each lead agent name so a fresh daemon process (seq back at 0)
|
||||
// cannot collide with a same-name lead that outlived a restart — the same scheme
|
||||
// HerdrPeerLauncher uses for members (fleetd #727).
|
||||
private final String nameNonce = String.format("%06x", new SecureRandom().nextInt(1 << 24));
|
||||
private final AtomicLong nameSeq = new AtomicLong();
|
||||
|
||||
/**
|
||||
* @param agents herdr agent control (start, list)
|
||||
@@ -90,9 +100,27 @@ public final class LeadLauncher {
|
||||
* @param cfg the loaded config — {@code fleet.leaders}, {@code profiles} and each lead's tab
|
||||
*/
|
||||
public LeadLauncher(AgentControl agents, WorkspaceControl spaces, FleetConfig cfg) {
|
||||
this(agents, spaces, cfg, () -> sleepUninterruptibly(300));
|
||||
}
|
||||
|
||||
/**
|
||||
* Test seam: as above, plus an injectable {@code sleeper} for the {@code agent_pane_busy}
|
||||
* retry (fleetd #727), so a test can prove the retry budget without a real sleep.
|
||||
*/
|
||||
LeadLauncher(AgentControl agents, WorkspaceControl spaces, FleetConfig cfg, Runnable sleeper) {
|
||||
this.agents = agents;
|
||||
this.spaces = spaces;
|
||||
this.cfg = cfg;
|
||||
this.sleeper = sleeper;
|
||||
}
|
||||
|
||||
/** Uninterruptible sleep — the production {@link #sleeper} between {@code agent_pane_busy} retries. */
|
||||
private static void sleepUninterruptibly(long ms) {
|
||||
try {
|
||||
Thread.sleep(ms);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -306,7 +334,16 @@ public final class LeadLauncher {
|
||||
return null;
|
||||
}
|
||||
|
||||
/** Start one lead. Returns false (having logged) rather than throwing on any failure. */
|
||||
/**
|
||||
* Start one lead. Returns false (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
|
||||
* types it into, a stale {@code agent_name_taken} (a crashed session's name the registry has
|
||||
* not yet released) is retried under a fresh per-attempt name rather than refusing the whole
|
||||
* 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) {
|
||||
String label = lead.tabLabel();
|
||||
String cwd = (lead.cwd() == null || lead.cwd().isBlank())
|
||||
@@ -323,8 +360,12 @@ public final class LeadLauncher {
|
||||
// Same shape as the member launchers: herdr resolves the executable from `kind`, so
|
||||
// argv[0] (the configured launcher, e.g. `ccs`) is dropped and only the rest is passed.
|
||||
List<String> argv = leadArgv(profile);
|
||||
Agent started = agents.start("lead-" + name, herdrKind(profile),
|
||||
argv.isEmpty() ? argv : argv.subList(1, argv.size()), tab.rootPaneId());
|
||||
ResilientAgentLaunch.checkFits(profile.profile(), argv);
|
||||
List<String> args = argv.isEmpty() ? argv : argv.subList(1, argv.size());
|
||||
Agent started = ResilientAgentLaunch.startUniquelyNamed(agents, herdrKind(profile), args,
|
||||
tab.rootPaneId(),
|
||||
attempt -> "lead-" + name + "-" + nameNonce + "-" + nameSeq.incrementAndGet(),
|
||||
ResilientAgentLaunch.NAME_RETRIES, ResilientAgentLaunch.SHELL_READY_RETRIES, sleeper);
|
||||
|
||||
// Label AFTER the start succeeds. A label written before would survive a failed start
|
||||
// and then read back as a live lead on the next boot, which is the exact staleness the
|
||||
|
||||
@@ -6,6 +6,7 @@ import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.ResilientAgentLaunch;
|
||||
import dev.ltms.fleet.herdr.Tab;
|
||||
import dev.ltms.fleet.herdr.Workspace;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
@@ -67,16 +68,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
|
||||
/** herdr rejects a duplicate agent {@code name}; we retry a bumped name this many times. */
|
||||
private static final int NAME_RETRIES = 8;
|
||||
|
||||
/**
|
||||
* Retries for {@code agent.start} against a seed pane whose shell has not reached its prompt
|
||||
* yet — {@code tab.create}/{@code pane.split} return as soon as the pane exists, and herdr
|
||||
* refuses to start an agent in a pane that is not "an available shell" ({@code agent_pane_busy}).
|
||||
*/
|
||||
private static final int SHELL_READY_RETRIES = 20;
|
||||
|
||||
private final String namePrefix; // label prefix: naming + reap scheme
|
||||
private final AgentControl agents;
|
||||
private final WorkspaceControl spaces;
|
||||
@@ -766,90 +757,26 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
// Protocol 19 resolves the executable from the agent kind (== namePrefix here), so
|
||||
// argv[0] — the configured executable — is dropped and only the extra args are passed.
|
||||
List<String> args = argv.isEmpty() ? argv : argv.subList(1, argv.size());
|
||||
checkPaneCommandFits(cfg, argv);
|
||||
HerdrException last = null;
|
||||
for (int attempt = 0; attempt < NAME_RETRIES; attempt++) {
|
||||
long seq = nameSeq.incrementAndGet();
|
||||
String name = namePrefix + "-" + cfg.profile() + "-" + nameNonce + "-" + seq;
|
||||
try {
|
||||
return new Started(startAwaitingShellPrompt(name, args, paneId), seq);
|
||||
} catch (HerdrException e) {
|
||||
if (!"agent_name_taken".equals(e.code())) throw e;
|
||||
log.debug("peer name '{}' taken, retrying", name);
|
||||
last = e;
|
||||
}
|
||||
try {
|
||||
ResilientAgentLaunch.checkFits(cfg.profile(), argv);
|
||||
} catch (ResilientAgentLaunch.TooLargeException e) {
|
||||
throw new PeerUnreachableException(e.getMessage());
|
||||
}
|
||||
throw last;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #220: herdr does not exec the launch command — it TYPES it into the pane as one line,
|
||||
* and a pty line buffer holds only {@value #PANE_COMMAND_BYTE_LIMIT} bytes (BSD/macOS {@code
|
||||
* MAX_CANON}). Everything past that byte is dropped. Nothing reports it: herdr answers "agent
|
||||
* started", the backend exits on the mangled argument it was handed, the pane closes, and the
|
||||
* only symptom is {@link #waitUntilInjectableOrThrow} timing out 20 seconds later with no
|
||||
* reason. That is exactly how #214 broke every claude-code spawn — one 50-byte flag pushed a
|
||||
* 978-byte command to 1028, and the tail that got cut was {@code --autocompact 250000}.
|
||||
*
|
||||
* <p>So measure it here and refuse, loudly and immediately, rather than spawn something that
|
||||
* cannot work. The estimate is deliberately conservative: fleetd cannot see herdr's quoting, so
|
||||
* every argument is charged its own bytes plus a separator and a quote pair. An over-estimate
|
||||
* costs a clear error at a length that was already unsafe; an under-estimate would let the
|
||||
* silent truncation back in.
|
||||
*
|
||||
* @throws PeerUnreachableException when the command cannot fit — the same failure the spawn
|
||||
* would have hit anyway, named at the point it is still
|
||||
* explainable
|
||||
*/
|
||||
private void checkPaneCommandFits(FleetConfig.Profile cfg, List<String> argv) {
|
||||
int bytes = 0;
|
||||
String longest = null;
|
||||
int longestBytes = 0;
|
||||
for (String arg : argv) {
|
||||
int argBytes = arg == null ? 0 : arg.getBytes(java.nio.charset.StandardCharsets.UTF_8).length;
|
||||
bytes += argBytes + QUOTING_OVERHEAD_PER_ARG;
|
||||
if (argBytes > longestBytes) {
|
||||
longestBytes = argBytes;
|
||||
longest = arg;
|
||||
}
|
||||
}
|
||||
if (bytes <= PANE_COMMAND_BYTE_LIMIT) {
|
||||
return;
|
||||
}
|
||||
String culprit = longest == null ? "<none>"
|
||||
: longest.substring(0, Math.min(longest.length(), 60)) + (longest.length() > 60 ? "…" : "");
|
||||
throw new PeerUnreachableException(
|
||||
"launch command for profile " + cfg.profile() + " is about " + bytes + " bytes, over the "
|
||||
+ PANE_COMMAND_BYTE_LIMIT + "-byte limit of the pane line herdr types it into. "
|
||||
+ "The pty would drop the tail silently and the backend would exit on a mangled "
|
||||
+ "argument. Longest argument is " + longestBytes + " bytes: " + culprit
|
||||
+ " — move it off the command line (a file flag) or shorten it.");
|
||||
long[] lastSeq = {0};
|
||||
Agent agent = ResilientAgentLaunch.startUniquelyNamed(agents, namePrefix, args, paneId,
|
||||
attempt -> {
|
||||
lastSeq[0] = nameSeq.incrementAndGet();
|
||||
return namePrefix + "-" + cfg.profile() + "-" + nameNonce + "-" + lastSeq[0];
|
||||
},
|
||||
ResilientAgentLaunch.NAME_RETRIES, ResilientAgentLaunch.SHELL_READY_RETRIES, sleeper);
|
||||
return new Started(agent, lastSeq[0]);
|
||||
}
|
||||
|
||||
/**
|
||||
* The pty line buffer herdr types a launch command into: BSD/macOS {@code MAX_CANON}. Not a
|
||||
* fleetd choice and not configurable — see {@link #checkPaneCommandFits}.
|
||||
* fleetd choice and not configurable — see {@link ResilientAgentLaunch#checkFits}.
|
||||
*/
|
||||
static final int PANE_COMMAND_BYTE_LIMIT = 1024;
|
||||
|
||||
/** Per-argument allowance for the separating space and a shell quote pair fleetd cannot see. */
|
||||
private static final int QUOTING_OVERHEAD_PER_ARG = 3;
|
||||
|
||||
/** Start the agent into {@code paneId}, waiting out the seed shell's boot with the sleeper. */
|
||||
private Agent startAwaitingShellPrompt(String name, List<String> args, String paneId) {
|
||||
HerdrException busy = null;
|
||||
for (int attempt = 0; attempt < SHELL_READY_RETRIES; attempt++) {
|
||||
try {
|
||||
return agents.start(name, namePrefix, args, paneId);
|
||||
} catch (HerdrException e) {
|
||||
if (!"agent_pane_busy".equals(e.code())) throw e;
|
||||
log.debug("pane {} not at its shell prompt yet, retrying agent.start", paneId);
|
||||
busy = e;
|
||||
sleeper.run();
|
||||
}
|
||||
}
|
||||
throw busy;
|
||||
}
|
||||
static final int PANE_COMMAND_BYTE_LIMIT = ResilientAgentLaunch.PANE_COMMAND_BYTE_LIMIT;
|
||||
|
||||
// --- discovery + reap ----------------------------------------------------------------------
|
||||
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -26,6 +26,7 @@ import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
/**
|
||||
* Authoritative in-daemon registry of the worker sessions this {@code fleetd} process spawned.
|
||||
@@ -57,6 +58,18 @@ public final class SessionManager implements TurnListener {
|
||||
* Populated on every spawn path, removed on {@link #release}.
|
||||
*/
|
||||
private final ConcurrentHashMap<String /*paneId*/, PeerHandle> handles = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* fleetd #702: a pane mid-teardown, keyed by paneId, held from just before its registry entry
|
||||
* is removed until {@link #releaseRemoved} finishes. {@link #spawnedMemberRole} consults this
|
||||
* alongside the registry, so a caller resolving the pane's terminal during that window still
|
||||
* sees a live member and never falls through to a tab map.
|
||||
*
|
||||
* <p>Depth-counted rather than a plain set: two threads can be tearing down the same pane at
|
||||
* once (the CAS in {@link #releaseIfCurrent} exists for exactly that race), and with a set the
|
||||
* loser's {@code finally} would unmark the pane while the winner is still mid-teardown,
|
||||
* reopening the window this exists to close.
|
||||
*/
|
||||
private final ConcurrentHashMap<String /*paneId*/, Releasing> releasing = new ConcurrentHashMap<>();
|
||||
private final MemberPresence presence;
|
||||
private final SecureRandom nonceRandom = new SecureRandom();
|
||||
private final AtomicLong nonceSeq = new AtomicLong();
|
||||
@@ -303,9 +316,11 @@ public final class SessionManager implements TurnListener {
|
||||
* is a logged path an operator can reclaim, the cost of a deleted one is unrecoverable work.
|
||||
*/
|
||||
private MemberSession release(String paneId, ReleaseCause cause) {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
|
||||
return removed;
|
||||
return releaseWindow(paneId, registry.get(paneId), () -> {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
|
||||
return removed;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -314,16 +329,85 @@ public final class SessionManager implements TurnListener {
|
||||
* DONE record from stopping a worker that delivery has made BUSY.
|
||||
*/
|
||||
private boolean releaseIfCurrent(MemberSession expected, ReleaseCause cause) {
|
||||
if (!registry.remove(expected.paneId(), expected)) {
|
||||
// A lifecycle transition replaced the record between the caller's check and this remove.
|
||||
// Log it: this race is by definition unobservable otherwise, and a reaper that silently
|
||||
// declines to reap is the hardest kind of behaviour to diagnose after the fact.
|
||||
log.debug("skipping reap of pane={}: its registry record changed after the idle check "
|
||||
+ "(most likely a delivery made it BUSY)", expected.paneId());
|
||||
return false;
|
||||
return releaseWindow(expected.paneId(), expected, () -> {
|
||||
if (!registry.remove(expected.paneId(), expected)) {
|
||||
// A lifecycle transition replaced the record between the caller's check and this
|
||||
// remove. Log it: this race is by definition unobservable otherwise, and a reaper
|
||||
// that silently declines to reap is the hardest kind of behaviour to diagnose
|
||||
// after the fact.
|
||||
log.debug("skipping reap of pane={}: its registry record changed after the idle "
|
||||
+ "check (most likely a delivery made it BUSY)", expected.paneId());
|
||||
return false;
|
||||
}
|
||||
releaseRemoved(expected.paneId(), expected, handles.remove(expected.paneId()), cause);
|
||||
return true;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #702: mark {@code paneId} as mid-teardown — using {@code known}'s terminal/role when
|
||||
* it is available — for the whole of {@code teardown}, which removes the registry entry and
|
||||
* then runs {@link #releaseRemoved}. Shared by both registry-removal sites ({@link #release}'s
|
||||
* unconditional remove and {@link #releaseIfCurrent}'s CAS remove) so neither can leave the
|
||||
* other's window unmarked.
|
||||
*
|
||||
* <p>The mark is written before {@code teardown} runs — so it covers the removal itself, not
|
||||
* only what comes after it — and cleared in a {@code finally}, so an unchecked throw out of
|
||||
* {@code teardown} (including one from {@link PeerLauncher#stop}, which declares nothing) can
|
||||
* never leave the pane marked for the rest of the daemon's life.
|
||||
*/
|
||||
private <T> T releaseWindow(String paneId, MemberSession known, Supplier<T> teardown) {
|
||||
releasing.compute(paneId, (_, prior) -> Releasing.enter(prior, known));
|
||||
try {
|
||||
return teardown.get();
|
||||
} finally {
|
||||
releasing.compute(paneId, (_, prior) -> prior == null ? null : prior.leave());
|
||||
}
|
||||
releaseRemoved(expected.paneId(), expected, handles.remove(expected.paneId()), cause);
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Depth count plus the terminal/role a mid-teardown pane belongs to, for
|
||||
* {@link #spawnedMemberRole}. The terminal/role come from whichever call into
|
||||
* {@link #releaseWindow} first knew them: a call that finds the registry entry already gone
|
||||
* passes a {@code null} session, and must not blank out what the first call recorded.
|
||||
*/
|
||||
record Releasing(int depth, String terminalId, MemberRole role) {
|
||||
static Releasing enter(Releasing prior, MemberSession known) {
|
||||
int depth = (prior == null ? 0 : prior.depth()) + 1;
|
||||
String terminalId = known != null ? known.terminalId() : prior == null ? null : prior.terminalId();
|
||||
MemberRole role = known != null ? known.role() : prior == null ? null : prior.role();
|
||||
return new Releasing(depth, terminalId, role);
|
||||
}
|
||||
|
||||
Releasing leave() {
|
||||
return depth <= 1 ? null : new Releasing(depth - 1, terminalId, role);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The role of the live spawned member occupying {@code terminal} — whether it is currently in
|
||||
* the registry, or mid-teardown between {@link #release} removing its registry entry and
|
||||
* {@link #releaseRemoved} actually stopping its pane (fleetd #702). {@code null} for a terminal
|
||||
* that is neither: this method is the one reader a caller resolver consults before any tab
|
||||
* map, so a live or releasing member's identity never falls back to a tab label.
|
||||
*
|
||||
* <p>Checks the registry directly via {@link #findByTerminal} rather than {@link #roster()},
|
||||
* so this hot-path lookup (consulted on every resolve) never pays for a list copy or a stream.
|
||||
*/
|
||||
public MemberRole spawnedMemberRole(String terminal) {
|
||||
MemberSession session = findByTerminal(terminal);
|
||||
if (session != null) {
|
||||
return session.role();
|
||||
}
|
||||
if (terminal == null) {
|
||||
return null;
|
||||
}
|
||||
for (Releasing r : releasing.values()) {
|
||||
if (terminal.equals(r.terminalId())) {
|
||||
return r.role();
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
private void releaseRemoved(String paneId, MemberSession removed, PeerHandle removedHandle,
|
||||
|
||||
@@ -1,13 +1,25 @@
|
||||
package dev.ltms.fleet.auth;
|
||||
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.mcp.ConnectionIdentity;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.session.WorktreeRequest;
|
||||
import dev.ltms.fleet.session.Worktrees;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -541,6 +553,219 @@ class CallerResolverTest {
|
||||
assertEquals("term_a", p.terminal());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #702: between {@link SessionManager#release} removing the registry entry and the
|
||||
* pane actually stopping, a resolve for that terminal must still see the live member and
|
||||
* never fall through to a tab map — wired through the real {@link SessionManager}, not a
|
||||
* hand-rolled stand-in for its {@code spawnedMemberRole}.
|
||||
*
|
||||
* <p>Reuses the ticket's own test idea: an injected {@code hasUncommitted} resolves the
|
||||
* releasing terminal from inside {@code release}'s window — a real call landing inside the
|
||||
* window, so no sleep and no race.
|
||||
*
|
||||
* <p>The property asserted is "no tab map is consulted", never "the same role is returned".
|
||||
* The mandatory control is the lead tab map: it names this exact terminal, and the same
|
||||
* resolve taken <em>outside</em> the window (before release runs) must still return the lead
|
||||
* role — without that control, the in-window assertion would also pass on an empty map and
|
||||
* prove nothing.
|
||||
*/
|
||||
@Test
|
||||
void aPaneMidTeardownResolvesAsItsOwnRoleConsultingNoTabMap() {
|
||||
FakeHerdr herdr = new FakeHerdr().pinNextStarts(1, "term_a", "w2:p7");
|
||||
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);
|
||||
|
||||
AtomicReference<Principal> duringWindow = new AtomicReference<>();
|
||||
AtomicReference<CallerResolver> resolverRef = new AtomicReference<>();
|
||||
Worktrees worktrees = new Worktrees() {
|
||||
@Override
|
||||
public String add(String repoRoot, String branch, String baseRef) {
|
||||
return "/wt/" + branch.replace('/', '_');
|
||||
}
|
||||
|
||||
@Override
|
||||
public void remove(String repoRoot, String worktreePath) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteBranch(String repoRoot, String branch) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
// Runs from INSIDE release()'s git-status shell-out: the registry entry is
|
||||
// already gone, but the pane has not stopped yet.
|
||||
duringWindow.set(resolverRef.get().resolve("127.0.0.1", 42, null));
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String repoRoot(String cwd) {
|
||||
return "/repo";
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<String> snapshot(String worktreePath, String branch, String message) {
|
||||
return Optional.empty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public WipRefStats wipRefs(String repoRoot) {
|
||||
return new WipRefStats(0, 0L);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shareWithGroup(String repoRoot, String worktreePath) {
|
||||
}
|
||||
};
|
||||
SessionManager sessions = new SessionManager(workers, worktrees);
|
||||
|
||||
// The lead tab map names "term_a" before anything is ever spawned onto it — the control
|
||||
// this test needs. Built up front so the SAME resolver answers every resolve() call below.
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> FakeHerdr.WORKER_PID);
|
||||
CallerResolver resolver = CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
() -> Map.of("term_a", "the-lead"), null, sessions::spawnedMemberRole, Map::of);
|
||||
resolverRef.set(resolver);
|
||||
|
||||
Principal before = resolver.resolve("127.0.0.1", 42, null);
|
||||
assertEquals(Role.PRIMARY, before.role(),
|
||||
"control: with no live or releasing member on this terminal, the lead tab map must "
|
||||
+ "win — this is what proves the in-window assertion below is not passing on "
|
||||
+ "an empty map");
|
||||
assertEquals("the-lead", before.name());
|
||||
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("fleetd-702", null));
|
||||
assertEquals("term_a", s.terminalId(), "sanity: the spawn resolved to the pinned pane");
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertNotNull(duringWindow.get(), "the dirty check must have run and captured a resolve");
|
||||
assertEquals(Role.WORKER, duringWindow.get().role(),
|
||||
"inside the window the pane must resolve as its own live-member role, consulting no "
|
||||
+ "tab map — a lead tab naming the same terminal must not win");
|
||||
assertEquals("term_a", duringWindow.get().terminal());
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link SessionManager#releaseRemoved} unbinds the architect slot before the git-status
|
||||
* shell-out that opens the teardown window, so a resolve landing inside that window must see
|
||||
* the slot already unbound and resolve {@link Role#WORKER} — never {@link Role#ARCHITECT},
|
||||
* and never by checking role equality against the live session, which would hold even if the
|
||||
* unbind ran too late.
|
||||
*
|
||||
* <p>The control is the same resolve taken outside the window, while the slot is still bound,
|
||||
* which must return {@link Role#ARCHITECT} — without it this test would also pass against a
|
||||
* slot that was never bound, and prove nothing.
|
||||
*/
|
||||
@Test
|
||||
void aReleasingArchitectIsDemotedToWorkerInsideTheTeardownWindow() {
|
||||
FakeHerdr herdr = new FakeHerdr().pinNextStarts(1, "term_a", "w2:p7");
|
||||
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);
|
||||
|
||||
AtomicReference<Principal> duringWindow = new AtomicReference<>();
|
||||
AtomicReference<CallerResolver> resolverRef = new AtomicReference<>();
|
||||
Worktrees worktrees = new Worktrees() {
|
||||
@Override
|
||||
public String add(String repoRoot, String branch, String baseRef) {
|
||||
return "/wt/" + branch.replace('/', '_');
|
||||
}
|
||||
|
||||
@Override
|
||||
public void remove(String repoRoot, String worktreePath) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteBranch(String repoRoot, String branch) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
// Runs from INSIDE release()'s git-status shell-out: the architect slot is already
|
||||
// unbound by this point, but the pane has not stopped yet.
|
||||
duringWindow.set(resolverRef.get().resolve("127.0.0.1", 42, null));
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String repoRoot(String cwd) {
|
||||
return "/repo";
|
||||
}
|
||||
|
||||
@Override
|
||||
public Optional<String> snapshot(String worktreePath, String branch, String message) {
|
||||
return Optional.empty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public WipRefStats wipRefs(String repoRoot) {
|
||||
return new WipRefStats(0, 0L);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shareWithGroup(String repoRoot, String worktreePath) {
|
||||
}
|
||||
};
|
||||
SessionManager sessions = new SessionManager(workers, worktrees);
|
||||
|
||||
MemberRegistry members = new MemberRegistry(new FleetConfig.Fleet(Map.of(),
|
||||
Map.of("lead-designer", new FleetConfig.Slot("ltms-local")), Map.of(), Map.of(), null));
|
||||
sessions.setMemberLifecycle(members);
|
||||
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> FakeHerdr.WORKER_PID);
|
||||
CallerResolver resolver = CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
Map::of, members, sessions::spawnedMemberRole, Map::of);
|
||||
resolverRef.set(resolver);
|
||||
|
||||
MemberSession s = sessions.acquire("ltms-local", MemberRole.ARCHITECT, null, "/caller/proj", null,
|
||||
new WorktreeRequest("fleetd-702d", null));
|
||||
assertEquals("term_a", s.terminalId(), "sanity: the spawn resolved to the pinned pane");
|
||||
assertEquals(MemberRole.ARCHITECT, s.role(), "sanity: the slot bind succeeded");
|
||||
|
||||
Principal before = resolver.resolve("127.0.0.1", 42, null);
|
||||
assertEquals(Role.ARCHITECT, before.role(),
|
||||
"control: with the slot still bound, the pane must resolve as an architect — this "
|
||||
+ "is what proves the in-window assertion below is not passing against a "
|
||||
+ "slot that was never bound");
|
||||
assertEquals("lead-designer", before.name());
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertNotNull(duringWindow.get(), "the dirty check must have run and captured a resolve");
|
||||
assertEquals(Role.WORKER, duringWindow.get().role(),
|
||||
"inside the window the architect slot is already unbound, so the result must be "
|
||||
+ "WORKER — asserting role-equality with the live session here would tempt "
|
||||
+ "moving the unbind earlier or later, which would be wrong either way");
|
||||
assertEquals("term_a", duringWindow.get().terminal());
|
||||
}
|
||||
|
||||
/** A spawned architect in the roster resolves ARCHITECT, carrying its bound slot's name. */
|
||||
@Test
|
||||
void aSpawnedArchitectInTheRosterResolvesArchitectWithItsSlotName() {
|
||||
|
||||
@@ -1,10 +1,16 @@
|
||||
package dev.ltms.fleet.lead;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.ResilientAgentLaunch;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
@@ -64,6 +70,21 @@ class LeadLauncherTest {
|
||||
return new LeadLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), cfg);
|
||||
}
|
||||
|
||||
/** As {@link #launcher}, plus a fast no-op sleeper so a busy-retry test never real-sleeps. */
|
||||
private static LeadLauncher fastLauncher(FakeHerdr herdr, FleetConfig cfg) {
|
||||
return new LeadLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), cfg, () -> { });
|
||||
}
|
||||
|
||||
/** {@link #opusProfile()} with one argv element long enough to overflow the pane line limit. */
|
||||
private static FleetConfig.Profile hugeArgvProfile() {
|
||||
return new FleetConfig.Profile(
|
||||
"opus", null, "claude-opus-5", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "x".repeat(1500)), "tab", "fleet", null,
|
||||
"http://127.0.0.1:8765/mcp", null, null,
|
||||
null, null, null,
|
||||
Map.of("CLAUDE_CODE_AUTO_COMPACT_WINDOW", "300000"), null, null, true, null);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static List<String> startedArgs(FakeHerdr herdr) {
|
||||
return (List<String>) ((Map<String, Object>) herdr.lastCall("agent.start").params()).get("args");
|
||||
@@ -86,7 +107,11 @@ class LeadLauncherTest {
|
||||
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
|
||||
assertTrue(herdr.called("agent.start"), "a lead must actually be started");
|
||||
assertEquals("lead-opus", startedName(herdr));
|
||||
// fleetd #727: the name carries a per-process nonce and a per-start sequence number — the
|
||||
// same unique-naming scheme HerdrPeerLauncher uses for members — rather than the fixed
|
||||
// "lead-opus" a stale registry entry could block a legitimate relaunch under.
|
||||
assertTrue(startedName(herdr).matches("lead-opus-[0-9a-f]{6}-\\d+"),
|
||||
"name is lead-<name>-<nonce>-<seq>: " + startedName(herdr));
|
||||
}
|
||||
|
||||
/** The tab is labelled with the configured `tab:` so the scanner finds the lead on the next resolve. */
|
||||
@@ -436,4 +461,114 @@ class LeadLauncherTest {
|
||||
assertEquals(0, launcher(herdr, cfg).ensureLeads());
|
||||
assertTrue(herdr.calls.isEmpty(), "nothing declared ⇒ nothing scanned");
|
||||
}
|
||||
|
||||
// ── fleetd #727: the same three launch protections every member spawn gets ───────────────────
|
||||
|
||||
/**
|
||||
* A freshly created pane may not have redrawn its prompt yet, so herdr answers
|
||||
* {@code agent_pane_busy}. The lead launch must wait it out rather than fail on the first miss
|
||||
* — exactly the retry {@code HerdrPeerLauncher} already gives every member.
|
||||
*/
|
||||
@Test
|
||||
void aLeadLaunchRetriesWhileTheSeedShellBoots() {
|
||||
FakeHerdr herdr = new FakeHerdr().agentPaneBusyTimes(2);
|
||||
|
||||
assertEquals(1, fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
|
||||
"the lead must still start once the shell is ready");
|
||||
assertEquals(3, herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
|
||||
"two busy rejections, then the successful start");
|
||||
}
|
||||
|
||||
/**
|
||||
* The busy retry is bounded, not an infinite poll. If the pane never becomes ready the launch
|
||||
* must eventually give up and log the failure, not hang the daemon's reconcile loop forever —
|
||||
* proven here by a budget that would still be busy on attempt
|
||||
* {@value ResilientAgentLaunch#SHELL_READY_RETRIES} and a call count that stops exactly there.
|
||||
*/
|
||||
@Test
|
||||
void aLeadLaunchGivesUpAfterTheBoundedBusyBudgetRatherThanLoopingForever() {
|
||||
FakeHerdr herdr = new FakeHerdr().agentPaneBusyTimes(999);
|
||||
|
||||
assertEquals(0, fastLauncher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
|
||||
"a pane that never becomes ready must not be reported as a started lead");
|
||||
assertEquals(ResilientAgentLaunch.SHELL_READY_RETRIES,
|
||||
herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
|
||||
"the retry budget is bounded: it stops after exactly SHELL_READY_RETRIES attempts");
|
||||
}
|
||||
|
||||
/**
|
||||
* herdr refuses a duplicate agent {@code name} outright. A stale registry entry — a crashed
|
||||
* lead session, or a name the registry has not yet released — must not permanently block a
|
||||
* legitimate relaunch, so each retry attempt carries a fresh per-start name.
|
||||
*/
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
void aLeadLaunchRetriesUnderAFreshNameWhenTheOldNameIsStillTaken() {
|
||||
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(2);
|
||||
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
|
||||
"the lead must still start once a free name is found");
|
||||
List<String> names = herdr.calls.stream()
|
||||
.filter(c -> c.method().equals("agent.start"))
|
||||
.map(c -> ((Map<String, Object>) c.params()).get("name").toString())
|
||||
.toList();
|
||||
assertEquals(3, names.size(), "2 rejected + 1 success");
|
||||
assertEquals(3, Set.copyOf(names).size(), "each attempt must use a distinct name");
|
||||
}
|
||||
|
||||
/**
|
||||
* The name-collision retry is bounded too. If the name is taken on every attempt, the launch
|
||||
* must give up rather than keep minting new names forever — the per-start naming scheme means
|
||||
* every failed attempt was refused outright by herdr (no process exists under a name herdr
|
||||
* refused), so a bounded, exhausted retry never leaves a second process running: nothing is
|
||||
* running at all.
|
||||
*/
|
||||
@Test
|
||||
void aLeadLaunchNeverEndsUpWithASecondProcessWhenTheNameStaysTaken() {
|
||||
FakeHerdr herdr = new FakeHerdr().agentNameTakenTimes(999);
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
|
||||
"a name that is never free must not be reported as a started lead");
|
||||
assertEquals(ResilientAgentLaunch.NAME_RETRIES,
|
||||
herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count(),
|
||||
"the retry budget is bounded: it stops after exactly NAME_RETRIES attempts, "
|
||||
+ "never racing a duplicate into existence");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #220/#727: herdr types the launch command into the pane as one line, and a pty line
|
||||
* buffer holds only 1024 bytes — past that the tail is dropped with no error at all, and the
|
||||
* backend exits on a mangled argument. The lead launch must refuse an over-long command outright
|
||||
* rather than let it be typed and silently truncated.
|
||||
*/
|
||||
@Test
|
||||
void anOverlongLeadArgvIsRefusedRatherThanTypedAndTruncated() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(LeadLauncher.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
|
||||
int started;
|
||||
try {
|
||||
started = launcher(herdr, configWith(lead("opus", "lead: opus", 1), hugeArgvProfile()))
|
||||
.ensureLeads();
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertEquals(0, started, "an over-long command must never be reported as a started lead");
|
||||
assertFalse(herdr.called("agent.start"),
|
||||
"nothing may be started — a truncated command is worse than no spawn");
|
||||
String warn = appender.list.stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains("failed to launch"))
|
||||
.findFirst()
|
||||
.orElseThrow(() -> new AssertionError("expected a WARN naming the launch failure: "
|
||||
+ appender.list));
|
||||
assertTrue(warn.contains("1024"), "names the limit: " + warn);
|
||||
assertTrue(warn.contains("opus"), "names the profile: " + warn);
|
||||
assertTrue(warn.contains("x".repeat(60)), "names the culprit argument: " + warn);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -2436,4 +2436,160 @@ class SessionManagerTest {
|
||||
return "threw:" + e.getClass().getName() + ":" + e.getMessage();
|
||||
}
|
||||
}
|
||||
|
||||
// --- fleetd #702: spawnedMemberRole must still answer for a pane mid-teardown ----------------
|
||||
|
||||
/**
|
||||
* A {@link Worktrees} test double whose {@code hasUncommitted} runs an injected hook before
|
||||
* answering. This is what lets a test resolve a releasing pane's terminal from inside the
|
||||
* window {@link SessionManager#release} opens between removing the registry entry and
|
||||
* actually stopping the pane — the hook runs synchronously on the release call's own thread,
|
||||
* at the exact point {@code release} shells out to {@code git status}, so there is no sleep
|
||||
* and no race to land in it.
|
||||
*/
|
||||
private static final class HookedWorktrees implements Worktrees {
|
||||
private Runnable hook;
|
||||
private boolean dirty = false;
|
||||
|
||||
HookedWorktrees onHasUncommitted(Runnable hook) {
|
||||
this.hook = hook;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String add(String repoRoot, String branch, String baseRef) {
|
||||
return "/wt/" + branch.replace('/', '_');
|
||||
}
|
||||
|
||||
@Override
|
||||
public void remove(String repoRoot, String worktreePath) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void deleteBranch(String repoRoot, String branch) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
if (hook != null) {
|
||||
hook.run();
|
||||
}
|
||||
return dirty;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public String repoRoot(String cwd) {
|
||||
return "/repo";
|
||||
}
|
||||
|
||||
@Override
|
||||
public java.util.Optional<String> snapshot(String worktreePath, String branch, String message) {
|
||||
return java.util.Optional.empty();
|
||||
}
|
||||
|
||||
@Override
|
||||
public WipRefStats wipRefs(String repoRoot) {
|
||||
return new WipRefStats(0, 0L);
|
||||
}
|
||||
|
||||
@Override
|
||||
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void shareWithGroup(String repoRoot, String worktreePath) {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnedMemberRoleResolvesTheLiveRegisteredRole() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller", null);
|
||||
|
||||
assertEquals(s.role(), sessions.spawnedMemberRole(s.terminalId()),
|
||||
"a registered session resolves to its own role");
|
||||
assertNull(sessions.spawnedMemberRole("term_unknown"),
|
||||
"a terminal with no session at all resolves to null");
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnedMemberRoleIsNullOnceReleaseFullyCompletes() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr, new HookedWorktrees());
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("fleetd-702a", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertNull(sessions.spawnedMemberRole(s.terminalId()),
|
||||
"once release has fully finished, the terminal is neither registered nor releasing");
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnedMemberRoleStillAnswersBetweenTheRegistryRemovalAndThePaneStop() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
java.util.concurrent.atomic.AtomicReference<MemberRole> duringWindow = new java.util.concurrent.atomic.AtomicReference<>();
|
||||
HookedWorktrees worktrees = new HookedWorktrees();
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("fleetd-702b", null));
|
||||
worktrees.onHasUncommitted(() -> {
|
||||
// This runs from INSIDE release()'s git-status shell-out: the registry entry is
|
||||
// already gone, but the pane has not stopped yet — a real call landing inside the
|
||||
// exact window fleetd #702 reports, so no sleep and no race is needed to reach it.
|
||||
assertTrue(sessions.get(s.paneId()).isEmpty(),
|
||||
"sanity: the registry entry is already gone at this point");
|
||||
duringWindow.set(sessions.spawnedMemberRole(s.terminalId()));
|
||||
});
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(s.role(), duringWindow.get(),
|
||||
"spawnedMemberRole must still answer the live role while the pane is mid-teardown, "
|
||||
+ "not only while the session is still in the registry");
|
||||
assertNull(sessions.spawnedMemberRole(s.terminalId()),
|
||||
"and once release has fully finished, the window is closed too");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code Releasing.enter} and {@code Releasing.leave} are called directly here, never through
|
||||
* {@link SessionManager#release} or {@link SessionManager#spawnedMemberRole}, so a test naming
|
||||
* one of them exercises only that one — a regression in the other can never hide behind it.
|
||||
*/
|
||||
@Test
|
||||
void releasingLeaveStepsDownADepthGreaterThanOneInsteadOfRemovingIt() {
|
||||
SessionManager.Releasing depthTwo = new SessionManager.Releasing(2, "term_a", MemberRole.DEV);
|
||||
|
||||
SessionManager.Releasing afterLeave = depthTwo.leave();
|
||||
|
||||
assertNotNull(afterLeave,
|
||||
"depth 2 means another release of the SAME pane is still mid-teardown; leave() must "
|
||||
+ "step the depth down, never remove the marker outright — removing it here "
|
||||
+ "is what a plain Set would do, and would reopen the window the still-in-"
|
||||
+ "flight release is relying on staying closed");
|
||||
assertEquals(1, afterLeave.depth());
|
||||
assertEquals("term_a", afterLeave.terminalId());
|
||||
assertEquals(MemberRole.DEV, afterLeave.role());
|
||||
}
|
||||
|
||||
@Test
|
||||
void releasingEnterPreservesThePriorTerminalWhenTheOverlappingCallHasNoSessionOfItsOwn() {
|
||||
SessionManager.Releasing prior = new SessionManager.Releasing(1, "term_a", MemberRole.DEV);
|
||||
|
||||
SessionManager.Releasing afterEnter = SessionManager.Releasing.enter(prior, null);
|
||||
|
||||
assertEquals(2, afterEnter.depth(), "depth still increments whether or not this enter knows its session");
|
||||
assertEquals("term_a", afterEnter.terminalId(),
|
||||
"an overlapping release that finds the registry entry already gone has no session of "
|
||||
+ "its own to pass as known, and must not blank out the terminal the first "
|
||||
+ "call already recorded — that terminal is what spawnedMemberRole matches "
|
||||
+ "against");
|
||||
assertEquals(MemberRole.DEV, afterEnter.role());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user