Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha ee5f8b932b #137: complete an async ticket's own reply after answer() times out
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Successful in 1m44s
fleet_send{turnId} (MessageService.answer) blocks the primary only for its
own bounded MCP-call window (25s default, 120s max) — far shorter than a
worker's resumed turn can genuinely take. When that window expires, answer()
closes its rendezvous waiter, so the worker's eventual fleet_reply has no
live waiter to resolve and falls back to the session inbox. The async
ticket's future was never completed by that path, so fleet_poll{ticket}
stayed PENDING until fleet_stop's abandon() forced it FAILED with a
misleading "the worker session was released before it replied" reason,
even though the reply had genuinely arrived.

- MessageService.reply(): before falling to the inbox, look for the async
  task this exact turn belongs to (already answered — question cleared,
  turnId still stamped — but not yet resolved) and complete it directly
  with the real reply, so fleet_poll{ticket} returns it.
- MessageService.abandon(): defense in depth, independent of the above —
  never write a false WORKER_FAILED once a reply reached the inbox for
  this target; recover and use its real content instead.
- Two new tests drive the full delegation path (async send -> ask ->
  answer with a short timeout -> reply -> poll/abandon), not a reply sink
  directly; both fail with the fix disabled and pass with it restored.
2026-08-31 14:12:23 +07:00
16 changed files with 317 additions and 635 deletions
+166
View File
@@ -0,0 +1,166 @@
# CB-137 / fleetd issue #137 — report
## Real root cause (not the hypothesis in the ticket)
I read `MessageService.java` and `Rendezvous.java` before changing anything. The mechanism is real,
but the exact place it happens is `MessageService.answer()`, not "the reply goes to the inbox on
purpose" in general.
1. A lead delegates with `fleet_send{wait:false}` → `sendAsync()` creates a `Task` and runs `send()`
on a background virtual thread with a 30-minute internal budget (`ASYNC_TIMEOUT_MS`).
2. The worker calls `fleet_ask`. That resolves the open rendezvous waiter with `Kind.QUESTION`, so
`send()` returns immediately and the `Task` is left open (its `future` stays unresolved — see the
comment in `sendAsync`'s lambda: "Keep the accepted owner until answer() finishes it").
3. The lead answers with `fleet_send{turnId, content}`. This calls `FleetMcp.answer()` →
`MessageService.answer(turnId, content, timeout)`. The `timeout` here is **not** the generous
30-minute async budget — it is the MCP tool's own bounded wait: `DEFAULT_TIMEOUT_MS = 25_000`,
clamped to at most `MAX_TIMEOUT_MS = 120_000` (`FleetMcp.java:71-72,495,512`). This is the same
~60–120s window every blocking `fleet_send` call is capped at (documented elsewhere as "the
caller's own MCP client call timeout").
4. `answer()` opens a **fresh** rendezvous waiter for the worker session and blocks on it for at most
that window. If the worker's resumed turn takes longer than that to actually finish (very
plausible — the resumed turn can mean more edits, a build, a commit, a push, opening a PR), the
wait times out. On timeout, `answer()`'s `finally` block unconditionally calls
`rendezvous.close(workerSession, reply)`, **removing the waiter from the map**, and returns
`Outcome.TIMED_OUT_WORKING` to the lead.
5. The worker keeps working, unaware anything happened, and eventually calls `fleet_reply`. That
reaches `MessageService.reply(session, content)`, which tries `rendezvous.resolve(session,
content)` — but the waiter was already closed in step 4, so `resolve` returns `false`. `reply()`
then falls back to `inbox.publish(...)` and marks `strandedReplies.put(session, true)`
(CB-640 bookkeeping) — the reply is safely held, but **the async `Task`'s `future` is never
completed**.
6. `fleet_poll{ticket}` keeps returning `PENDING` forever (the `Task` never resolves) — until the
lead eventually calls `fleet_stop`. That fires `sessions.onRelease` → `messages.abandon(target,
reason)` (`Fleetd.java:481-497`), where `reason` is built with the exact text from the bug report
("the worker session was released before it replied; worktree=... branch=... snapshot=...",
`Fleetd.java:484-487`). `abandon()`'s loop finds the still-open `Task` (`question == null`, future
not done) and completes it as `WORKER_FAILED` with that misleading reason — even though the
worker's real reply is sitting, intact, in the inbox the whole time.
So: the reported behaviour is correct, and the specific trigger is `answer()`'s own bounded wait
being shorter than the worker's real resumed-turn time — not anything to do with the ~55s
`fleet_ask` window itself (that part, issue #61, is untouched).
## Fix
Two changes in `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`, both scoped to the
ticket/reply routing and the terminal-state text — `fleet_ask`'s own window and mechanics are
untouched.
**1. `reply()` — priority 1 (the ticket resolves with the real reply).**
Before falling back to the inbox, `reply()` now looks for an async `Task` that is specifically in the
"already answered but not yet resolved" state (`question == null`, `turnId != null` — set once
`answer()` has cleared the question but before anything completed the future, `!future.isDone()`).
If one exists for this `target`, the worker's reply completes that `Task`'s future directly as
`Outcome.REPLIED` with the real content, and the reply never touches the inbox at all. A task that
was never asked has `turnId == null` and can never match, so ordinary (no-`fleet_ask`) delegations
are unaffected — they already resolve through the pre-existing rendezvous fast path.
I chose this over leaving `answer()`'s own timeout behaviour untouched and instead keeping its
rendezvous waiter open in the background: that alternative works but reopens the "at most one
waiter per session" invariant (`Rendezvous.open` throws on a double-open) to a new class of races
with a fresh send arriving mid-window. The `send()` path already guards against sending into an
answered-but-still-resolving worker via `hasAsyncQuestion(target)` (checks `asyncTasksByTurn`,
which still holds the task until it resolves), so routing through `reply()` gets the same protection
without touching `answer()`'s waiter lifecycle at all — the smaller, safer diff.
**2. `abandon()` — priority 3 (required independently, "even if you fix (1)").**
Before marking any of a released target's still-open tasks `WORKER_FAILED`, `abandon()` now checks
`hasStrandedReply(target)` (the existing CB-640 fact — true whenever the *last* `reply()` for this
target fell through to the inbox). If true, it drains the inbox (`recoverStrandedReply`) and — if it
actually finds a message — completes the task as `REPLIED` with that real content instead of writing
the failure. This is deliberately a **separate** check from fix 1: fix 1 already prevents the
inbox-stranding from happening in the exact scenario this ticket describes, so by the time
`abandon()` runs the task is normally already resolved and `abandon()`'s `complete()` call is a
harmless no-op. This second check exists so that if some *other* future path ever strands a reply
in the inbox without resolving its ticket, `abandon()` still refuses to report a false failure —
"if a reply reached any sink for that turn, the terminal state is done," per the ticket. I verified
both are required by disabling each independently and confirming the two new tests fail (see below).
**Priority 4 (the snapshot/worktree hint).** Handled as a consequence of both fixes rather than a
separate branch: once a task resolves as `REPLIED` (via either fix), `abandon()` never calls
`new Reply(Outcome.WORKER_FAILED, reason)` for that task at all, so the "the worker session was
released before it replied; worktree=... branch=... snapshot=..." text is never constructed or
attached to that ticket's outcome. It still appears, correctly, for a task that never got a reply
(the existing `abandonFailsEveryPendingAsyncTicketForTheReleasedTarget` /
`anAbandonedAsyncTaskPollsAsFailedNotPending` tests still pass unchanged).
**Priority 2** was not needed — fix 1 makes `fleet_poll{ticket}` return the actual reply (the
higher-priority option), so I did not fall back to "the ticket merely resolves as done with no
content."
## Tests — driven through the real delegation path, not the reply sink directly
Both new tests in `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` go through
`sendAsync` → `injectDelivery` → `ask` → `answer` (with a short timeout, so it genuinely times out,
mirroring the ~25–120s real MCP-call bound vs. a longer resumed turn) → `reply` → `poll`/`abandon`.
No test constructs a `Reply` and hands it to a sink directly.
- `aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket` — asserts `fleet_poll{ticket}` (via
`messages.poll`) reaches `Phase.DONE` with the worker's actual reply text and
`replySource() == "reply"`, and that `hasStrandedReply(T)` stays `false` (proves the reply never
touched the inbox at all — fix 1 caught it).
- `fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket` — same setup, then calls `abandon(T, "the
worker session was released before it replied")` (what `fleet_stop` triggers) and asserts it
returns `false` (no failure recorded) and the ticket still polls `DONE` with the real reply.
**Proof both fail without the change.** I temporarily short-circuited both new private methods
(`askAnsweredAsyncTask` → always `null`, `recoverStrandedReply` → always `null`) — i.e. disabled
both fixes — and ran just these two tests:
```
[ERROR] Tests run: 2, Failures: 2, Errors: 0, Skipped: 0
dev.ltms.fleet.msg.MessageServiceTest.aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket
org.opentest4j.AssertionFailedError: expected: <DONE> but was: <PENDING>
dev.ltms.fleet.msg.MessageServiceTest.fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket
org.opentest4j.AssertionFailedError: a reply already arrived, so nothing here is a genuine failure
==> expected: <false> but was: <true>
```
This is the exact bug: the ticket stays `PENDING` forever, and `abandon()` reports `true` (a
failure) even though a reply had already arrived. I then restored both fixes (verified with
`grep -n "TEMP #137-proof"` finding nothing) and re-ran — both pass.
## Build
Ran from `fleetd/`, unpiped, full output read (not `| tail`):
```
mvn clean install
...
[INFO] Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0
[INFO] BUILD SUCCESS
[INFO] Total time: 36.315 s
```
Main was at 1037 tests; this branch adds the 2 new tests above → 1039, all green, `exit=0`.
## What I could NOT check
- No IDE tooling is mounted for me (worker), so no `ide_diagnostics`/IntelliJ inspection pass — only
`mvn clean install` (compiler + full test suite), as the worker procedure allows.
- I cannot restart the daemon or dogfood this live — I have no forge/daemon control. This is
unverified against a real herdr pane, a real MCP client's ~60s call cap, or a real worker session;
everything above is verified only through the JUnit fixture's simulated timing
(`FakeHerdr`/`injector.onStatus`/direct `messages.answer(...,150)` calls), not a live fleet.
A primary should still consider a short live dogfood (an async delegation that asks, gets answered,
and takes longer than ~2 minutes to reply) before calling this closed.
- I did not touch, and did not re-verify, the `fleet_ask` ~55s window itself (issue #61) — out of
scope per the brief.
## Scope note (not investigated further)
`answer()`'s nested/double-`fleet_ask` case (the worker asks a second question before ever
replying to the first answer) has some pre-existing behaviour around which `turnId` a `QUESTION`
resolution gets attributed to that I did not fully untangle — it predates this change, my fix does
not touch it, and it is unrelated to the reported defect. Flagging only; not investigated further.
## Handoff
- Branch: `worker/cb-137-ask-ticket-e7760c-2`
- Worktree root: `/Users/dai.ha/LTMS/.bridged-worktrees/734324-2`
- Files changed:
- `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`
- `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java`
- `REPORT-cb137.md` (this file)
- Build: `Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS` (verbatim above)
-18
View File
@@ -28,7 +28,6 @@
<testcontainers.version>1.20.4</testcontainers.version>
<commons-compress.version>1.27.1</commons-compress.version>
<commons-lang3.version>3.18.0</commons-lang3.version>
<sqlite-jdbc.version>3.53.4.0</sqlite-jdbc.version>
</properties>
<!--
@@ -45,12 +44,6 @@
3.0-rc5; bumping Jackson 3 to the patched 3.2.x breaks the SDK (annotation mismatch).
Only the loopback /mcp endpoint parses this JSON, from trusted local Claude clients.
The 11.0.23 -> 11.0.25 bump did clear jetty CVE-2024-8184 (5.9) and CVE-2024-6763.
fleetd #206: org.xerial:sqlite-jdbc 3.53.4.0 (added for OpenCodeSessionDiscovery) — the
only known advisory against this artifact is CVE-2023-32697 (RCE via an attacker-controlled
JDBC URL), fixed in 3.41.2.2; 3.53.4.0 is well past that fix and OSV.dev reports no open
advisory against it. Checked via the OSV.dev API (no Mend.io/JetBrains IDE MCP mount
available from this worktree) on 2026-08-31.
-->
<!-- Force the latest patched Jetty 11.x across all Javalin-pulled Jetty modules (no version
@@ -130,17 +123,6 @@
<version>${amqp.version}</version>
</dependency>
<!-- fleetd #206: opencode moved its session store from a JSON tree to SQLite
(opencode.db). This is the JDBC driver OpenCodeSessionDiscovery uses to read it
read-only. Ships bundled native libraries (linux/mac/windows, several archs), so it
is a heavier jar than most deps here — see the pom's dependency-security note below
for the size/CVE tradeoff actually measured. -->
<dependency>
<groupId>org.xerial</groupId>
<artifactId>sqlite-jdbc</artifactId>
<version>${sqlite-jdbc.version}</version>
</dependency>
<!-- Logging -->
<dependency>
<groupId>org.slf4j</groupId>
@@ -177,14 +177,14 @@ public final class Fleetd {
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet(),
() -> config.get().memberCredentials(), null, config::get));
() -> config.get().memberCredentials()));
}
if (!opencodeProfiles.isEmpty()) {
adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(),
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
() -> config.get().fleet(),
() -> config.get().memberCredentials(), config::get));
() -> config.get().memberCredentials()));
}
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
@@ -2,11 +2,8 @@ package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.OptionalLong;
import java.util.Set;
/**
* Resolves which herdr pane a process belongs to — the herdr half of connection-based MCP
@@ -15,46 +12,23 @@ import java.util.Set;
* is calling without the worker sending anything spoofable.
*
* <p>herdr owns the PID→pane truth: {@code pane.process_info} reports each pane's {@code shell_pid}
* and foreground process PIDs. A pid that is neither of those directly — e.g. a grandchild a
* worker spawned, such as a {@code python3} or {@code curl} helper that opens its own MCP
* connection — is resolved by walking its ancestry (via {@link ParentResolver}) up to the root and
* matching any ancestor against a pane's {@code shell_pid} or foreground pids (CB-161). Without
* this walk such a pid matches no pane, and the caller falls through to loopback-trust and is
* resolved as the primary — a worker→primary privilege escalation.
*
* <p>This scans agent panes; a spawn-time {@code pid→terminal} cache is the obvious optimization
* once wired into {@code ClaudeCodeLauncher}.
* and foreground process PIDs. This scans agent panes; a spawn-time {@code pid→terminal} cache is
* the obvious optimization once wired into {@code ClaudeCodeLauncher}.
*
* <p>CB-185 split the fleet across two herdr daemons — lead operations on one, members on the
* other ({@code memberHerdrSocket}). A caller's pane can live on <em>either</em> daemon (a lead's
* MCP connection resolves against the lead daemon; a member's against the member daemon), so this
* must be able to search more than one client. {@link #PaneLocator(HerdrClient, HerdrClient)}
* searches the lead client first, then the member client, and collapses to a single scan when the
* two are the same object (the historical single-daemon deployment). The caller's ancestor set is
* computed once per {@link #terminalForPid} call and reused across every client searched — it
* does not depend on which daemon a pane happens to live on.
* two are the same object (the historical single-daemon deployment).
*/
public final class PaneLocator {
/**
* Bound on how many ancestor generations {@link #ancestorsOf} walks. This runs on every MCP
* call, so a cycle or a pathologically deep process tree must not hang identity resolution;
* 32 generations is far more than any real worker→helper process tree needs.
*/
private static final int MAX_ANCESTRY_DEPTH = 32;
private final List<HerdrClient> herdrs;
private final ParentResolver parentResolver;
/** Search only this client — the single-daemon deployment. */
public PaneLocator(HerdrClient herdr) {
this(herdr, ParentResolver.PROCESS_HANDLE);
}
/** Search only this client, resolving ancestry through {@code parentResolver} — for tests. */
public PaneLocator(HerdrClient herdr, ParentResolver parentResolver) {
this.herdrs = List.of(herdr);
this.parentResolver = parentResolver;
}
/**
@@ -63,13 +37,7 @@ public final class PaneLocator {
* collapses to one client and one scan, exactly {@link #PaneLocator(HerdrClient)}'s behaviour.
*/
public PaneLocator(HerdrClient lead, HerdrClient member) {
this(lead, member, ParentResolver.PROCESS_HANDLE);
}
/** Two-daemon deployment, resolving ancestry through {@code parentResolver} — for tests. */
public PaneLocator(HerdrClient lead, HerdrClient member, ParentResolver parentResolver) {
this.herdrs = lead == member ? List.of(lead) : List.of(lead, member);
this.parentResolver = parentResolver;
}
/**
@@ -81,9 +49,8 @@ public final class PaneLocator {
if (pid <= 0) {
return null;
}
Set<Long> ancestry = ancestorsOf(pid);
for (HerdrClient herdr : herdrs) {
String terminal = terminalForPid(herdr, ancestry);
String terminal = terminalForPid(herdr, pid);
if (terminal != null) {
return terminal;
}
@@ -91,54 +58,28 @@ public final class PaneLocator {
return null;
}
/**
* {@code pid} itself plus its ancestor chain, walked through {@link #parentResolver} up to
* {@link #MAX_ANCESTRY_DEPTH} generations or pid 1, whichever comes first. A vanished ancestor
* ({@link ParentResolver#parentOf} returning empty) ends the walk without error — it just means
* the chain is shorter than the bound. A cycle in a fake resolver is caught by the "already
* seen" check and also ends the walk, so this can never loop.
*/
private Set<Long> ancestorsOf(long pid) {
Set<Long> ancestry = new LinkedHashSet<>();
long current = pid;
for (int depth = 0; depth < MAX_ANCESTRY_DEPTH; depth++) {
if (current <= 0 || !ancestry.add(current)) {
break; // vanished/invalid pid, or a cycle back to a pid already recorded
}
if (current == 1) {
break; // reached the root of the process tree
}
OptionalLong parent = parentResolver.parentOf(current);
if (parent.isEmpty()) {
break; // vanished ancestor — not an error, just the end of the chain
}
current = parent.getAsLong();
}
return ancestry;
}
private static String terminalForPid(HerdrClient herdr, Set<Long> ancestry) {
private static String terminalForPid(HerdrClient herdr, long pid) {
for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) {
String paneId = pane.path("pane_id").asText(null);
if (paneId != null && paneOwnsAnyOf(herdr, paneId, ancestry)) {
if (paneId != null && paneOwnsPid(herdr, paneId, pid)) {
return pane.path("terminal_id").asText(null);
}
}
return null;
}
private static boolean paneOwnsAnyOf(HerdrClient herdr, String paneId, Set<Long> ancestry) {
private static boolean paneOwnsPid(HerdrClient herdr, String paneId, long pid) {
JsonNode info;
try {
info = herdr.call("pane.process_info", Map.of("pane_id", paneId)).path("process_info");
} catch (HerdrException e) {
return false; // pane vanished mid-scan — just skip it
}
if (ancestry.contains(info.path("shell_pid").asLong(-1))) {
if (info.path("shell_pid").asLong(-1) == pid) {
return true;
}
for (JsonNode p : info.path("foreground_processes")) {
if (ancestry.contains(p.path("pid").asLong(-1))) {
if (p.path("pid").asLong(-1) == pid) {
return true;
}
}
@@ -1,21 +0,0 @@
package dev.ltms.fleet.herdr;
import java.util.OptionalLong;
/**
* Resolves a pid's parent pid — the seam {@link PaneLocator} walks a process's ancestry through,
* so its tests can drive the walk from a fake pid→parent map instead of spawning real processes.
*
* <p>{@link #PROCESS_HANDLE} is the production implementation, backed by {@link ProcessHandle}.
*/
public interface ParentResolver {
/** The parent pid of {@code pid}, or empty if {@code pid} is gone or has no known parent. */
OptionalLong parentOf(long pid);
/** Production resolver: asks the JVM's {@link ProcessHandle} view of the OS process tree. */
ParentResolver PROCESS_HANDLE = pid -> ProcessHandle.of(pid)
.flatMap(ProcessHandle::parent)
.map(parent -> OptionalLong.of(parent.pid()))
.orElse(OptionalLong.empty());
}
@@ -94,26 +94,13 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, guard, profiles, defaultProfile, env, spawnReadyTimeoutMs, spawnReadyPollMs,
fleet, memberCredentials, null, null);
}
/** Production constructor, plus the live config for URI environment exclusions. */
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames,
Supplier<FleetConfig> config) {
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, guard, profiles, defaultProfile, env,
spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
fleet, memberCredentials, hostEnvNames, config);
fleet, memberCredentials);
}
/**
@@ -166,22 +153,10 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, guard, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis, sleeper,
fleet, memberCredentials, null, null);
}
/** Full testability constructor, plus the live config for URI environment exclusions. */
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env, long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper, Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames,
Supplier<FleetConfig> config) {
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, config);
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
this.guard = guard;
}
@@ -199,7 +174,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, null);
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames);
this.guard = guard;
}
@@ -170,8 +170,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
/** Guards {@link #warnNonZsh} to one WARN per launcher instance, not one per spawn. */
private final AtomicBoolean nonZshShellWarned = new AtomicBoolean();
/** Live config provides URI environment names that must never enter member panes. */
private final Supplier<FleetConfig> config;
/**
* @param namePrefix label prefix for this peer kind (drives naming and reap)
@@ -240,19 +238,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames) {
this(namePrefix, agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis,
sleeper, fleet, memberCredentials, hostEnvNames, null);
}
/** As above, plus the live full config for secret-bearing URI environment names. */
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env, long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper, Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames, Supplier<FleetConfig> config) {
Supplier<Set<String>> hostEnvNames) {
this.fleet = fleet;
this.namePrefix = namePrefix;
this.agents = agents;
@@ -265,7 +252,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
this.sleeper = sleeper;
this.memberCredentials = memberCredentials;
this.hostEnvNames = hostEnvNames != null ? hostEnvNames : () -> System.getenv().keySet();
this.config = config;
}
// --- adapter seams -------------------------------------------------------------------------
@@ -1042,11 +1028,9 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
/** Put {@link #BLOCKED_CREDENTIAL_SENTINEL} over every blocked name in the pane-creation env map. */
private void overlayBlockedCredentials(Map<String, String> workerEnv,
FleetConfig.MemberCredentials creds) {
Set<String> blocked = new java.util.TreeSet<>(creds.blockedSet());
blocked.addAll(brokerUriEnvNames());
for (String name : blocked) {
private static void overlayBlockedCredentials(Map<String, String> workerEnv,
FleetConfig.MemberCredentials creds) {
for (String name : creds.blockedSet()) {
workerEnv.put(name, BLOCKED_CREDENTIAL_SENTINEL);
}
}
@@ -1113,21 +1097,15 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* ssh-agent handle when explicitly allowed, and the exact keys of THIS launch's own env map.
*/
private Set<String> derivedAllowedNames(FleetConfig.MemberCredentials creds, Launch launch) {
Set<String> brokerUriEnvNames = brokerUriEnvNames();
Set<String> allowed = new java.util.TreeSet<>(
MemberEnvAllowList.derive(profiles.values(), creds.allowSet(), brokerUriEnvNames));
MemberEnvAllowList.derive(profiles.values(), creds.allowSet()));
if (creds.sshAuthSockAllowed()) {
allowed.add(SSH_AUTH_SOCK);
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
allowed.addAll(launch.env().keySet());
allowed.removeAll(brokerUriEnvNames);
return allowed;
}
private Set<String> brokerUriEnvNames() {
return MemberEnvAllowList.brokerUriEnvNames(config == null ? null : config.get());
}
/**
* CB-633 follow-up: one INFO line per allow-list spawn WHOSE SCRUB ACTUALLY RUNS, so an operator
* can read a single log line and know the scrub ran and how much of the visible environment it
@@ -40,9 +40,8 @@ import java.util.TreeSet;
* operator's own explicit list. Before this, {@code policy: allow-list} silently ignored every name
* an operator wrote under {@code allow:} unless a profile happened to carry it too, which meant
* turning the policy on could blank credentials working members already depended on. {@code
* SSH_AUTH_SOCK} and configured broker URI environment names are exceptions: even when the operator
* lists them under {@code allow:}, they are excluded here. {@code SSH_AUTH_SOCK} is added back ONLY
* by the caller when {@code sshAuthSock: allow} is explicitly set
* SSH_AUTH_SOCK} is the one exception: even when the operator lists it under {@code allow:}, it is
* excluded here and added back ONLY by the caller when {@code sshAuthSock: allow} is explicitly set
* (see {@link #SSH_AUTH_SOCK}'s javadoc) — it is a live handle to the operator's own ssh-agent, not
* a value, so treating it like any other allow-listed name would hand a member every key the
* operator's agent holds the moment they typed the name under {@code allow:} for an unrelated
@@ -107,16 +106,6 @@ public final class MemberEnvAllowList {
* run-to-run.
*/
public static Set<String> derive(Collection<FleetConfig.Profile> profiles, Set<String> configuredAllow) {
return derive(profiles, configuredAllow, Set.of());
}
/**
* As {@link #derive(Collection, Set)}, while excluding names that fleetd knows carry credentials.
* A configured broker URI contains its AMQP password inline, so it must never reach a member,
* even when an operator put its variable name in {@code memberCredentials.allow:}.
*/
public static Set<String> derive(Collection<FleetConfig.Profile> profiles, Set<String> configuredAllow,
Set<String> excludedNames) {
Set<String> derived = new TreeSet<>(INFRASTRUCTURE_PASSTHROUGH);
if (profiles != null) {
for (FleetConfig.Profile p : profiles) {
@@ -135,37 +124,9 @@ public final class MemberEnvAllowList {
}
}
}
if (excludedNames != null) {
derived.removeAll(excludedNames);
}
return Set.copyOf(derived);
}
/**
* The host environment names whose values are AMQP URIs with inline passwords. Both broker
* connections belong to fleetd, never to a member pane. Blank and absent configuration changes
* nothing.
*
* <p><b>How strong this exclusion is depends on the policy, and the difference matters.</b>
* Under {@code policy: allow-list} it is enforced by the generated ZDOTDIR scrub, which runs
* AFTER the pane's shell has sourced the operator's chain — so a login shell that re-exports the
* name is still blanked. Under the deny-list policy there is no scrub: the name is only removed
* from the pre-shell env map, and a login shell that sources the operator's secret store
* re-exports it. That is the long-standing weakness of deny-list (a sourced file can undo it),
* not something this exclusion introduces, but it means deny-list deployments do NOT get this
* guarantee. The same caveat applies to the non-zsh path, which has no scrub at all — see
* {@code HerdrPeerLauncher#applyEnvironmentAllowListPolicy}.
*/
public static Set<String> brokerUriEnvNames(FleetConfig config) {
if (config == null) {
return Set.of();
}
Set<String> names = new TreeSet<>();
addUriEnvIfPresent(names, config.broker());
addUriEnvIfPresent(names, config.coordinator());
return Set.copyOf(names);
}
/**
* Whether {@code name} survives the scrub when {@code allowedNames} is the derived set: an exact
* match, or an infrastructure-prefixed name ({@code LC_*}). Prefix rules live ONLY here and in
@@ -185,16 +146,4 @@ public final class MemberEnvAllowList {
into.add(name);
}
}
private static void addUriEnvIfPresent(Set<String> into, FleetConfig.Broker broker) {
if (broker != null && broker.hasUriEnv()) {
into.add(broker.uriEnv());
}
}
private static void addUriEnvIfPresent(Set<String> into, FleetConfig.Coordinator coordinator) {
if (coordinator != null && coordinator.hasUriEnv()) {
into.add(coordinator.uriEnv());
}
}
}
@@ -116,23 +116,12 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, spawnReadyPollMs,
fleet, memberCredentials, null);
}
/** Production constructor, plus the live config for URI environment exclusions. */
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env, long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<FleetConfig> config) {
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, config);
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials);
}
/**
@@ -193,20 +182,8 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
Path configRoot, Path discoveryRoot,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis, sleeper,
configRoot, discoveryRoot, fleet, memberCredentials, null);
}
/** Full testability constructor, plus the live config for URI environment exclusions. */
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env, long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper, Path configRoot, Path discoveryRoot,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<FleetConfig> config) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, null, config);
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
this.configRoot = configRoot;
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
}
@@ -1,87 +1,58 @@
package dev.ltms.fleet.member;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.sqlite.SQLiteConfig;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.Stream;
/**
* Resolves the opencode session id for a fleetd worker from opencode's on-disk storage — the
* only place this adapter touches opencode's private layout, and deliberately the <em>only</em>
* class that does.
*
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not
* a stable contract: opencode persists its session state in a SQLite database at
* {@code <storageRoot>/opencode.db} (a {@code session} table, one row per session, keyed by id and
* carrying a {@code directory} column). That schema can move between opencode releases exactly
* like the JSON-file layout it replaced did (opencode migrated off a one-JSON-file-per-session
* tree under {@code <storageRoot>/storage/session/<projectID>/ses_*.json} in January 2026 — that
* tree is now a frozen migration artefact nothing writes, which is why this class no longer reads
* it). opencode also ships a headless HTTP server that may supersede both of these entirely.
* Everything this adapter knows about that private storage — its shape and column names — lives
* here, so a layout change, or a switch to the HTTP server, changes exactly one class and nothing
* in {@link OpenCodeLauncher}.
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not a
* stable contract: opencode writes one JSON file per session under
* {@code <storageRoot>/session/<projectID>/<ses_*.json>}, and each record carries a
* {@code "version"} field (e.g. {@code "1.1.31"}), so the exact directory shape, file naming, and
* field names can move between opencode releases. opencode also ships a headless HTTP server that
* may supersede file scanning entirely. Everything this adapter knows about that private storage —
* its shape, naming, and field names — lives here, so a layout change, or a switch to the HTTP
* server, changes exactly one class and nothing in {@link OpenCodeLauncher}.
*
* <p>The determinism that makes this useful is structural, not a guess: every fleetd worker runs
* in its own unique git worktree, so the row's {@code directory} (its project root) equals the
* in its own unique git worktree, so the record's {@code directory} (its project root) equals the
* worker's cwd identifies <em>its</em> session unambiguously. We match on {@code directory} rather
* than diffing {@code opencode session list} before/after — that races under concurrent spawns, and
* the CLI listing does not even show the directory.
*
* <p>All reads are best-effort and never throw: a missing or unreadable database, a query that
* fails, or a directory with no row yet all yield {@code null}, and the caller (the session
* handle) treats that as "identity not resolved yet" and retries later. The database is opened
* read-only and never written to: opencode itself may be running and writing it concurrently (WAL
* mode), and this class must never disturb that.
* <p>All reads are best-effort and never throw: a missing or unreadable storage root, a record that
* fails to parse, or a directory with no record yet all yield {@code null}, and the caller (the
* session handle) treats that as "identity not resolved yet" and retries later.
*/
final class OpenCodeSessionDiscovery {
private static final Logger log = LoggerFactory.getLogger(OpenCodeSessionDiscovery.class);
private final Path storageRoot; // e.g. ~/.local/share/opencode (injectable for tests)
private final Path databasePath;
private final AtomicBoolean warnedMissingDatabase = new AtomicBoolean(false);
private final ObjectMapper json;
OpenCodeSessionDiscovery(Path storageRoot) {
this.storageRoot = storageRoot;
this.databasePath = storageRoot.resolve("opencode.db");
this.json = new ObjectMapper();
}
/**
* A connection to {@link #databasePath} opened with SQLite's {@code SQLITE_OPEN_READONLY}
* flag: it never creates the file, never writes, and never touches WAL or journal mode.
* opencode may be running and writing this database concurrently, and this class must never
* disturb it.
*
* <p>Package-private so a test can hold the connection and prove it refuses a write. That is
* the only way to pin this property: making the file unwritable does <em>not</em> work,
* because SQLite silently downgrades a read-write open of an unwritable file to read-only, so
* such a test passes whether or not the flag is set.
*/
Connection openReadOnly() throws SQLException {
SQLiteConfig config = new SQLiteConfig();
config.setReadOnly(true);
return config.createConnection("jdbc:sqlite:" + databasePath);
}
/**
* The opencode session id whose row references {@code directory} (the worker's cwd), or
* {@code null} when no row matches yet. When several rows share the directory — e.g. repeated
* spawns into the same worktree — the row with the highest {@code time_updated} wins: it is
* The opencode session id whose record references {@code directory} (the worker's cwd), or
* {@code null} when no record matches yet. When several records share the directory — e.g.
* repeated spawns into the same worktree — the <em>most recently modified</em> one wins: it is
* the session the pane most likely corresponds to.
*
* <p>Never throws: a missing {@code opencode.db}, a locked/unreadable database, a query
* failure, or a directory that has not been persisted yet all resolve to {@code null} rather
* than failing a spawn. A fleetd worker's session row is written lazily (when the session is
* first persisted), so {@code null} here is the normal answer right after the pane is ready,
* and the caller retries later.
* <p>Never throws: a missing {@code storageRoot}, an unreadable/malformed record, or a
* directory that has not been persisted yet all resolve to {@code null} rather than failing a
* spawn. A fleetd worker's session record is written lazily (when the session is first
* persisted), so {@code null} here is the normal answer right after the pane is ready, and the
* caller retries later.
*
* @param directory the worker's cwd, as resolved for this spawn
* @return the matching session id, or {@code null} if none is known yet
@@ -90,31 +61,62 @@ final class OpenCodeSessionDiscovery {
if (directory == null || directory.isBlank()) {
return null;
}
if (!Files.isRegularFile(databasePath)) {
if (warnedMissingDatabase.compareAndSet(false, true)) {
log.warn("opencode session database not found at {} — opencode's on-disk layout "
+ "may have moved again; session discovery will keep returning null",
databasePath);
}
Path sessionRoot = storageRoot.resolve("session");
if (!Files.isDirectory(sessionRoot)) {
return null;
}
String sql = "SELECT id FROM session WHERE directory = ? ORDER BY time_updated DESC LIMIT 1";
try (Connection connection = openReadOnly();
PreparedStatement statement = connection.prepareStatement(sql)) {
statement.setString(1, directory);
try (ResultSet rows = statement.executeQuery()) {
if (rows.next()) {
return rows.getString("id");
String best = null;
long bestMtime = Long.MIN_VALUE;
try (Stream<Path> projectDirs = Files.list(sessionRoot)) {
for (Path projectDir : projectDirs.filter(Files::isDirectory).toList()) {
try (Stream<Path> records = Files.list(projectDir)) {
for (Path record : records.toList()) {
String id = matchId(record, directory);
if (id == null) {
continue;
}
long mtime = lastModifiedEpochMillis(record);
if (mtime > bestMtime) {
bestMtime = mtime;
best = id;
}
}
} catch (IOException ignored) {
// one project dir unreadable — skip it; another may still match
}
}
} catch (SQLException e) {
// Locked, corrupt, or otherwise unreadable — never fatal to a spawn. Not the
// "database moved" signal (the file exists), so this stays below WARN.
log.debug("opencode session database unreadable at {}: {}", databasePath, e.toString());
} catch (IOException ignored) {
// storage root vanished or became unreadable — "no session known yet"
return null;
}
log.debug("no opencode session row for directory (root={}, directory={})",
storageRoot, directory);
return null;
return best;
}
/**
* The record's session id when it references {@code directory}, else {@code null}. A record
* that is not JSON, lacks {@code id}/{@code directory}, or points at a different directory is
* simply not our session; a malformed one is skipped, never fatal.
*/
private String matchId(Path record, String directory) {
try {
JsonNode node = json.readTree(record.toFile());
JsonNode id = node == null ? null : node.get("id");
JsonNode dir = node == null ? null : node.get("directory");
if (id == null || dir == null || !directory.equals(dir.asText())) {
return null;
}
return id.asText();
} catch (IOException e) {
return null;
}
}
/** The record's last-modified epoch ms, or {@code Long.MIN_VALUE} if unreadable (never wins). */
private static long lastModifiedEpochMillis(Path record) {
try {
return Files.getLastModifiedTime(record).toMillis();
} catch (IOException e) {
return Long.MIN_VALUE;
}
}
}
@@ -1,27 +0,0 @@
package dev.ltms.fleet.herdr;
import java.util.HashMap;
import java.util.Map;
import java.util.OptionalLong;
/**
* Fake {@link ParentResolver} backed by an explicit pid→parent map — lets {@link PaneLocatorTest}
* drive {@link PaneLocator}'s ancestry walk (grandchild pids, cycles) without spawning real OS
* processes.
*/
final class FakeParentResolver implements ParentResolver {
private final Map<Long, Long> parents = new HashMap<>();
/** {@code pid}'s parent is {@code parentPid}. A pid with no entry here has no known parent. */
FakeParentResolver parent(long pid, long parentPid) {
parents.put(pid, parentPid);
return this;
}
@Override
public OptionalLong parentOf(long pid) {
Long parent = parents.get(pid);
return parent == null ? OptionalLong.empty() : OptionalLong.of(parent);
}
}
@@ -1,11 +1,7 @@
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.junit.jupiter.api.Test;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.*;
/** Unit tests for PID → pane resolution (the herdr half of connection-based MCP identity). */
@@ -67,123 +63,4 @@ class PaneLocatorTest {
long paneListCalls = shared.calls.stream().filter(c -> c.method().equals("pane.list")).count();
assertEquals(1, paneListCalls, "same-object lead/member must scan exactly once, not twice");
}
// --- ancestry walk (CB-161: grandchild pids matched no pane, resolving as primary) --------
@Test
void stillResolvesAPidThatIsExactlyThePaneShellPid() {
// Regression: a pid with no parent chain at all — no ancestry walk is needed to match it.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
PaneLocator loc = new PaneLocator(pane, new FakeParentResolver());
assertEquals("term_x", loc.terminalForPid(5000));
}
@Test
void stillResolvesAPidThatIsExactlyAForegroundPid() {
// Regression: same as above, but matching via the foreground-processes list.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
PaneLocator loc = new PaneLocator(pane, new FakeParentResolver());
assertEquals("term_x", loc.terminalForPid(6000));
}
@Test
void resolvesAGrandchildPidTwoLevelsBelowTheShellPid() {
// The bug: a helper process a worker spawns (python3, curl, ...) is a grandchild of the
// pane's shell — not the shell_pid and not a foreground pid directly. Before the fix,
// paneOwnsPid only checked direct pid equality, so this pid matched no pane and the
// caller fell through to loopback-trust as the primary.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
FakeParentResolver parents = new FakeParentResolver()
.parent(7002, 7001) // grandchild -> child
.parent(7001, 5000); // child -> shell (the pane's shell_pid)
PaneLocator loc = new PaneLocator(pane, parents);
assertEquals("term_x", loc.terminalForPid(7002));
}
@Test
void nullForAPidWhoseAncestryMatchesNoPane() {
// Must not break the other direction: a pid that truly belongs to nothing here (e.g. the
// real primary) must still resolve to null. Resolving everything to a worker would demote
// the actual lead and refuse every orchestration call.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
FakeParentResolver parents = new FakeParentResolver()
.parent(9002, 9001)
.parent(9001, 9000); // chain never reaches 5000 or 6000
PaneLocator loc = new PaneLocator(pane, parents);
assertNull(loc.terminalForPid(9002));
}
@Test
void ancestryWalkTerminatesOnACycleInsteadOfHanging() {
// A fake (or corrupted) parent map that cycles must not hang identity resolution, which
// runs on every MCP call. The walk must still terminate and correctly resolve to null.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
FakeParentResolver parents = new FakeParentResolver()
.parent(100, 101)
.parent(101, 100); // cycle, never reaches the pane's pids
PaneLocator loc = new PaneLocator(pane, parents);
assertNull(loc.terminalForPid(100));
}
@Test
void ancestrySetIsComputedOnceAcrossBothClientsInTheTwoDaemonConstructor() {
// CB-185: the two-daemon constructor searches lead then member. The ancestor set is
// per-caller, not per-client — it must be walked once and reused, not recomputed for
// each client searched.
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
FakeParentResolver parents = new FakeParentResolver()
.parent(7002, 7001)
.parent(7001, 5000);
AtomicInteger calls = new AtomicInteger();
ParentResolver counting = pid -> {
calls.incrementAndGet();
return parents.parentOf(pid);
};
HerdrClient noPanes = new FakeHerdr().withNoPanes();
PaneLocator two = new PaneLocator(noPanes, pane, counting);
assertEquals("term_x", two.terminalForPid(7002));
assertEquals(3, calls.get(), "ancestry must be walked once (3 lookups: 7002, 7001, 5000), "
+ "not re-walked per herdr client");
}
/** Minimal single-pane {@link HerdrClient} fake, purpose-built for the ancestry tests above. */
private static final class OnePaneHerdr implements HerdrClient {
private final ObjectMapper mapper = new ObjectMapper();
private final String terminalId;
private final String paneId;
private final long shellPid;
private final long foregroundPid;
OnePaneHerdr(String terminalId, String paneId, long shellPid, long foregroundPid) {
this.terminalId = terminalId;
this.paneId = paneId;
this.shellPid = shellPid;
this.foregroundPid = foregroundPid;
}
@Override
public JsonNode call(String method, Object params) {
try {
return switch (method) {
case "pane.list" -> mapper.readTree(("""
{"type":"pane_list","panes":[
{"pane_id":"%s","terminal_id":"%s","workspace_id":"w1","tab_id":"w1:t1","agent":"claude"}]}""")
.formatted(paneId, terminalId));
case "pane.process_info" -> mapper.readTree(("""
{"type":"pane_process_info","process_info":{"pane_id":"%s","shell_pid":%d,
"foreground_processes":[{"pid":%d,"name":"node","argv0":"claude"}]}}""")
.formatted(paneId, shellPid, foregroundPid));
default -> throw new HerdrException("OnePaneHerdr has no canned response for " + method);
};
} catch (HerdrException e) {
throw e;
} catch (Exception e) {
throw new HerdrException("OnePaneHerdr decode failed for " + method, e);
}
}
@Override
public void close() {
}
}
}
@@ -161,34 +161,6 @@ class HerdrPeerLauncherAllowListWiringTest {
+ "it under allow: — sshAuthSock is unset here, so it defaults to block");
}
@Test
void brokerUriEnvStaysBlockedWhenListedInMemberCredentialsAllow() {
FakeHerdr herdr = new FakeHerdr();
WiringLauncher launcher = new WiringLauncher(herdr,
allowListWithAllow(List.of("BROKER_CONNECTION_URI")), "/bin/zsh", null,
() -> config("BROKER_CONNECTION_URI"));
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
Path dir = Path.of(launcher.env.get("ZDOTDIR"));
assertFalse(readAll(dir.resolve(EnvAllowListScrub.SCRUB_FILE)).contains("'BROKER_CONNECTION_URI'"),
"broker.uriEnv must not reach a member even when listed in memberCredentials.allow:");
}
@Test
void brokerUriEnvIsDeniedUnderTheDenyListPolicyEvenWhenAllowed() {
FakeHerdr herdr = new FakeHerdr();
WiringLauncher launcher = new WiringLauncher(herdr,
() -> new FleetConfig.MemberCredentials(null, List.of("BROKER_CONNECTION_URI"), List.of(), null),
"/bin/bash", null,
() -> config("BROKER_CONNECTION_URI"));
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
assertEquals("blocked-by-fleetd-cb596-see-gitea-issue-82", launcher.env.get("BROKER_CONNECTION_URI"),
"the deny-list overlay must deny broker.uriEnv even when allow: names it");
}
/**
* CB-633 follow-up criterion 3: on every allow-list spawn the daemon logs one INFO line, shaped
* "member credentials: allowed N of M", with real counts — not constants. Real path: the count
@@ -288,24 +260,15 @@ class HerdrPeerLauncherAllowListWiringTest {
/** Plus an injectable {@code hostEnvNames} source, for the "allowed N of M" log line test. */
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
Supplier<Set<String>> hostEnvNames) {
this(herdr, creds, shell, hostEnvNames, null);
}
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
Supplier<Set<String>> hostEnvNames, Supplier<FleetConfig> config) {
Supplier<Set<String>> hostEnvNames) {
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
Map.of("test", profile()), "test",
name -> "SHELL".equals(name) ? shell : null,
0, () -> 0L, () -> { }, null, creds, hostEnvNames, config);
0, () -> 0L, () -> { }, null, creds, hostEnvNames);
}
@Override
protected Launch buildLaunch(FleetConfig.Profile cfg, LaunchSpec spec) {
Map<String, String> launchEnv = baseEnv(cfg);
launchEnv.putAll(env);
env.clear();
env.putAll(launchEnv);
return new Launch(env, List.of("test"));
}
@@ -315,12 +278,6 @@ class HerdrPeerLauncherAllowListWiringTest {
}
}
private static FleetConfig config(String brokerUriEnv) {
return new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
new FleetConfig.Broker(null, brokerUriEnv, null), null, null, null, null, null,
null, null, null, null, null).withDefaults();
}
/** The generated directory is a temp directory; make sure the test does not leave a pile. */
@Test
void theGeneratedDirectoryIsRemovedWhenThePaneIsStopped() {
@@ -121,34 +121,6 @@ class MemberEnvAllowListTest {
assertTrue(derived.contains("OTHER_NAME"), "other allow: names are unaffected");
}
@Test
void configuredBrokerAndCoordinatorUriEnvNamesAreExcludedEvenWhenAllowed() {
FleetConfig config = config("BROKER_CONNECTION_URI", "COORDINATOR_CONNECTION_URI");
Set<String> excluded = MemberEnvAllowList.brokerUriEnvNames(config);
Set<String> derived = MemberEnvAllowList.derive(List.of(),
Set.of("BROKER_CONNECTION_URI", "COORDINATOR_CONNECTION_URI", "OTHER_NAME"), excluded);
assertFalse(derived.contains("BROKER_CONNECTION_URI"),
"broker.uriEnv is secret-bearing and must not ride in on allow:");
assertFalse(derived.contains("COORDINATOR_CONNECTION_URI"),
"coordinator.uriEnv has the same inline-password shape");
assertTrue(derived.contains("OTHER_NAME"), "unrelated allow: entries are unaffected");
}
@Test
void absentOrBlankBrokerUriEnvAddsNoExclusions() {
assertTrue(MemberEnvAllowList.brokerUriEnvNames(config(null, null)).isEmpty());
assertTrue(MemberEnvAllowList.brokerUriEnvNames(config(" ", "")).isEmpty());
}
private static FleetConfig config(String brokerUriEnv, String coordinatorUriEnv) {
return new FleetConfig(null, null, null, Map.of(), null, null, null, null, null,
new FleetConfig.Broker(null, brokerUriEnv, null), null, null, null, null, null,
null, null, null, null,
new FleetConfig.Coordinator(null, coordinatorUriEnv, null, null)).withDefaults();
}
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
@Test
void keepsMatchesExactlyPlusTheLocalePrefixRule() {
@@ -334,7 +334,8 @@ class OpenCodeLauncherTest {
assertNull(handle.agentSessionId(), "no record yet → null, not a spawn-time block");
// Once the record appears (here: same cwd), lazy discovery resolves it — the handle's
// session id matches its own worktree, not another's.
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_resolved", "/work/dir", 1000L);
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "p1", "ses_a.json",
"ses_resolved", "/work/dir", 1000L);
assertEquals("ses_resolved", handle.agentSessionId(),
"agentSessionId() re-scans and picks up a record that has since been written");
}
@@ -5,89 +5,68 @@ import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.SQLException;
import java.sql.Statement;
import java.nio.file.attribute.FileTime;
import static org.junit.jupiter.api.Assertions.*;
/**
* {@link OpenCodeSessionDiscovery} matches an opencode session row by the worker's cwd (its
* {@code directory}) against opencode's {@code opencode.db} SQLite database. These tests build a
* SYNTHETIC database themselves, in a JUnit temp directory — never the operator's real
* {@code ~/.local/share/opencode/opencode.db}, which a live opencode process may be writing.
* {@link OpenCodeSessionDiscovery} matches an opencode session record by the worker's cwd (its
* {@code directory}) against opencode's on-disk storage. These tests populate a TEMP storage root
* themselves — never the operator's real {@code ~/.local/share/opencode}.
*/
class OpenCodeSessionDiscoveryTest {
/**
* Create {@code <root>/opencode.db} with a minimal {@code session} table (just the columns
* {@link OpenCodeSessionDiscovery} reads: {@code id}, {@code directory}, {@code time_updated})
* and insert one row. Static so {@link OpenCodeLauncherTest} can reuse it.
* Write a session record {@code {"id":..., "directory":...}} under
* {@code <root>/session/<projectID>/<fileName>} and stamp it with a known last-modified time,
* so "most recently modified wins" is deterministic. Static so the launcher test can reuse it.
*/
static void writeRecord(Path root, String id, String directory, long timeUpdated) throws Exception {
Path db = root.resolve("opencode.db");
try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + db)) {
try (Statement statement = connection.createStatement()) {
statement.execute("CREATE TABLE IF NOT EXISTS session ("
+ "id TEXT PRIMARY KEY, directory TEXT, time_updated INTEGER)");
}
// Bound parameters, not string interpolation: the class under test uses a
// PreparedStatement, and a hand-escaped INSERT here is a pattern someone copies out.
try (PreparedStatement insert = connection.prepareStatement(
"INSERT INTO session (id, directory, time_updated) VALUES (?, ?, ?)")) {
insert.setString(1, id);
insert.setString(2, directory);
insert.setLong(3, timeUpdated);
insert.executeUpdate();
}
}
static void writeRecord(Path root, String projectId, String fileName, String id,
String directory, long lastModifiedEpochMillis) throws Exception {
Path dir = root.resolve("session").resolve(projectId);
Files.createDirectories(dir);
Path file = dir.resolve(fileName);
Files.writeString(file, "{\"id\":\"" + id + "\",\"directory\":\"" + directory
+ "\",\"projectID\":\"" + projectId + "\",\"version\":\"1.1.31\"}");
Files.setLastModifiedTime(file, FileTime.fromMillis(lastModifiedEpochMillis));
}
@Test
void findsTheRowWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L);
writeRecord(root, "ses_bbb", "/w/b", 2000L);
void findsTheRecordWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
writeRecord(root, "p2", "ses_b.json", "ses_bbb", "/w/b", 2000L);
assertEquals("ses_bbb", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/b"),
"the row whose directory equals the cwd is the one found");
"the record whose directory equals the cwd is the one found");
assertEquals("ses_aaa", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
}
@Test
void aNonMatchingDirectoryYieldsNullRatherThanAMismatch(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L);
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/other"),
"no row for this cwd yet → null, not a wrong session");
"no record for this cwd yet → null, not a wrong session");
}
@Test
void prefersTheMostRecentlyUpdatedRowWhenSeveralMatch(@TempDir Path root) throws Exception {
writeRecord(root, "ses_old", "/w/a", 1000L);
writeRecord(root, "ses_new", "/w/a", 5000L);
void prefersTheMostRecentlyModifiedRecordWhenSeveralMatch(@TempDir Path root) throws Exception {
writeRecord(root, "p1", "old.json", "ses_old", "/w/a", 1000L);
writeRecord(root, "p2", "new.json", "ses_new", "/w/a", 5000L);
assertEquals("ses_new", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
"the row with the highest time_updated for the cwd wins");
"the freshest record for the cwd wins");
}
@Test
void aMissingDatabaseYieldsNullWithoutThrowing(@TempDir Path root) {
// No opencode.db at all under the root.
void aMissingOrEmptyStorageRootYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
// Missing: no session dir at all under the root.
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
}
@Test
void anEmptyDatabaseYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
Path db = root.resolve("opencode.db");
try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + db);
Statement statement = connection.createStatement()) {
statement.execute("CREATE TABLE session (id TEXT PRIMARY KEY, directory TEXT, "
+ "time_updated INTEGER)");
}
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
// Present but empty: a session dir with nothing in it produces no match, not a throw.
Path emptyRoot = root.resolve("empty");
Files.createDirectories(emptyRoot.resolve("session"));
assertNull(new OpenCodeSessionDiscovery(emptyRoot).sessionIdForDirectory("/w/a"));
}
@Test
@@ -97,41 +76,15 @@ class OpenCodeSessionDiscoveryTest {
assertNull(discovery.sessionIdForDirectory(" "));
}
/**
* The one line standing between fleetd and writing the operator's live {@code opencode.db} —
* 841MB, with a running opencode writing it — is {@code config.setReadOnly(true)} in
* {@link OpenCodeSessionDiscovery#openReadOnly()}. Delete it and every other test in this class
* still passes, so this is the test that guards it.
*
* <p>It asks the connection to write, and requires a refusal. The obvious alternative — make
* the database file unwritable and check the read still works — proves nothing: SQLite silently
* downgrades a read-write open of an unwritable file to read-only, so that test passes either
* way. It was tried and watched pass with the flag removed.
*/
@Test
void theDatabaseIsOpenedReadOnly(@TempDir Path root) throws Exception {
writeRecord(root, "ses_aaa", "/w/a", 1000L);
void aMalformedRecordIsSkippedRatherThanFatal(@TempDir Path root) throws Exception {
// A record that fails to parse must not abort the scan of its siblings.
Path dir = root.resolve("session").resolve("p1");
Files.createDirectories(dir);
Files.writeString(dir.resolve("broken.json"), "{not valid json");
writeRecord(root, "p1", "good.json", "ses_good", "/w/a", 1000L);
try (Connection connection = new OpenCodeSessionDiscovery(root).openReadOnly();
Statement statement = connection.createStatement()) {
SQLException refused = assertThrows(SQLException.class,
() -> statement.executeUpdate("INSERT INTO session (id, directory, time_updated) "
+ "VALUES ('ses_zzz', '/w/z', 1)"),
"the connection must REFUSE a write — opencode is writing this database live");
assertTrue(refused.getMessage().toLowerCase().contains("readonly")
|| refused.getMessage().toLowerCase().contains("read-only"),
"the refusal must be about read-only, not some other error: " + refused.getMessage());
}
}
@Test
void aCorruptDatabaseFileYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
// A file at opencode.db that is not a SQLite database at all — the open/query must fail
// safe, never fatal to a spawn.
Path db = root.resolve("opencode.db");
Files.writeString(db, "this is not a sqlite database");
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
"an unreadable database resolves to null, not an exception");
assertEquals("ses_good", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
"an unreadable record is skipped; a later valid one still matches");
}
}