Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 1fc9e85bf1 | |||
| 809b7d9b20 | |||
| 8cf7215d56 | |||
| 1ef93e57cc |
@@ -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 ----------------------------------------------------------------------
|
||||
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
|
||||
Reference in New Issue
Block a user