Compare commits

..

2 Commits

Author SHA1 Message Date
Dai Ha 743377d6cd fleetd #149 review round 2: make the trust-dialog seed atomic and lock-protected
CI / build (pull_request) Successful in 1m21s
CI / contract (pull_request) Successful in 1m25s
Files.writeString truncates the target in place before writing, so there
was a window where .claude.json could be observed empty or half-written
- exactly the shape of the incident this ticket already hit once, but
reachable in production too: a crash/kill mid-write, or two concurrent
claude-code spawns (normal here - several run in parallel routinely)
racing a naive read-modify-write and silently discarding one spawn's
entry.

Two independent fixes, each with its own dedicated test proving it (not
the other):

- ClaudeCodeLauncher.writeAtomically: serialise to a sibling temp file in
  the same directory, then Files.move with ATOMIC_MOVE + REPLACE_EXISTING,
  preserving the target's existing POSIX permissions (.claude.json ships
  0600). A reader now only ever observes the fully-old or fully-new file,
  never a torn one. Package-visible so a test can drive it directly.
- TRUST_JSON_LOCK: a process-wide lock around seedTrustDialog's whole
  read-modify-write, so two concurrent spawns for different cwds both
  keep their entry instead of the second write discarding the first.
  Sufficient because every spawn on this daemon runs in one JVM; it does
  NOT protect against a second daemon process or the operator's own live
  Claude Code writing at the same instant - writeAtomically covers that
  case instead.

Both fail soft, same as before: any I/O failure here must never block a
spawn.

Four new tests: a large (30-project) existing file survives without
collapsing (asserted on the restored key set, not just that the result
parses); two concurrent spawns for different cwds both keep their entry
(CountDownLatch-synchronised, not a sleep); existing 0600 permissions
survive the write; and a direct test of writeAtomically with a busy-poll
reader thread proving a concurrent reader never observes a torn file.

See PR body for the full mutation-testing table, including an
honest note on which of these tests the atomicity mutation actually
caught (not the one implied by the numbering in review) and why.
2026-09-03 11:38:10 +07:00
Dai Ha a89dcc9b7e fleetd #149: seed the workspace-trust entry before a claude-code spawn
CI / contract (pull_request) Successful in 51s
CI / build (pull_request) Successful in 1m43s
A claude-code member spawned into a fresh worktree hits an interactive,
un-timed workspace-trust prompt on its first start in a directory it has
never seen. It never reaches its first turn and never mounts the bridge.

Fix: ClaudeCodeLauncher.seedTrustDialog writes
projects.<cwd>.hasTrustDialogAccepted / hasCompletedProjectOnboarding into
the profile's configDir/.claude.json (or ~/.claude.json when configDir is
unset) BEFORE the herdr spawn call, additively (existing keys/projects are
preserved). Gated to isProvisionedWorktree(cwd) - a .git that is a regular
gitdir-pointer file, never a real checkout's .git directory - the same
signal writeIdeOverlay already used, now shared between both.

That gate is a fix for a real incident hit while building this: an
earlier ungated version ran against this file's own pre-existing tests
(configDir=null, no cwd -> falls back to the real user.dir and
~/.claude.json) and corrupted the operator's actual ~/.claude.json down
to a single entry during a mutation-testing run. See PR body for the
full incident report.

FakeHerdr gained onAgentStart(Runnable) so a test can assert the seed
is on disk at the exact instant herdr's agent.start call is reached -
i.e. strictly before the peer process itself would start.
2026-09-03 11:23:06 +07:00
10 changed files with 674 additions and 217 deletions
@@ -390,12 +390,7 @@ public final class Fleetd {
// now that `sessions` exists to resolve target -> session -> profile.
exhaustionSinkRef.set(exhaustionSink);
AgentControl agents = router.memberAgents();
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink,
target -> sessions.roster().stream()
.filter(session -> target.equals(session.terminalId()))
.findFirst()
.map(session -> new CompletionResolver.WorktreeBranch(session.worktree(), session.branch()))
.orElse(null));
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink);
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
@@ -14,7 +14,6 @@ import java.util.TreeSet;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.LongSupplier;
import java.util.function.Function;
import java.util.regex.Pattern;
/**
@@ -86,17 +85,6 @@ public final class CompletionResolver implements TurnListener {
private static final String CLIPPED_PANE_TAIL_MARKER =
"[Pane tail clipped: member did not call fleet_reply.]";
/** A pane echo must be this large before it can replace a completion report. */
static final int ECHO_MIN_CHARS = 400;
/** Normalised TUI chrome may add this many characters to an otherwise echoed brief. */
static final int MAX_ECHO_EXCESS_CHARS = 160;
/** The explicit result returned instead of a lead's echoed injected brief. */
public static final String NO_REPORT_PREFIX = "[no report — the member ended its turn without fleet_reply, "
+ "and the pane still shows the injected brief. Nothing was produced on the pane. Check the "
+ "member's worktree and branch for committed work before re-delegating.";
private final AgentControl agents;
private final Rendezvous rendezvous;
private final ExhaustedPatternLookup exhaustedPatterns;
@@ -104,11 +92,6 @@ public final class CompletionResolver implements TurnListener {
private final BackendErrorPatternLookup backendErrorPatterns;
private final BackendErrorSink backendErrorSink;
private final LongSupplier nowNanos;
private final Function<String, WorktreeBranch> worktreeBranches;
/** Known member location, used only to guide a lead after an echoed brief. */
public record WorktreeBranch(String worktree, String branch) {
}
/**
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its
@@ -125,8 +108,7 @@ public final class CompletionResolver implements TurnListener {
* the other half of the {@link #MIN_TURN_NANOS} floor check, compared against a fresh reading at
* resolution time.
*/
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos,
String injectedText) {
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos) {
/**
* Convenience for tests exercising scrape/suppression logic that don't care about turn
@@ -134,11 +116,7 @@ public final class CompletionResolver implements TurnListener {
* Not used by production code — {@link #captureBaseline} always records a real reading.
*/
InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline) {
this(waiter, baseline, Long.MIN_VALUE / 2, null);
}
InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos) {
this(waiter, baseline, deliveredAtNanos, null);
this(waiter, baseline, Long.MIN_VALUE / 2);
}
}
@@ -161,18 +139,11 @@ public final class CompletionResolver implements TurnListener {
* (fleetd#201 Unit 5).
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink) {
ExhaustionSink exhaustionSink) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink,
BackendErrorPatternLookup.legacy(), BackendErrorSink.none(), System::nanoTime);
}
/** Production constructor with a lookup for the member worktree and branch. */
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, Function<String, WorktreeBranch> worktreeBranches) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink,
BackendErrorPatternLookup.legacy(), BackendErrorSink.none(), System::nanoTime, worktreeBranches);
}
/**
* Transition constructor (fleetd#201 Unit 1): same legacy backend-error defaults as the 4-arg
* constructor above, but with the injectable clock. Kept so existing fleetd#164 timing tests
@@ -201,7 +172,7 @@ public final class CompletionResolver implements TurnListener {
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
BackendErrorSink backendErrorSink) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink, backendErrorPatterns, backendErrorSink,
System::nanoTime, _ -> null);
System::nanoTime);
}
/**
@@ -214,17 +185,8 @@ public final class CompletionResolver implements TurnListener {
* package.
*/
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
BackendErrorSink backendErrorSink, LongSupplier nowNanos) {
this(agents, rendezvous, exhaustedPatterns, exhaustionSink, backendErrorPatterns, backendErrorSink,
nowNanos, _ -> null);
}
/** Full constructor with injectable clock and member location lookup. */
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
BackendErrorSink backendErrorSink, LongSupplier nowNanos,
Function<String, WorktreeBranch> worktreeBranches) {
ExhaustionSink exhaustionSink, BackendErrorPatternLookup backendErrorPatterns,
BackendErrorSink backendErrorSink, LongSupplier nowNanos) {
this.agents = agents;
this.rendezvous = rendezvous;
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
@@ -232,7 +194,6 @@ public final class CompletionResolver implements TurnListener {
this.backendErrorPatterns = Objects.requireNonNull(backendErrorPatterns, "backendErrorPatterns");
this.backendErrorSink = Objects.requireNonNull(backendErrorSink, "backendErrorSink");
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
this.worktreeBranches = Objects.requireNonNull(worktreeBranches, "worktreeBranches");
}
@Override
@@ -262,7 +223,7 @@ public final class CompletionResolver implements TurnListener {
baseline = null; // fail open: no baseline ⇒ no suppression
log.debug("delivery baseline for {} failed: {}", target, e.getMessage());
}
inFlight.put(target, new InFlight(waiter, baseline, nowNanos.getAsLong(), token.injectedText()));
inFlight.put(target, new InFlight(waiter, baseline, nowNanos.getAsLong()));
}
/** The turn currently baselined for {@code target}, or {@code null} — a test hook for the captureBaseline path. */
@@ -408,9 +369,6 @@ public final class CompletionResolver implements TurnListener {
return;
}
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
if (echoesInjectedBrief(tail, turn.injectedText())) {
completion = noReportMessage(target) + (clipped ? "\n" + CLIPPED_PANE_TAIL_MARKER : "");
}
if (rendezvous.resolveCompletion(waiter, completion)) {
inFlight.remove(target, turn);
if (clipped) {
@@ -423,40 +381,6 @@ public final class CompletionResolver implements TurnListener {
}
}
/**
* A full echoed brief is at least 400 normalised characters. A scrape that contains the brief may
* add no more than 160 normalised characters of TUI chrome. This accepts harmless status text, but
* preserves a real report that restates the full brief before adding substantive content.
*/
static boolean echoesInjectedBrief(String scrape, String injectedText) {
String normalScrape = normalize(scrape);
String normalInjected = normalize(injectedText);
if (normalScrape.length() < ECHO_MIN_CHARS || normalInjected.length() < ECHO_MIN_CHARS) {
return false;
}
if (normalInjected.contains(normalScrape)) {
return true;
}
return normalScrape.contains(normalInjected)
&& normalScrape.length() - normalInjected.length() <= MAX_ECHO_EXCESS_CHARS;
}
private static String normalize(String text) {
return text == null ? "" : text.toLowerCase().replaceAll("[^a-z0-9]+", "");
}
private String noReportMessage(String target) {
WorktreeBranch location = worktreeBranches.apply(target);
if (location == null || (location.worktree() == null && location.branch() == null)) {
return NO_REPORT_PREFIX + "]";
}
String locationText = location.worktree() == null ? "" : " worktree=" + location.worktree();
if (location.branch() != null) {
locationText += " branch=" + location.branch();
}
return NO_REPORT_PREFIX + locationText + "]";
}
/**
* fleetd#211: the raw-scrape fallback classification, run only when {@link #lastAssistantBlock}
* found nothing usable (see the call site in {@link #resolve}). Mirrors the two classifications
@@ -10,7 +10,6 @@ import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.msg.LeadChannel;
import dev.ltms.fleet.msg.LeadMessage;
@@ -544,9 +543,8 @@ public final class FleetMcp {
case REPLIED -> text(r.text());
// The worker's turn finished but it never called fleet_reply — hand back the scraped
// transcript tail, flagged so the primary knows it isn't a structured reply.
case COMPLETED_UNREPLIED -> text(r.text().startsWith(CompletionResolver.NO_REPORT_PREFIX)
? r.text()
: "[worker finished without a structured fleet_reply — transcript tail follows]\n" + r.text());
case COMPLETED_UNREPLIED -> text(
"[worker finished without a structured fleet_reply — transcript tail follows]\n" + r.text());
// The worker ran the turn then wedged (CB-109) — surface the error context.
case WORKER_FAILED -> text("[worker failed — turn ended in an unrecoverable state]\n" + r.text());
// The backend refused on a subscription usage limit (CB-578 stage A) — the worker's
@@ -661,7 +659,6 @@ public final class FleetMcp {
}
return switch (v.phase()) {
case DONE -> text(v.replySource() != null && v.replySource().equals("transcript")
&& !v.reply().startsWith(CompletionResolver.NO_REPORT_PREFIX)
? "[done — worker finished without a structured fleet_reply; transcript tail follows]\n" + v.reply()
: v.reply());
case PENDING -> text("[pending — " + v.detail() + "]");
@@ -1,5 +1,8 @@
package dev.ltms.fleet.member;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.Agent;
@@ -14,6 +17,8 @@ import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.nio.file.attribute.PosixFileAttributeView;
import java.util.EnumSet;
import java.util.List;
import java.util.Map;
@@ -47,6 +52,21 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
private static final Logger log = LoggerFactory.getLogger(ClaudeCodeLauncher.class);
/** JSON codec for the additive workspace-trust seed (fleetd #149) — Jackson's default settings. */
private static final ObjectMapper TRUST_JSON = new ObjectMapper();
/**
* Serialises every {@link #seedTrustDialog} read-modify-write for the whole daemon process.
* Two claude-code spawns starting at once are normal (fleetd runs several members in parallel
* routinely) and both would otherwise read the same {@code .claude.json}, add their own entry
* to their own in-memory copy, and write — the second write wins and the first spawn's trust
* entry silently disappears. A single process-wide lock is enough because every spawn on this
* daemon runs in this one JVM; it does not protect against a second daemon process or the
* operator's own Claude Code process writing at the same instant, which {@link #writeAtomically}
* covers instead (each writer only ever sees a fully-old or fully-new file, never a torn one).
*/
private static final Object TRUST_JSON_LOCK = new Object();
private final SubscriptionGuard guard;
/**
@@ -263,6 +283,10 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
putIfPresent(workerEnv, "ANTHROPIC_MODEL", cfg.model());
putIfPresent(workerEnv, "CLAUDE_CONFIG_DIR", cfg.configDir());
applyGitToken(workerEnv, cfg);
// fleetd #149: seed the workspace-trust entry BEFORE this spawn ever reaches herdr — see
// seedTrustDialog for why, and isProvisionedWorktree for why this is gated to a worktree
// fleetd itself provisioned (never a real checkout, never an un-configured fallback cwd).
seedTrustDialog(cfg.configDir(), spec.cwd());
// CB-547a: Claude Code can MINT its own session id, so fleetd chooses it — a fresh spawn
// gets a UUID we pass as --session-id and return from agentSessionId(), so the resume
@@ -412,12 +436,12 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
*/
private static void writeIdeOverlay(String cwd, String projectPath) {
try {
Path dotGit = Path.of(cwd, ".git");
if (!Files.isRegularFile(dotGit)) {
if (!isProvisionedWorktree(cwd)) {
// Not a provisioned worktree (primary's real checkout has a .git directory, or the
// cwd is not a repo at all). Never write into it.
return;
}
Path dotGit = Path.of(cwd, ".git");
// The overlay FILE lives at the worktree root (claude-code's cwd), but its CONTENT pins
// project_path to the module dir the IDE opened (projectPath), not the worktree root.
Files.writeString(Path.of(cwd, "CLAUDE.local.md"), PeerLauncher.ideOverlayText(projectPath));
@@ -447,6 +471,189 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
}
}
/**
* fleetd #149: pre-seed the workspace-trust entry for {@code cwd} in Claude Code's own config
* file, BEFORE this launch ever reaches herdr (called from {@link #buildLaunch}, which always
* runs before the base class starts the process). Claude Code asks an interactive, un-timed
* "Is this a project you created or one you trust?" the first time it starts in a directory it
* has not seen before, and every {@code worktree: true} spawn lands in a brand-new directory —
* so without this seed the member sits on that dialog forever, never mounts the bridge MCP, and
* never calls {@code fleet_reply}. herdr reports it as healthy the whole time ({@code
* agent_status: blocked}, {@code interactive_ready: true}), so nothing else catches it. Measured
* live on fleet01 2026-08-23: an unseeded fresh cwd sat on the dialog indefinitely; a seeded one
* reached {@code idle} clean.
*
* <p>This is not a new grant — the operator already trusted this repo by configuring the
* profile against it, and a worktree is a checkout of that same repo.
*
* <p>The key Claude Code reads is per-project, in {@code .claude.json}: {@code
* projects.<cwd>.hasTrustDialogAccepted}. The file lives at {@code <configDir>/.claude.json}
* when the profile sets {@code CLAUDE_CONFIG_DIR} (mirrors this launcher's own env var above),
* else the default {@code ~/.claude.json} — the same file Claude Code itself would read either
* way, so this seeds exactly what the spawned peer is about to open.
*
* <p><b>Additive, not a rewrite.</b> {@code .claude.json} is large (tens of KB, dozens of
* projects) and Claude Code itself rewrites it while running, so this reads the file as a JSON
* tree (missing or unreadable → treated as an empty object) and changes only
* {@code projects.<cwd>.hasTrustDialogAccepted} / {@code .hasCompletedProjectOnboarding} —
* every other top-level key and every other project entry is written back untouched. Only the
* one project entry for {@code cwd} is replaced/created; an existing entry for a DIFFERENT cwd
* (or the operator's own project history) is never touched.
*
* <p>Best-effort, like {@link #writeIdeOverlay}: a failure here (unwritable configDir, a
* corrupt existing file, …) must never fail the spawn — it is logged at debug and swallowed. A
* peer that starts without the seed still starts; it just may hit the dialog fleetd #149
* describes.
*
* <p><b>Gated to a provisioned worktree</b> ({@link #isProvisionedWorktree}) — see that
* method's javadoc for the incident that made this gate mandatory, not optional: this must
* never run against a real checkout or an un-configured fallback cwd, only the exact
* always-fresh-directory population fleetd #149 describes.
*
* <p><b>Atomic and lock-protected.</b> {@code .claude.json} is a live file — Claude Code itself
* rewrites it while running, and this daemon routinely spawns several members at once, each
* calling this method for its own cwd. Every write goes through {@link #writeAtomically} (a
* sibling-temp-file + {@code ATOMIC_MOVE}, never a truncate-in-place) so a crash mid-write or a
* concurrent reader never observes a half-written file, and through {@link #TRUST_JSON_LOCK} so
* two concurrent spawns' entries both survive instead of the second write silently discarding
* the first. Both exist because of a real incident: see {@link #isProvisionedWorktree}'s javadoc
* and {@link #writeAtomically}'s javadoc.
*
* @param configDir the profile's {@code CLAUDE_CONFIG_DIR} ({@code cfg.configDir()}), or
* {@code null}/blank to target the default {@code ~/.claude.json}
* @param cwd the spawn's resolved working directory — the exact key Claude Code will look
* up for itself once it starts there
*/
private static void seedTrustDialog(String configDir, String cwd) {
if (!isProvisionedWorktree(cwd)) {
return;
}
Path target = (configDir == null || configDir.isBlank())
? Path.of(System.getProperty("user.home"), ".claude.json")
: Path.of(configDir, ".claude.json");
synchronized (TRUST_JSON_LOCK) {
try {
if (target.getParent() != null) {
Files.createDirectories(target.getParent());
}
ObjectNode root = null;
if (Files.isRegularFile(target)) {
JsonNode existing = TRUST_JSON.readTree(target.toFile());
if (existing instanceof ObjectNode existingObject) {
root = existingObject;
}
}
if (root == null) {
root = TRUST_JSON.createObjectNode();
}
JsonNode projectsNode = root.get("projects");
ObjectNode projects = projectsNode instanceof ObjectNode projectsObject
? projectsObject : TRUST_JSON.createObjectNode();
if (!(projectsNode instanceof ObjectNode)) {
root.set("projects", projects);
}
JsonNode projectNode = projects.get(cwd);
ObjectNode project = projectNode instanceof ObjectNode projectObject
? projectObject : TRUST_JSON.createObjectNode();
if (!(projectNode instanceof ObjectNode)) {
projects.set(cwd, project);
}
project.put("hasTrustDialogAccepted", true);
project.put("hasCompletedProjectOnboarding", true);
writeAtomically(target, TRUST_JSON.writerWithDefaultPrettyPrinter().writeValueAsString(root));
} catch (Exception e) {
log.debug("cannot seed workspace-trust entry for cwd '{}' into '{}'", cwd, target, e);
}
}
}
/**
* Write {@code content} to {@code target} atomically: serialise to a sibling temp file in the
* <strong>same directory</strong> as {@code target} (an atomic move is only guaranteed within
* one filesystem — a different directory could mean a different filesystem), then
* {@link StandardCopyOption#ATOMIC_MOVE} it into place. A reader — Claude Code itself, or
* another {@code seedTrustDialog} call — only ever observes the fully-old file or the
* fully-new one, never a truncated or half-written one.
*
* <p><b>fleetd #149 incident.</b> The original implementation used
* {@code Files.writeString(target, content)} directly, which truncates {@code target} in place
* before writing the replacement bytes. Combined with an ungated {@code cwd} (see
* {@link #isProvisionedWorktree}'s javadoc), a mutation-testing run hit that truncation window
* against the operator's real {@code ~/.claude.json} and left it at 178 bytes. The gate closes
* <em>which file</em> this can ever target; this closes <em>how</em> the target is written, so
* that even a legitimate write against a real, live, concurrently-read {@code .claude.json}
* cannot leave it observably empty or partial.
*
* <p>Preserves {@code target}'s existing POSIX permissions (Claude Code ships {@code
* .claude.json} as {@code 0600}) when the filesystem reports them; a freshly created temp file
* already defaults to owner-only permissions on a POSIX filesystem, so a first-ever write (no
* existing {@code target}) is no less private without this. On a non-POSIX filesystem (e.g.
* Windows) the permission copy is a silent no-op rather than a failure.
*
* <p>Package-visible (not {@code private}) so a test can drive it directly with a concurrent
* reader thread and prove the torn-file property this method exists for — the LOCK in
* {@link #seedTrustDialog} already fully serialises every call this launcher itself makes, so a
* test that only ever goes through {@code seedTrustDialog}/{@code spawn()} could never observe
* a torn file regardless of whether this method is atomic; it would be proving the lock, not
* this method. Atomicity's actual job is protecting against a writer the lock cannot reach at
* all — a second daemon process, or the operator's own live Claude Code — so the test for it
* has to reach this method on its own.
*/
static void writeAtomically(Path target, String content) throws IOException {
Path parent = target.getParent();
Path tmp = Files.createTempFile(parent, target.getFileName() + ".", ".tmp");
try {
Files.writeString(tmp, content);
copyPosixPermissionsIfPresent(target, tmp);
Files.move(tmp, target, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING);
} catch (IOException e) {
Files.deleteIfExists(tmp);
throw e;
}
}
/** Copy {@code target}'s POSIX permissions onto {@code tmp}, or no-op where either is unsupported. */
private static void copyPosixPermissionsIfPresent(Path target, Path tmp) {
try {
if (!Files.isRegularFile(target)) {
return; // nothing to inherit from — first-ever write, temp file's own default stands
}
PosixFileAttributeView view = Files.getFileAttributeView(target, PosixFileAttributeView.class);
if (view == null) {
return; // non-POSIX filesystem — nothing this JVM can read/set here
}
Files.setPosixFilePermissions(tmp, Files.getPosixFilePermissions(target));
} catch (IOException e) {
log.debug("cannot preserve permissions of '{}' onto its replacement", target, e);
}
}
/**
* Whether {@code cwd} is a fleetd-provisioned git worktree — signalled the same way
* {@link #writeIdeOverlay} already gates on: a {@code .git} that is a <strong>regular
* file</strong> holding a {@code gitdir:} pointer, as opposed to a real checkout's {@code .git}
* <strong>directory</strong>. {@code null}/blank never qualifies.
*
* <p>Shared by every write that must land only in a worktree fleetd itself created for a
* member — never in a real checkout, an arbitrary configured directory, or (see the incident
* below) the daemon's own fallback cwd.
*
* <p><b>fleetd #149 incident.</b> {@link #seedTrustDialog} originally ran unconditionally on
* any non-blank {@code cwd}. Most of this launcher's OWN tests spawn a profile with no
* {@code cwd} configured, so the base class's {@code resolveCwd} falls through to the real
* {@code user.dir} — and with no {@code configDir} either (also the common case in this
* file's fixtures), the seed's target falls through the same way to the real
* {@code ~/.claude.json}. Running this repo's own test suite corrupted the operator's actual
* config file (it shrank from ~72 KB to a single seeded entry) the first time a mutation
* happened to make the write non-additive. Gating both cwd-targeted writes on "this is a
* worktree fleetd provisioned" — exactly the population fleetd #149 describes
* ({@code worktree: true} always lands in a brand-new directory) — makes that class of write
* impossible against a real checkout or an untouched fallback cwd, in production or in tests.
*/
private static boolean isProvisionedWorktree(String cwd) {
return cwd != null && !cwd.isBlank() && Files.isRegularFile(Path.of(cwd, ".git"));
}
/** {@code s}, or {@code null} when {@code s} is null/blank — the charter-presence test used above. */
private static String nonBlank(String s) {
return (s == null || s.isBlank()) ? null : s;
@@ -746,7 +746,7 @@ public final class MessageService {
if (task != null) {
asyncTasksByWaiter.put(reply, task);
}
TurnToken token = new TurnToken(target, reply, content);
TurnToken token = new TurnToken(target, reply);
// The send has won the lock; the accepted-delivery hook records delegator ownership
// here (CB-548). It runs BEFORE enqueue so a throwing hook — onAccepted is now a
// public callback — fails the send without queuing a message that would orphan.
@@ -9,19 +9,12 @@ import java.util.concurrent.CompletableFuture;
public final class TurnToken {
private final String target;
private final CompletableFuture<Rendezvous.Resolution> waiter;
private final String injectedText;
public TurnToken(String target, CompletableFuture<Rendezvous.Resolution> waiter) {
this(target, waiter, null);
}
public TurnToken(String target, CompletableFuture<Rendezvous.Resolution> waiter, String injectedText) {
this.target = target;
this.waiter = waiter;
this.injectedText = injectedText;
}
public String target() { return target; }
public CompletableFuture<Rendezvous.Resolution> waiter() { return waiter; }
public String injectedText() { return injectedText; }
}
@@ -49,6 +49,7 @@ public final class FakeHerdr implements HerdrClient {
private int pinnedStarts = 0; // how many upcoming agent.start calls report a fixed pane
private String pinnedStartTerminal;
private String pinnedStartPane;
private Runnable onAgentStart; // fires the instant agent.start is called — see onAgentStart(Runnable)
private volatile int agentGetOkCalls = Integer.MAX_VALUE; // how many agent.get calls succeed first
private volatile String agentGetFailCode = null; // error code every agent.get call after that reports
@@ -152,6 +153,18 @@ public final class FakeHerdr implements HerdrClient {
return this;
}
/**
* Run {@code hook} synchronously the instant an {@code agent.start} call reaches this fake —
* i.e. the instant the peer PROCESS would start against a real herdr daemon. A test uses this
* to assert something is already true at that exact point (rather than merely true once
* {@code spawn()} returns), e.g. fleetd #149's trust-dialog seed having already been written to
* disk before the process herdr would launch ever starts.
*/
public FakeHerdr onAgentStart(Runnable hook) {
this.onAgentStart = hook;
return this;
}
/**
* Seed a named agent into {@code agent.list} (e.g. an orphaned worker for CB-117 reaper tests).
@@ -250,6 +263,9 @@ public final class FakeHerdr implements HerdrClient {
case "agent.read" -> mapper.readTree(mapper.writeValueAsString(
java.util.Map.of("type", "agent_read", "read", java.util.Map.of("text", readText))));
case "agent.start" -> {
if (onAgentStart != null) {
onAgentStart.run();
}
// Protocol 19: kind and pane_id are required — reject like the real daemon.
java.util.Map<?, ?> p = params instanceof java.util.Map<?, ?> m ? m : java.util.Map.of();
for (String required : new String[]{"kind", "pane_id"}) {
@@ -196,26 +196,6 @@ class CompletionResolverTest {
assertEquals("complete report", waiter.getNow(null).text());
}
@Test
void suppressesABareEchoWithOnlyTuiChrome() {
String injected = "Load the implementer skill. You own fleetd #999. ".repeat(12);
String scrape = injected + "\nDev auto - GPT-5.6 Terra OpenAI";
assertTrue(CompletionResolver.echoesInjectedBrief(scrape, injected));
}
@Test
void pinsTheMaximumTuiChromeExcess() {
String injected = "a".repeat(CompletionResolver.ECHO_MIN_CHARS);
String underMargin = injected + "b".repeat(CompletionResolver.MAX_ECHO_EXCESS_CHARS);
String overMargin = injected + "b".repeat(CompletionResolver.MAX_ECHO_EXCESS_CHARS + 1);
assertTrue(CompletionResolver.echoesInjectedBrief(underMargin, injected),
"the configured excess itself remains an echoed brief");
assertFalse(CompletionResolver.echoesInjectedBrief(overMargin, injected),
"one character beyond the excess must preserve the scrape as a real report");
}
@Test
void resolvesSynchronouslyBeforePostTurnContextClearing() {
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
@@ -4,6 +4,9 @@ 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 com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.GuardException;
import dev.ltms.fleet.guard.SubscriptionGuard;
@@ -23,11 +26,16 @@ import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.attribute.PosixFileAttributeView;
import java.nio.file.attribute.PosixFilePermissions;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
import java.util.function.Supplier;
@@ -2098,4 +2106,433 @@ class ClaudeCodeLauncherTest {
assertTrue(charterFile.getFileName().toString().startsWith("fleetd-role-charter-"),
"same file-naming scheme as before this fix (no wrapping directory): " + charterFile);
}
// --- fleetd #149: workspace-trust dialog seed ------------------------------------------------
//
// Claude Code asks an interactive, un-timed "Is this a project you created or one you trust?"
// the first time it starts in a directory it has not seen. Every worktree: true spawn lands in
// a brand-new directory, so without a seed the member sits on that dialog forever — herdr still
// reports it healthy (agent_status: blocked, interactive_ready: true) — and never mounts the
// bridge MCP or calls fleet_reply. These tests start the REAL launcher (only the herdr transport
// is faked) so the seed is proven to run inside buildLaunch, before agent.start (the
// process-starting call) ever fires — a test that only checked the JSON writer in isolation
// would prove nothing about whether the launcher actually calls it at the right time.
//
// INCIDENT: the seed originally ran on ANY non-blank cwd. Running this file's own test suite —
// most of whose fixtures spawn with no cwd/configDir set, so both fall back to the real
// user.dir / ~/.claude.json — corrupted the operator's actual ~/.claude.json (it shrank from
// ~72 KB to a single seeded entry) the first time a manual mutation run made the write
// non-additive. The fix restricts the seed to isProvisionedWorktree(cwd) — a real .git FILE
// (not directory) — exactly the writeIdeOverlay gate already used for the same category of
// risk. Every test below marks its own worktree fixture with that .git file, and
// seedTrustDialogNeverWritesWhenCwdIsNotAProvisionedWorktree is the regression test for the
// incident itself.
private FleetConfig.Profile trustProfile(String configDir, String cwd) {
return new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", configDir, "FLEETD_WORKER_TOKEN",
List.of("claude"), "tab", "fleetd-workers", "w #{n}", "http://127.0.0.1:8765/mcp",
cwd, null);
}
/**
* Give {@code dir} the exact signature {@link ClaudeCodeLauncher#isProvisionedWorktree} (and
* {@code writeIdeOverlay} before it) checks for: a {@code .git} REGULAR FILE, never a
* directory. The content is never parsed by the trust seed, so any {@code gitdir:} pointer is
* fine.
*/
private static void markAsProvisionedWorktree(Path dir) throws IOException {
Files.writeString(dir.resolve(".git"), "gitdir: /tmp/not-a-real-gitdir");
}
@Test
void seedsWorkspaceTrustForTheCwdBeforeTheProcessStarts(
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
FakeHerdr herdr = new FakeHerdr();
Path claudeJson = configDir.resolve(".claude.json");
AtomicReference<Boolean> seededBeforeStart = new AtomicReference<>(false);
herdr.onAgentStart(() -> {
try {
if (!Files.exists(claudeJson)) {
return;
}
JsonNode root = new ObjectMapper().readTree(claudeJson.toFile());
JsonNode project = root.path("projects").path(worktree.toString());
seededBeforeStart.set(project.path("hasTrustDialogAccepted").asBoolean(false)
&& project.path("hasCompletedProjectOnboarding").asBoolean(false));
} catch (IOException e) {
seededBeforeStart.set(false);
}
});
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
assertTrue(herdr.called("agent.start"), "the hook must actually have fired during the spawn");
assertEquals(Boolean.TRUE, seededBeforeStart.get(),
"the trust entry for the cwd must already exist at the instant agent.start (the "
+ "process-starting herdr call) fires — not merely once spawn() returns");
}
@Test
void seedTrustDialogWritesBothTrustFlagsForTheResolvedCwd(
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
Path claudeJson = configDir.resolve(".claude.json");
assertTrue(Files.exists(claudeJson), "seeded into <configDir>/.claude.json");
JsonNode project = new ObjectMapper().readTree(claudeJson.toFile())
.path("projects").path(worktree.toString());
assertTrue(project.path("hasTrustDialogAccepted").asBoolean(false));
assertTrue(project.path("hasCompletedProjectOnboarding").asBoolean(false));
}
/**
* Criterion 3: seeding is additive. An existing {@code .claude.json} carries the operator's own
* project history and unrelated top-level settings — the seed must change only
* {@code projects.<cwd>} for THIS cwd and leave everything else, including a different project's
* own unrelated data, exactly as it was.
*/
@Test
void seedTrustDialogIsAdditiveAndPreservesUnknownKeysAndOtherProjects(
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
Path claudeJson = configDir.resolve(".claude.json");
Files.writeString(claudeJson, """
{
"numStartups": 42,
"oauthAccount": {"emailAddress": "operator@example.com"},
"projects": {
"/some/other/project": {
"hasTrustDialogAccepted": true,
"mcpServers": {"foo": {"type": "stdio", "command": "foo-mcp"}}
}
}
}
""");
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
JsonNode root = new ObjectMapper().readTree(claudeJson.toFile());
assertEquals(42, root.path("numStartups").asInt(), "unrelated top-level key survives untouched");
assertEquals("operator@example.com", root.path("oauthAccount").path("emailAddress").asText(),
"an unrelated nested top-level key survives untouched");
JsonNode other = root.path("projects").path("/some/other/project");
assertTrue(other.path("hasTrustDialogAccepted").asBoolean(false),
"a different project's own trust entry survives");
assertEquals("foo-mcp", other.path("mcpServers").path("foo").path("command").asText(),
"a different project's own unrelated nested data survives");
JsonNode mine = root.path("projects").path(worktree.toString());
assertTrue(mine.path("hasTrustDialogAccepted").asBoolean(false));
assertTrue(mine.path("hasCompletedProjectOnboarding").asBoolean(false));
}
/**
* Criterion 2: where the profile sets no {@code configDir}, the seed goes to the default
* {@code ~/.claude.json}. {@code user.home} is redirected to a {@code @TempDir} for the
* duration of this test and restored in a {@code finally} — the real operator {@code
* ~/.claude.json} must never be touched by a test.
*/
@Test
void seedTrustDialogTargetsDefaultClaudeJsonWhenConfigDirIsUnset(
@TempDir Path fakeHome, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
String originalHome = System.getProperty("user.home");
System.setProperty("user.home", fakeHome.toString());
try {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(null, worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
Path claudeJson = fakeHome.resolve(".claude.json");
assertTrue(Files.exists(claudeJson),
"no configDir set — the default target is ~/.claude.json, here the redirected fake home");
JsonNode project = new ObjectMapper().readTree(claudeJson.toFile())
.path("projects").path(worktree.toString());
assertTrue(project.path("hasTrustDialogAccepted").asBoolean(false));
} finally {
System.setProperty("user.home", originalHome);
}
}
/**
* Regression test for the fleetd #149 incident itself: a cwd that is NOT a fleetd-provisioned
* worktree (no {@code .git} FILE — the exact shape a real checkout, or an un-configured
* fallback cwd, has) must never be written to, however {@code configDir} is set. This is the
* fix for the exact defect that corrupted the operator's real {@code ~/.claude.json}.
*/
@Test
void seedTrustDialogNeverWritesWhenCwdIsNotAProvisionedWorktree(
@TempDir Path configDir, @TempDir Path plainCwd) {
// plainCwd deliberately carries NO .git file — the same shape a real checkout's cwd
// fallback has.
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), plainCwd.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
assertFalse(Files.exists(configDir.resolve(".claude.json")),
"a non-worktree cwd must never get a .claude.json written for it — this is the fix "
+ "for the incident where writing unconditionally corrupted the operator's own "
+ "real ~/.claude.json via this file's own no-cwd/no-configDir test fixtures");
}
/**
* Same regression, for the default (no {@code configDir}) path — the exact combination (no
* {@code configDir}, no worktree-shaped {@code cwd}) that hit the operator's real
* {@code ~/.claude.json} during the incident. {@code user.home} is still redirected to a
* {@code @TempDir} out of caution, so even a reintroduced bug here cannot touch the real file.
*/
@Test
void seedTrustDialogNeverWritesToDefaultHomeWhenCwdIsNotAProvisionedWorktree(
@TempDir Path fakeHome, @TempDir Path plainCwd) {
String originalHome = System.getProperty("user.home");
System.setProperty("user.home", fakeHome.toString());
try {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(null, plainCwd.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
assertFalse(Files.exists(fakeHome.resolve(".claude.json")),
"the exact incident combination — no configDir, non-worktree cwd — must never "
+ "write, even to the (redirected) default ~/.claude.json");
} finally {
System.setProperty("user.home", originalHome);
}
}
// --- fleetd #149 review round 2: the write must be atomic and lock-protected -----------------
//
// The first round fixed WHICH cwd this can ever target. This round fixes HOW the target is
// written: .claude.json is large (tens of KB, dozens of projects on a real host), Claude Code
// itself rewrites it while running, and this daemon spawns several members in parallel as a
// matter of routine. The pre-fix Files.writeString(target, content) truncates target in place
// before writing the replacement bytes — a crash mid-write, or another writer's read landing in
// that window, loses data. Two more failure modes follow directly: (a) a crash/kill mid-write
// leaves target truncated, and (b) two concurrent spawns racing a naive read-modify-write let
// the second writer's write silently discard the first spawn's entry. The fix is
// ClaudeCodeLauncher.writeAtomically (sibling temp file + ATOMIC_MOVE, package-visible for the
// test below that proves it directly) plus TRUST_JSON_LOCK (a process-wide lock serialising
// every seedTrustDialog call this launcher itself makes).
/**
* Requested test 1: an existing, large-ish (not just a two-key fixture) .claude.json must never
* collapse. Asserts on the restored KEY SET — not merely that the result still parses as JSON,
* which the incident's 178-byte file also did.
*/
@Test
void seedTrustDialogPreservesALargeExistingFileWithoutCollapsing(
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
Path claudeJson = configDir.resolve(".claude.json");
ObjectMapper mapper = new ObjectMapper();
ObjectNode root = mapper.createObjectNode();
root.put("numStartups", 4200);
root.put("firstStartTime", "2025-01-01T00:00:00.000Z");
root.putObject("oauthAccount").put("emailAddress", "operator@example.com");
Set<String> otherProjectPaths = new HashSet<>();
ObjectNode projects = root.putObject("projects");
for (int i = 0; i < 30; i++) {
String path = "/Users/operator/code/project-" + i;
otherProjectPaths.add(path);
ObjectNode project = projects.putObject(path);
project.put("hasTrustDialogAccepted", true);
project.putObject("mcpServers").putObject("server-" + i).put("command", "server-" + i + "-mcp");
}
String before = mapper.writerWithDefaultPrettyPrinter().writeValueAsString(root);
Files.writeString(claudeJson, before);
long sizeBefore = Files.size(claudeJson);
assertTrue(sizeBefore > 4096,
"fixture must actually be large-ish to be a meaningful proof: " + sizeBefore + " bytes");
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
long sizeAfter = Files.size(claudeJson);
assertTrue(sizeAfter >= sizeBefore,
"the file must never collapse below its pre-seed size — before=" + sizeBefore
+ " after=" + sizeAfter + " bytes (the incident shrank ~72 KB to 178 bytes)");
JsonNode after = mapper.readTree(claudeJson.toFile());
assertEquals(4200, after.path("numStartups").asInt(), "unrelated top-level key survives");
assertEquals("operator@example.com", after.path("oauthAccount").path("emailAddress").asText());
JsonNode afterProjects = after.path("projects");
for (String path : otherProjectPaths) {
assertTrue(afterProjects.path(path).path("hasTrustDialogAccepted").asBoolean(false),
"pre-existing project entry " + path + " must survive");
}
assertEquals(otherProjectPaths.size() + 1, afterProjects.size(),
"exactly one NEW project entry (this worktree's) must be added, none dropped");
assertTrue(afterProjects.path(worktree.toString()).path("hasTrustDialogAccepted").asBoolean(false));
}
/**
* Requested test 2: two concurrent spawns for DIFFERENT cwd values sharing one configDir must
* both end up present in the final file — a naive concurrent read-modify-write would let the
* second writer's read (taken before the first writer's write lands) silently discard the
* first. A {@link CountDownLatch} — not a sleep — lines both threads up at the starting line so
* this does not depend on scheduling luck to be meaningful.
*
* <p>This is a test of {@code TRUST_JSON_LOCK}, not of {@code writeAtomically}: the lock fully
* serialises every {@code seedTrustDialog} call this launcher itself makes, so this test would
* pass even without atomicity. See {@link #writeAtomicallyNeverExposesATornFileToAConcurrentReader}
* for the test that exercises atomicity specifically.
*/
@Test
void concurrentSeedsForDifferentCwdsBothSurvive(
@TempDir Path configDir, @TempDir Path worktreeA, @TempDir Path worktreeB) throws Exception {
markAsProvisionedWorktree(worktreeA);
markAsProvisionedWorktree(worktreeB);
CountDownLatch ready = new CountDownLatch(2);
CountDownLatch go = new CountDownLatch(1);
AtomicReference<Exception> failureA = new AtomicReference<>();
AtomicReference<Exception> failureB = new AtomicReference<>();
Thread ta = new Thread(spawnTask(configDir, worktreeA, ready, go, failureA), "spawn-a");
Thread tb = new Thread(spawnTask(configDir, worktreeB, ready, go, failureB), "spawn-b");
ta.start();
tb.start();
assertTrue(ready.await(5, TimeUnit.SECONDS), "both threads must reach the starting line");
go.countDown();
ta.join(5000);
tb.join(5000);
assertFalse(ta.isAlive(), "spawn A must finish within the timeout");
assertFalse(tb.isAlive(), "spawn B must finish within the timeout");
assertNull(failureA.get(), "spawn A must not throw: " + failureA.get());
assertNull(failureB.get(), "spawn B must not throw: " + failureB.get());
JsonNode root = new ObjectMapper().readTree(configDir.resolve(".claude.json").toFile());
assertTrue(root.path("projects").path(worktreeA.toString())
.path("hasTrustDialogAccepted").asBoolean(false),
"worktree A's entry must survive the race");
assertTrue(root.path("projects").path(worktreeB.toString())
.path("hasTrustDialogAccepted").asBoolean(false),
"worktree B's entry must survive the race");
}
private Runnable spawnTask(Path configDir, Path worktree, CountDownLatch ready, CountDownLatch go,
AtomicReference<Exception> failure) {
return () -> {
try {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(herdr),
new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
ready.countDown();
go.await();
launcher.spawn();
} catch (Exception e) {
failure.set(e);
}
};
}
/**
* Requested test 3: the atomic write must preserve {@code .claude.json}'s existing {@code 0600}
* permissions, not silently widen them via a fresh temp file's own defaults landing on top of a
* file that had different (e.g. group-readable) permissions. Skips rather than fails on a
* filesystem with no POSIX permissions (e.g. Windows) — the same {@code assumeTrue} pattern this
* file already uses for {@code memberHerdrSocketWithWorktreeRootAndGroupPutsCharterUnderWorktreeRootAndSharesIt}.
*/
@Test
void seedTrustDialogPreservesExisting0600Permissions(
@TempDir Path configDir, @TempDir Path worktree) throws Exception {
markAsProvisionedWorktree(worktree);
Path claudeJson = configDir.resolve(".claude.json");
Files.writeString(claudeJson, "{}");
assumeTrue(Files.getFileAttributeView(claudeJson, PosixFileAttributeView.class) != null,
"no POSIX permissions on this filesystem — skipping rather than failing");
Files.setPosixFilePermissions(claudeJson, PosixFilePermissions.fromString("rw-------"));
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = trustProfile(configDir.toString(), worktree.toString());
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null)
.spawn();
assertEquals("rw-------", PosixFilePermissions.toString(Files.getPosixFilePermissions(claudeJson)),
"the atomic write must preserve .claude.json's existing 0600 permissions, not widen "
+ "them via a fresh temp file's own defaults");
}
/**
* Direct proof that {@link ClaudeCodeLauncher#writeAtomically} — not {@code TRUST_JSON_LOCK} —
* is what keeps a concurrent reader of {@code target} from ever observing a truncated or
* partially-written file. Calls {@code writeAtomically} directly (package-visible for exactly
* this test) rather than going through {@code seedTrustDialog}/{@code spawn()}, because {@link
* #concurrentSeedsForDifferentCwdsBothSurvive} above is protected by the lock and would pass
* even without atomicity — it proves the lock, not the atomic move. This test proves the atomic
* move specifically: a reader racing a writer OUTSIDE that lock (a second daemon process, or the
* operator's own live Claude Code — exactly what the lock cannot reach) must still never see a
* torn file.
*
* <p>The new content is made large (tens of MB) so a naive truncate-then-write has a real,
* non-instantaneous window for the busy-poll reader thread to land in — this is inherently a
* race, not a guaranteed-deterministic assertion, but it uses no {@code Thread.sleep} and
* reliably reproduced the torn read when run against the pre-fix
* {@code Files.writeString(target, content)} implementation (see the PR's mutation table).
*/
@Test
void writeAtomicallyNeverExposesATornFileToAConcurrentReader(@TempDir Path dir) throws Exception {
Path target = dir.resolve(".claude.json");
String oldContent = "{\"marker\":\"OLD\"}";
Files.writeString(target, oldContent);
String newContent = "{\"marker\":\"NEW\",\"pad\":\"" + "x".repeat(20_000_000) + "\"}";
AtomicReference<String> tornRead = new AtomicReference<>();
AtomicBoolean stop = new AtomicBoolean(false);
Thread reader = new Thread(() -> {
while (!stop.get()) {
try {
String seen = Files.readString(target);
if (!seen.equals(oldContent) && !seen.equals(newContent)) {
tornRead.compareAndSet(null, "torn read of length " + seen.length() + ": "
+ seen.substring(0, Math.min(seen.length(), 80)));
stop.set(true);
}
} catch (IOException ignored) {
// ATOMIC_MOVE guarantees the path always resolves to old-or-new content once
// readable at all — a transient "briefly missing during the rename" is fine to
// ignore and keep sampling.
}
}
});
reader.start();
try {
ClaudeCodeLauncher.writeAtomically(target, newContent);
} finally {
stop.set(true);
reader.join(5000);
}
assertNull(tornRead.get(), "a concurrent reader must never observe a partially-written file: "
+ tornRead.get());
assertEquals(newContent, Files.readString(target), "the final content must be the new content");
}
}
@@ -56,11 +56,7 @@ class MessageServiceTest {
/** Run {@code send} on a background thread; the current thread drives the worker's turn. */
private CompletableFuture<MessageService.Reply> sendAsync() {
return sendAsync("do the task");
}
private CompletableFuture<MessageService.Reply> sendAsync(String content) {
return CompletableFuture.supplyAsync(() -> messages.send(T, content, 5000));
return CompletableFuture.supplyAsync(() -> messages.send(T, "do the task", 5000));
}
private void awaitWaiting() throws InterruptedException {
@@ -90,94 +86,6 @@ class MessageServiceTest {
assertTrue(reply.completed(), "a scraped completion still counts as completed");
}
@Test
void completionFallbackReplacesAnEchoedInjectedBriefWithNoReportOutcome() throws Exception {
String brief = "Implement the requested change. ".repeat(20);
CompletableFuture<MessageService.Reply> send = sendAsync(brief);
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
herdr.readText("⏺ " + brief + "\n❯ ");
injector.onStatus(T, AgentStatus.IDLE);
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome());
assertEquals(CompletionResolver.NO_REPORT_PREFIX + "]", reply.text(),
"the real injector -> completion fallback path must not return the lead's brief");
}
@Test
void completionFallbackKeepsARealReportThatRestatesTheWholeBrief() throws Exception {
String brief = "Load the implementer skill. You own fleetd #999. ".repeat(12);
String report = brief + "\n\n## Report\n"
+ ("I implemented the fix in GitWorktrees.java, added five tests, ran mvn clean install "
+ "and got 1168 tests with 0 failures. Commit 321d8dc pushed. ").repeat(8);
CompletableFuture<MessageService.Reply> send = sendAsync(brief);
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
herdr.readText("⏺ " + report + "\n❯ ");
injector.onStatus(T, AgentStatus.IDLE);
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.COMPLETED_UNREPLIED, reply.outcome());
assertEquals(report.strip(), reply.text(), "a report that restates the whole brief must survive unchanged");
}
@Test
void completionFallbackNamesTheKnownWorktreeAndBranchForAnEchoedBrief() throws Exception {
FakeHerdr localHerdr = new FakeHerdr();
Rendezvous localRendezvous = new Rendezvous();
java.util.concurrent.atomic.AtomicLong localClock = new java.util.concurrent.atomic.AtomicLong();
CompletionResolver localCompletion = new CompletionResolver(new AgentControl(localHerdr), localRendezvous,
ExhaustedPatternLookup.none(), ExhaustionSink.none(),
dev.ltms.fleet.inject.BackendErrorPatternLookup.legacy(),
dev.ltms.fleet.inject.BackendErrorSink.none(),
() -> localClock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1),
target -> new CompletionResolver.WorktreeBranch("/tmp/member-worktree", "worker/cb241"));
Injector localInjector = new Injector(new AgentControl(localHerdr), localCompletion);
MessageService localMessages = new MessageService(new AgentControl(localHerdr), localInjector,
localRendezvous, new InMemoryReplyInbox());
String brief = "Implement the requested change. ".repeat(20);
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> localMessages.send(T, brief, 5000));
long deadline = System.currentTimeMillis() + 2000;
while (!localRendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(localRendezvous.isWaiting(T));
localHerdr.readText("$ prompt");
localInjector.onStatus(T, AgentStatus.IDLE);
localInjector.onStatus(T, AgentStatus.WORKING);
localHerdr.readText("⏺ " + brief + "\n❯ ");
localInjector.onStatus(T, AgentStatus.IDLE);
assertEquals(CompletionResolver.NO_REPORT_PREFIX + " worktree=/tmp/member-worktree "
+ "branch=worker/cb241]", send.get(5, TimeUnit.SECONDS).text());
}
@Test
void completionFallbackKeepsTheClippedMarkerWhenAnEchoedBriefIsTooLong() throws Exception {
String brief = "a".repeat(4_001); // CompletionResolver's 4,000-character scrape cap
CompletableFuture<MessageService.Reply> send = sendAsync(brief);
awaitWaiting();
herdr.readText("$ prompt");
injector.onStatus(T, AgentStatus.IDLE);
injector.onStatus(T, AgentStatus.WORKING);
herdr.readText("⏺ " + brief + "\n❯ ");
injector.onStatus(T, AgentStatus.IDLE);
assertEquals(CompletionResolver.NO_REPORT_PREFIX + "]\n"
+ "[Pane tail clipped: member did not call fleet_reply.]",
send.get(5, TimeUnit.SECONDS).text());
}
@Test
void backendErrorScrapeThroughMessageServiceFailsInsteadOfBecomingReplyText() throws Exception {
// fleetd#164 (part 2 addendum): a scrape that reads cleanly but is only the backend's own