Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a052975420 | |||
| 3743789e8d | |||
| 388aba7632 | |||
| ad587eafa3 | |||
| 08ce9aef11 |
@@ -28,6 +28,7 @@
|
|||||||
<testcontainers.version>1.20.4</testcontainers.version>
|
<testcontainers.version>1.20.4</testcontainers.version>
|
||||||
<commons-compress.version>1.27.1</commons-compress.version>
|
<commons-compress.version>1.27.1</commons-compress.version>
|
||||||
<commons-lang3.version>3.18.0</commons-lang3.version>
|
<commons-lang3.version>3.18.0</commons-lang3.version>
|
||||||
|
<sqlite-jdbc.version>3.53.4.0</sqlite-jdbc.version>
|
||||||
</properties>
|
</properties>
|
||||||
|
|
||||||
<!--
|
<!--
|
||||||
@@ -44,6 +45,12 @@
|
|||||||
3.0-rc5; bumping Jackson 3 to the patched 3.2.x breaks the SDK (annotation mismatch).
|
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.
|
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.
|
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
|
<!-- Force the latest patched Jetty 11.x across all Javalin-pulled Jetty modules (no version
|
||||||
@@ -123,6 +130,17 @@
|
|||||||
<version>${amqp.version}</version>
|
<version>${amqp.version}</version>
|
||||||
</dependency>
|
</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 -->
|
<!-- Logging -->
|
||||||
<dependency>
|
<dependency>
|
||||||
<groupId>org.slf4j</groupId>
|
<groupId>org.slf4j</groupId>
|
||||||
|
|||||||
@@ -177,14 +177,14 @@ public final class Fleetd {
|
|||||||
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||||
() -> config.get().fleet(),
|
() -> config.get().fleet(),
|
||||||
() -> config.get().memberCredentials()));
|
() -> config.get().memberCredentials(), null, config::get));
|
||||||
}
|
}
|
||||||
if (!opencodeProfiles.isEmpty()) {
|
if (!opencodeProfiles.isEmpty()) {
|
||||||
adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(),
|
adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(),
|
||||||
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||||
() -> config.get().fleet(),
|
() -> config.get().fleet(),
|
||||||
() -> config.get().memberCredentials()));
|
() -> config.get().memberCredentials(), config::get));
|
||||||
}
|
}
|
||||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
||||||
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
||||||
|
|||||||
@@ -2,8 +2,11 @@ package dev.ltms.fleet.herdr;
|
|||||||
|
|
||||||
import com.fasterxml.jackson.databind.JsonNode;
|
import com.fasterxml.jackson.databind.JsonNode;
|
||||||
|
|
||||||
|
import java.util.LinkedHashSet;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
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
|
* Resolves which herdr pane a process belongs to — the herdr half of connection-based MCP
|
||||||
@@ -12,23 +15,46 @@ import java.util.Map;
|
|||||||
* is calling without the worker sending anything spoofable.
|
* 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}
|
* <p>herdr owns the PID→pane truth: {@code pane.process_info} reports each pane's {@code shell_pid}
|
||||||
* and foreground process PIDs. This scans agent panes; a spawn-time {@code pid→terminal} cache is
|
* and foreground process PIDs. A pid that is neither of those directly — e.g. a grandchild a
|
||||||
* the obvious optimization once wired into {@code ClaudeCodeLauncher}.
|
* 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}.
|
||||||
*
|
*
|
||||||
* <p>CB-185 split the fleet across two herdr daemons — lead operations on one, members on the
|
* <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
|
* 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
|
* 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)}
|
* 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
|
* 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).
|
* 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.
|
||||||
*/
|
*/
|
||||||
public final class PaneLocator {
|
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 List<HerdrClient> herdrs;
|
||||||
|
private final ParentResolver parentResolver;
|
||||||
|
|
||||||
/** Search only this client — the single-daemon deployment. */
|
/** Search only this client — the single-daemon deployment. */
|
||||||
public PaneLocator(HerdrClient herdr) {
|
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.herdrs = List.of(herdr);
|
||||||
|
this.parentResolver = parentResolver;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -37,7 +63,13 @@ public final class PaneLocator {
|
|||||||
* collapses to one client and one scan, exactly {@link #PaneLocator(HerdrClient)}'s behaviour.
|
* collapses to one client and one scan, exactly {@link #PaneLocator(HerdrClient)}'s behaviour.
|
||||||
*/
|
*/
|
||||||
public PaneLocator(HerdrClient lead, HerdrClient member) {
|
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.herdrs = lead == member ? List.of(lead) : List.of(lead, member);
|
||||||
|
this.parentResolver = parentResolver;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -49,8 +81,9 @@ public final class PaneLocator {
|
|||||||
if (pid <= 0) {
|
if (pid <= 0) {
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
Set<Long> ancestry = ancestorsOf(pid);
|
||||||
for (HerdrClient herdr : herdrs) {
|
for (HerdrClient herdr : herdrs) {
|
||||||
String terminal = terminalForPid(herdr, pid);
|
String terminal = terminalForPid(herdr, ancestry);
|
||||||
if (terminal != null) {
|
if (terminal != null) {
|
||||||
return terminal;
|
return terminal;
|
||||||
}
|
}
|
||||||
@@ -58,28 +91,54 @@ public final class PaneLocator {
|
|||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
private static String terminalForPid(HerdrClient herdr, long pid) {
|
/**
|
||||||
|
* {@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) {
|
||||||
for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) {
|
for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) {
|
||||||
String paneId = pane.path("pane_id").asText(null);
|
String paneId = pane.path("pane_id").asText(null);
|
||||||
if (paneId != null && paneOwnsPid(herdr, paneId, pid)) {
|
if (paneId != null && paneOwnsAnyOf(herdr, paneId, ancestry)) {
|
||||||
return pane.path("terminal_id").asText(null);
|
return pane.path("terminal_id").asText(null);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
private static boolean paneOwnsPid(HerdrClient herdr, String paneId, long pid) {
|
private static boolean paneOwnsAnyOf(HerdrClient herdr, String paneId, Set<Long> ancestry) {
|
||||||
JsonNode info;
|
JsonNode info;
|
||||||
try {
|
try {
|
||||||
info = herdr.call("pane.process_info", Map.of("pane_id", paneId)).path("process_info");
|
info = herdr.call("pane.process_info", Map.of("pane_id", paneId)).path("process_info");
|
||||||
} catch (HerdrException e) {
|
} catch (HerdrException e) {
|
||||||
return false; // pane vanished mid-scan — just skip it
|
return false; // pane vanished mid-scan — just skip it
|
||||||
}
|
}
|
||||||
if (info.path("shell_pid").asLong(-1) == pid) {
|
if (ancestry.contains(info.path("shell_pid").asLong(-1))) {
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
for (JsonNode p : info.path("foreground_processes")) {
|
for (JsonNode p : info.path("foreground_processes")) {
|
||||||
if (p.path("pid").asLong(-1) == pid) {
|
if (ancestry.contains(p.path("pid").asLong(-1))) {
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,21 @@
|
|||||||
|
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,13 +94,26 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
|||||||
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
||||||
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
|
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
|
||||||
Function<String, String> env,
|
Function<String, String> env,
|
||||||
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
||||||
Supplier<FleetConfig.Fleet> fleet,
|
Supplier<FleetConfig.Fleet> fleet,
|
||||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
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) {
|
||||||
this(agents, spaces, guard, profiles, defaultProfile, env,
|
this(agents, spaces, guard, profiles, defaultProfile, env,
|
||||||
spawnReadyTimeoutMs,
|
spawnReadyTimeoutMs,
|
||||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||||
fleet, memberCredentials);
|
fleet, memberCredentials, hostEnvNames, config);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -153,10 +166,22 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
|||||||
Function<String, String> env,
|
Function<String, String> env,
|
||||||
long spawnReadyTimeoutMs,
|
long spawnReadyTimeoutMs,
|
||||||
LongSupplier nowMillis, Runnable sleeper,
|
LongSupplier nowMillis, Runnable sleeper,
|
||||||
Supplier<FleetConfig.Fleet> fleet,
|
Supplier<FleetConfig.Fleet> fleet,
|
||||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
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) {
|
||||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
|
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, config);
|
||||||
this.guard = guard;
|
this.guard = guard;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -174,7 +199,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
|||||||
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
Supplier<FleetConfig.MemberCredentials> memberCredentials,
|
||||||
Supplier<Set<String>> hostEnvNames) {
|
Supplier<Set<String>> hostEnvNames) {
|
||||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames);
|
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, null);
|
||||||
this.guard = guard;
|
this.guard = guard;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -170,6 +170,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
|
|
||||||
/** Guards {@link #warnNonZsh} to one WARN per launcher instance, not one per spawn. */
|
/** Guards {@link #warnNonZsh} to one WARN per launcher instance, not one per spawn. */
|
||||||
private final AtomicBoolean nonZshShellWarned = new AtomicBoolean();
|
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)
|
* @param namePrefix label prefix for this peer kind (drives naming and reap)
|
||||||
@@ -238,8 +240,19 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
long spawnReadyTimeoutMs,
|
long spawnReadyTimeoutMs,
|
||||||
LongSupplier nowMillis, Runnable sleeper,
|
LongSupplier nowMillis, Runnable sleeper,
|
||||||
Supplier<FleetConfig.Fleet> fleet,
|
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<FleetConfig.MemberCredentials> memberCredentials,
|
||||||
Supplier<Set<String>> hostEnvNames) {
|
Supplier<Set<String>> hostEnvNames, Supplier<FleetConfig> config) {
|
||||||
this.fleet = fleet;
|
this.fleet = fleet;
|
||||||
this.namePrefix = namePrefix;
|
this.namePrefix = namePrefix;
|
||||||
this.agents = agents;
|
this.agents = agents;
|
||||||
@@ -252,6 +265,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
this.sleeper = sleeper;
|
this.sleeper = sleeper;
|
||||||
this.memberCredentials = memberCredentials;
|
this.memberCredentials = memberCredentials;
|
||||||
this.hostEnvNames = hostEnvNames != null ? hostEnvNames : () -> System.getenv().keySet();
|
this.hostEnvNames = hostEnvNames != null ? hostEnvNames : () -> System.getenv().keySet();
|
||||||
|
this.config = config;
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- adapter seams -------------------------------------------------------------------------
|
// --- adapter seams -------------------------------------------------------------------------
|
||||||
@@ -1028,9 +1042,11 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/** Put {@link #BLOCKED_CREDENTIAL_SENTINEL} over every blocked name in the pane-creation env map. */
|
/** Put {@link #BLOCKED_CREDENTIAL_SENTINEL} over every blocked name in the pane-creation env map. */
|
||||||
private static void overlayBlockedCredentials(Map<String, String> workerEnv,
|
private void overlayBlockedCredentials(Map<String, String> workerEnv,
|
||||||
FleetConfig.MemberCredentials creds) {
|
FleetConfig.MemberCredentials creds) {
|
||||||
for (String name : creds.blockedSet()) {
|
Set<String> blocked = new java.util.TreeSet<>(creds.blockedSet());
|
||||||
|
blocked.addAll(brokerUriEnvNames());
|
||||||
|
for (String name : blocked) {
|
||||||
workerEnv.put(name, BLOCKED_CREDENTIAL_SENTINEL);
|
workerEnv.put(name, BLOCKED_CREDENTIAL_SENTINEL);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -1097,15 +1113,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
* ssh-agent handle when explicitly allowed, and the exact keys of THIS launch's own env map.
|
* 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) {
|
private Set<String> derivedAllowedNames(FleetConfig.MemberCredentials creds, Launch launch) {
|
||||||
|
Set<String> brokerUriEnvNames = brokerUriEnvNames();
|
||||||
Set<String> allowed = new java.util.TreeSet<>(
|
Set<String> allowed = new java.util.TreeSet<>(
|
||||||
MemberEnvAllowList.derive(profiles.values(), creds.allowSet()));
|
MemberEnvAllowList.derive(profiles.values(), creds.allowSet(), brokerUriEnvNames));
|
||||||
if (creds.sshAuthSockAllowed()) {
|
if (creds.sshAuthSockAllowed()) {
|
||||||
allowed.add(SSH_AUTH_SOCK);
|
allowed.add(SSH_AUTH_SOCK);
|
||||||
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
|
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
|
||||||
allowed.addAll(launch.env().keySet());
|
allowed.addAll(launch.env().keySet());
|
||||||
|
allowed.removeAll(brokerUriEnvNames);
|
||||||
return allowed;
|
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
|
* 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
|
* can read a single log line and know the scrub ran and how much of the visible environment it
|
||||||
|
|||||||
@@ -40,8 +40,9 @@ import java.util.TreeSet;
|
|||||||
* operator's own explicit list. Before this, {@code policy: allow-list} silently ignored every name
|
* 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
|
* 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
|
* turning the policy on could blank credentials working members already depended on. {@code
|
||||||
* SSH_AUTH_SOCK} is the one exception: even when the operator lists it under {@code allow:}, it is
|
* SSH_AUTH_SOCK} and configured broker URI environment names are exceptions: even when the operator
|
||||||
* excluded here and added back ONLY by the caller when {@code sshAuthSock: allow} is explicitly set
|
* 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
|
||||||
* (see {@link #SSH_AUTH_SOCK}'s javadoc) — it is a live handle to the operator's own ssh-agent, not
|
* (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
|
* 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
|
* operator's agent holds the moment they typed the name under {@code allow:} for an unrelated
|
||||||
@@ -106,6 +107,16 @@ public final class MemberEnvAllowList {
|
|||||||
* run-to-run.
|
* run-to-run.
|
||||||
*/
|
*/
|
||||||
public static Set<String> derive(Collection<FleetConfig.Profile> profiles, Set<String> configuredAllow) {
|
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);
|
Set<String> derived = new TreeSet<>(INFRASTRUCTURE_PASSTHROUGH);
|
||||||
if (profiles != null) {
|
if (profiles != null) {
|
||||||
for (FleetConfig.Profile p : profiles) {
|
for (FleetConfig.Profile p : profiles) {
|
||||||
@@ -124,9 +135,37 @@ public final class MemberEnvAllowList {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
if (excludedNames != null) {
|
||||||
|
derived.removeAll(excludedNames);
|
||||||
|
}
|
||||||
return Set.copyOf(derived);
|
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
|
* 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
|
* match, or an infrastructure-prefixed name ({@code LC_*}). Prefix rules live ONLY here and in
|
||||||
@@ -146,4 +185,16 @@ public final class MemberEnvAllowList {
|
|||||||
into.add(name);
|
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,12 +116,23 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
|||||||
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
||||||
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
|
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
|
||||||
Function<String, String> env,
|
Function<String, String> env,
|
||||||
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
||||||
Supplier<FleetConfig.Fleet> fleet,
|
Supplier<FleetConfig.Fleet> fleet,
|
||||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
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) {
|
||||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials);
|
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, config);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -182,8 +193,20 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
|||||||
Path configRoot, Path discoveryRoot,
|
Path configRoot, Path discoveryRoot,
|
||||||
Supplier<FleetConfig.Fleet> fleet,
|
Supplier<FleetConfig.Fleet> fleet,
|
||||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
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,
|
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
|
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, null, config);
|
||||||
this.configRoot = configRoot;
|
this.configRoot = configRoot;
|
||||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,58 +1,87 @@
|
|||||||
package dev.ltms.fleet.member;
|
package dev.ltms.fleet.member;
|
||||||
|
|
||||||
import com.fasterxml.jackson.databind.JsonNode;
|
import org.slf4j.Logger;
|
||||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
import org.slf4j.LoggerFactory;
|
||||||
|
import org.sqlite.SQLiteConfig;
|
||||||
|
|
||||||
import java.io.IOException;
|
|
||||||
import java.nio.file.Files;
|
import java.nio.file.Files;
|
||||||
import java.nio.file.Path;
|
import java.nio.file.Path;
|
||||||
import java.util.stream.Stream;
|
import java.sql.Connection;
|
||||||
|
import java.sql.PreparedStatement;
|
||||||
|
import java.sql.ResultSet;
|
||||||
|
import java.sql.SQLException;
|
||||||
|
import java.util.concurrent.atomic.AtomicBoolean;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Resolves the opencode session id for a fleetd worker from opencode's on-disk storage — the
|
* 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>
|
* only place this adapter touches opencode's private layout, and deliberately the <em>only</em>
|
||||||
* class that does.
|
* class that does.
|
||||||
*
|
*
|
||||||
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not a
|
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not
|
||||||
* stable contract: opencode writes one JSON file per session under
|
* a stable contract: opencode persists its session state in a SQLite database at
|
||||||
* {@code <storageRoot>/session/<projectID>/<ses_*.json>}, and each record carries a
|
* {@code <storageRoot>/opencode.db} (a {@code session} table, one row per session, keyed by id and
|
||||||
* {@code "version"} field (e.g. {@code "1.1.31"}), so the exact directory shape, file naming, and
|
* carrying a {@code directory} column). That schema can move between opencode releases exactly
|
||||||
* field names can move between opencode releases. opencode also ships a headless HTTP server that
|
* like the JSON-file layout it replaced did (opencode migrated off a one-JSON-file-per-session
|
||||||
* may supersede file scanning entirely. Everything this adapter knows about that private storage —
|
* tree under {@code <storageRoot>/storage/session/<projectID>/ses_*.json} in January 2026 — that
|
||||||
* its shape, naming, and field names — lives here, so a layout change, or a switch to the HTTP
|
* tree is now a frozen migration artefact nothing writes, which is why this class no longer reads
|
||||||
* server, changes exactly one class and nothing in {@link OpenCodeLauncher}.
|
* 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>The determinism that makes this useful is structural, not a guess: every fleetd worker runs
|
* <p>The determinism that makes this useful is structural, not a guess: every fleetd worker runs
|
||||||
* in its own unique git worktree, so the record's {@code directory} (its project root) equals the
|
* in its own unique git worktree, so the row's {@code directory} (its project root) equals the
|
||||||
* worker's cwd identifies <em>its</em> session unambiguously. We match on {@code directory} rather
|
* 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
|
* than diffing {@code opencode session list} before/after — that races under concurrent spawns, and
|
||||||
* the CLI listing does not even show the directory.
|
* the CLI listing does not even show the directory.
|
||||||
*
|
*
|
||||||
* <p>All reads are best-effort and never throw: a missing or unreadable storage root, a record that
|
* <p>All reads are best-effort and never throw: a missing or unreadable database, a query that
|
||||||
* fails to parse, or a directory with no record yet all yield {@code null}, and the caller (the
|
* fails, or a directory with no row yet all yield {@code null}, and the caller (the session
|
||||||
* session handle) treats that as "identity not resolved yet" and retries later.
|
* 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.
|
||||||
*/
|
*/
|
||||||
final class OpenCodeSessionDiscovery {
|
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 storageRoot; // e.g. ~/.local/share/opencode (injectable for tests)
|
||||||
private final ObjectMapper json;
|
private final Path databasePath;
|
||||||
|
private final AtomicBoolean warnedMissingDatabase = new AtomicBoolean(false);
|
||||||
|
|
||||||
OpenCodeSessionDiscovery(Path storageRoot) {
|
OpenCodeSessionDiscovery(Path storageRoot) {
|
||||||
this.storageRoot = storageRoot;
|
this.storageRoot = storageRoot;
|
||||||
this.json = new ObjectMapper();
|
this.databasePath = storageRoot.resolve("opencode.db");
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* The opencode session id whose record references {@code directory} (the worker's cwd), or
|
* A connection to {@link #databasePath} opened with SQLite's {@code SQLITE_OPEN_READONLY}
|
||||||
* {@code null} when no record matches yet. When several records share the directory — e.g.
|
* flag: it never creates the file, never writes, and never touches WAL or journal mode.
|
||||||
* repeated spawns into the same worktree — the <em>most recently modified</em> one wins: it is
|
* 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 session the pane most likely corresponds to.
|
* the session the pane most likely corresponds to.
|
||||||
*
|
*
|
||||||
* <p>Never throws: a missing {@code storageRoot}, an unreadable/malformed record, or a
|
* <p>Never throws: a missing {@code opencode.db}, a locked/unreadable database, a query
|
||||||
* directory that has not been persisted yet all resolve to {@code null} rather than failing a
|
* failure, or a directory that has not been persisted yet all resolve to {@code null} rather
|
||||||
* spawn. A fleetd worker's session record is written lazily (when the session is first
|
* than failing a spawn. A fleetd worker's session row is written lazily (when the session is
|
||||||
* persisted), so {@code null} here is the normal answer right after the pane is ready, and the
|
* first persisted), so {@code null} here is the normal answer right after the pane is ready,
|
||||||
* caller retries later.
|
* and the caller retries later.
|
||||||
*
|
*
|
||||||
* @param directory the worker's cwd, as resolved for this spawn
|
* @param directory the worker's cwd, as resolved for this spawn
|
||||||
* @return the matching session id, or {@code null} if none is known yet
|
* @return the matching session id, or {@code null} if none is known yet
|
||||||
@@ -61,62 +90,31 @@ final class OpenCodeSessionDiscovery {
|
|||||||
if (directory == null || directory.isBlank()) {
|
if (directory == null || directory.isBlank()) {
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
Path sessionRoot = storageRoot.resolve("session");
|
if (!Files.isRegularFile(databasePath)) {
|
||||||
if (!Files.isDirectory(sessionRoot)) {
|
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);
|
||||||
|
}
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
String best = null;
|
String sql = "SELECT id FROM session WHERE directory = ? ORDER BY time_updated DESC LIMIT 1";
|
||||||
long bestMtime = Long.MIN_VALUE;
|
try (Connection connection = openReadOnly();
|
||||||
try (Stream<Path> projectDirs = Files.list(sessionRoot)) {
|
PreparedStatement statement = connection.prepareStatement(sql)) {
|
||||||
for (Path projectDir : projectDirs.filter(Files::isDirectory).toList()) {
|
statement.setString(1, directory);
|
||||||
try (Stream<Path> records = Files.list(projectDir)) {
|
try (ResultSet rows = statement.executeQuery()) {
|
||||||
for (Path record : records.toList()) {
|
if (rows.next()) {
|
||||||
String id = matchId(record, directory);
|
return rows.getString("id");
|
||||||
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 (IOException ignored) {
|
} catch (SQLException e) {
|
||||||
// storage root vanished or became unreadable — "no session known yet"
|
// 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());
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
return best;
|
log.debug("no opencode session row for directory (root={}, directory={})",
|
||||||
}
|
storageRoot, directory);
|
||||||
|
return null;
|
||||||
/**
|
|
||||||
* 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;
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -366,21 +366,40 @@ public final class MessageService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Route a worker's explicit {@code fleet_reply}: resolve an open send, or queue it in the
|
* Route a worker's explicit {@code fleet_reply}: resolve an open send, complete an async ticket
|
||||||
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
|
* still parked waiting on this exact turn's answer, or — only once neither applies — queue it in
|
||||||
* result is <em>not</em> a failure — the reply is held for later drain.
|
* the inbox. Unlike the bare {@link Rendezvous#resolve}, a no-waiter result is <em>not</em> a
|
||||||
|
* failure — the reply is held for later drain.
|
||||||
*
|
*
|
||||||
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code fleet_ask} /
|
* <p><strong>Do NOT use this for mid-turn questions.</strong> {@code fleet_ask} /
|
||||||
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
|
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
|
||||||
* are interactive and must never be queued.
|
* are interactive and must never be queued.
|
||||||
*
|
*
|
||||||
* @return always {@code true} — the reply either resolved a live send or was queued
|
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
|
||||||
|
* was queued
|
||||||
*/
|
*/
|
||||||
public boolean reply(String session, String content) {
|
public boolean reply(String session, String content) {
|
||||||
if (rendezvous.resolve(session, content)) {
|
if (rendezvous.resolve(session, content)) {
|
||||||
count(FleetMetrics.REPLIES, "path", "rendezvous");
|
count(FleetMetrics.REPLIES, "path", "rendezvous");
|
||||||
return true; // a live send took it — unchanged fast path
|
return true; // a live send took it — unchanged fast path
|
||||||
}
|
}
|
||||||
|
// #137: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming a
|
||||||
|
// turn that {@link #answer} already gave up waiting on. answer()'s own bounded wait (the
|
||||||
|
// primary's fleet_send{turnId} call, capped well under a minute) can time out and close its
|
||||||
|
// waiter long before the worker — now actually resuming real work — finishes and replies. That
|
||||||
|
// reply used to have nowhere to land but the session inbox, leaving the async ticket's future
|
||||||
|
// unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's abandon() forced it
|
||||||
|
// FAILED with a misleading "session released before it replied" reason, even though the reply
|
||||||
|
// had, in fact, arrived. Completing the matching ticket directly here means fleet_poll{ticket}
|
||||||
|
// sees the real reply instead.
|
||||||
|
Task orphan = askAnsweredAsyncTask(session);
|
||||||
|
if (orphan != null && orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
|
||||||
|
if (orphan.turnId != null) {
|
||||||
|
asyncTasksByTurn.remove(orphan.turnId, orphan);
|
||||||
|
}
|
||||||
|
count(FleetMetrics.REPLIES, "path", "async-recovered");
|
||||||
|
return true; // the ticket itself took it — no inbox stranding at all
|
||||||
|
}
|
||||||
inbox.publish(session, UUID.randomUUID().toString(), content);
|
inbox.publish(session, UUID.randomUUID().toString(), content);
|
||||||
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
|
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
|
||||||
// worker whose replies keep missing their waiter, not only the queue depth this leaves behind.
|
// worker whose replies keep missing their waiter, not only the queue depth this leaves behind.
|
||||||
@@ -394,6 +413,24 @@ public final class MessageService {
|
|||||||
return true; // held, not lost
|
return true; // held, not lost
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The still-open async task on {@code target} whose {@code fleet_ask} was already answered — its
|
||||||
|
* {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer} —
|
||||||
|
* yet whose future is not resolved yet (#137). {@code null} if no such task exists, including the
|
||||||
|
* common case where {@code target}'s worker never used {@code fleet_ask} at all (a task that was
|
||||||
|
* never asked has {@code turnId == null}, so it can never match here and only ever completes
|
||||||
|
* through the ordinary rendezvous fast path in {@link #reply}).
|
||||||
|
*/
|
||||||
|
private Task askAnsweredAsyncTask(String target) {
|
||||||
|
for (Task task : tasks.values()) {
|
||||||
|
if (target.equals(task.target) && task.question == null && task.turnId != null
|
||||||
|
&& !task.future.isDone()) {
|
||||||
|
return task;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
|
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
|
||||||
private void count(String name, String... labels) {
|
private void count(String name, String... labels) {
|
||||||
if (metrics != null) {
|
if (metrics != null) {
|
||||||
@@ -435,19 +472,43 @@ public final class MessageService {
|
|||||||
* <p>Resolving the waiter as a failure — rather than letting it time out — also means the
|
* <p>Resolving the waiter as a failure — rather than letting it time out — also means the
|
||||||
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
|
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
|
||||||
*
|
*
|
||||||
* @return true if a live waiter was failed
|
* <p><strong>#137 defence in depth.</strong> {@link #reply} already hands a worker's real
|
||||||
|
* {@code fleet_reply} straight to the async ticket it belongs to whenever one is still parked
|
||||||
|
* waiting for it (see {@link #askAnsweredAsyncTask}), so by the time a session is released its
|
||||||
|
* tasks are normally already resolved — this loop's {@code complete} calls are then harmless
|
||||||
|
* no-ops (a {@link CompletableFuture} can only resolve once). But should some other path someday
|
||||||
|
* strand a reply in the inbox without completing its ticket, checking
|
||||||
|
* {@link #hasStrandedReply(String)} here — before ever writing a failure — means a torn-down
|
||||||
|
* session whose worker in fact replied is still reported {@code REPLIED} with that reply's own
|
||||||
|
* text, never the misleading "the worker session was released before it replied" (which also
|
||||||
|
* means the snapshot/worktree recovery hint that follows it never prints once a reply exists).
|
||||||
|
*
|
||||||
|
* @return true if a live waiter or an async task was failed (never true for one recovered as a
|
||||||
|
* reply — see the note above)
|
||||||
*/
|
*/
|
||||||
public boolean abandon(String target, String reason) {
|
public boolean abandon(String target, String reason) {
|
||||||
|
boolean hadStrandedReply = hasStrandedReply(target);
|
||||||
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
|
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
|
||||||
strandedReplies.remove(target);
|
strandedReplies.remove(target);
|
||||||
queuedDeliveries.remove(target);
|
queuedDeliveries.remove(target);
|
||||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
||||||
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
|
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
|
||||||
boolean asyncFailed = false;
|
boolean asyncFailed = false;
|
||||||
|
Reply recovered = null; // lazily drained at most once, only if a task actually needs it
|
||||||
for (Task task : tasks.values()) {
|
for (Task task : tasks.values()) {
|
||||||
if (target.equals(task.target) && task.question == null
|
if (!target.equals(task.target) || task.question != null || task.future.isDone()) {
|
||||||
&& task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))) {
|
continue;
|
||||||
asyncFailed = true;
|
}
|
||||||
|
if (hadStrandedReply && recovered == null) {
|
||||||
|
recovered = recoverStrandedReply(target);
|
||||||
|
}
|
||||||
|
Reply outcome = recovered != null ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
|
||||||
|
if (task.future.complete(outcome)) {
|
||||||
|
if (outcome.outcome() == Outcome.WORKER_FAILED) {
|
||||||
|
asyncFailed = true;
|
||||||
|
} else if (task.turnId != null) {
|
||||||
|
asyncTasksByTurn.remove(task.turnId, task);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (failed) {
|
if (failed) {
|
||||||
@@ -456,6 +517,21 @@ public final class MessageService {
|
|||||||
return failed || asyncFailed;
|
return failed || asyncFailed;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Drain {@code target}'s inbox and hand its content back as a {@link Outcome#REPLIED} result
|
||||||
|
* (#137 defence in depth for {@link #abandon}) — {@code null} if it turned out empty (the
|
||||||
|
* stranding fact raced away, e.g. a lead's own {@code fleet_poll} on the raw session already
|
||||||
|
* drained it first). When more than one message is queued, only the newest is the worker's actual
|
||||||
|
* final answer ({@link #drainReplies} returns them oldest-first).
|
||||||
|
*/
|
||||||
|
private Reply recoverStrandedReply(String target) {
|
||||||
|
var messages = drainReplies(target);
|
||||||
|
if (messages.isEmpty()) {
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
return new Reply(Outcome.REPLIED, messages.get(messages.size() - 1).content());
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
|
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
|
||||||
* so that a subsequent drain or peek no longer returns it.
|
* so that a subsequent drain or peek no longer returns it.
|
||||||
|
|||||||
@@ -0,0 +1,27 @@
|
|||||||
|
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,7 +1,11 @@
|
|||||||
package dev.ltms.fleet.herdr;
|
package dev.ltms.fleet.herdr;
|
||||||
|
|
||||||
|
import com.fasterxml.jackson.databind.JsonNode;
|
||||||
|
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||||
import org.junit.jupiter.api.Test;
|
import org.junit.jupiter.api.Test;
|
||||||
|
|
||||||
|
import java.util.concurrent.atomic.AtomicInteger;
|
||||||
|
|
||||||
import static org.junit.jupiter.api.Assertions.*;
|
import static org.junit.jupiter.api.Assertions.*;
|
||||||
|
|
||||||
/** Unit tests for PID → pane resolution (the herdr half of connection-based MCP identity). */
|
/** Unit tests for PID → pane resolution (the herdr half of connection-based MCP identity). */
|
||||||
@@ -63,4 +67,123 @@ class PaneLocatorTest {
|
|||||||
long paneListCalls = shared.calls.stream().filter(c -> c.method().equals("pane.list")).count();
|
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");
|
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() {
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+45
-2
@@ -161,6 +161,34 @@ class HerdrPeerLauncherAllowListWiringTest {
|
|||||||
+ "it under allow: — sshAuthSock is unset here, so it defaults to block");
|
+ "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
|
* 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
|
* "member credentials: allowed N of M", with real counts — not constants. Real path: the count
|
||||||
@@ -260,15 +288,24 @@ class HerdrPeerLauncherAllowListWiringTest {
|
|||||||
|
|
||||||
/** Plus an injectable {@code hostEnvNames} source, for the "allowed N of M" log line test. */
|
/** Plus an injectable {@code hostEnvNames} source, for the "allowed N of M" log line test. */
|
||||||
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
|
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
|
||||||
Supplier<Set<String>> hostEnvNames) {
|
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) {
|
||||||
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
|
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
Map.of("test", profile()), "test",
|
Map.of("test", profile()), "test",
|
||||||
name -> "SHELL".equals(name) ? shell : null,
|
name -> "SHELL".equals(name) ? shell : null,
|
||||||
0, () -> 0L, () -> { }, null, creds, hostEnvNames);
|
0, () -> 0L, () -> { }, null, creds, hostEnvNames, config);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
protected Launch buildLaunch(FleetConfig.Profile cfg, LaunchSpec spec) {
|
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"));
|
return new Launch(env, List.of("test"));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -278,6 +315,12 @@ 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. */
|
/** The generated directory is a temp directory; make sure the test does not leave a pile. */
|
||||||
@Test
|
@Test
|
||||||
void theGeneratedDirectoryIsRemovedWhenThePaneIsStopped() {
|
void theGeneratedDirectoryIsRemovedWhenThePaneIsStopped() {
|
||||||
|
|||||||
@@ -121,6 +121,34 @@ class MemberEnvAllowListTest {
|
|||||||
assertTrue(derived.contains("OTHER_NAME"), "other allow: names are unaffected");
|
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. */
|
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
|
||||||
@Test
|
@Test
|
||||||
void keepsMatchesExactlyPlusTheLocalePrefixRule() {
|
void keepsMatchesExactlyPlusTheLocalePrefixRule() {
|
||||||
|
|||||||
@@ -334,8 +334,7 @@ class OpenCodeLauncherTest {
|
|||||||
assertNull(handle.agentSessionId(), "no record yet → null, not a spawn-time block");
|
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
|
// Once the record appears (here: same cwd), lazy discovery resolves it — the handle's
|
||||||
// session id matches its own worktree, not another's.
|
// session id matches its own worktree, not another's.
|
||||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "p1", "ses_a.json",
|
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_resolved", "/work/dir", 1000L);
|
||||||
"ses_resolved", "/work/dir", 1000L);
|
|
||||||
assertEquals("ses_resolved", handle.agentSessionId(),
|
assertEquals("ses_resolved", handle.agentSessionId(),
|
||||||
"agentSessionId() re-scans and picks up a record that has since been written");
|
"agentSessionId() re-scans and picks up a record that has since been written");
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,68 +5,89 @@ import org.junit.jupiter.api.io.TempDir;
|
|||||||
|
|
||||||
import java.nio.file.Files;
|
import java.nio.file.Files;
|
||||||
import java.nio.file.Path;
|
import java.nio.file.Path;
|
||||||
import java.nio.file.attribute.FileTime;
|
import java.sql.Connection;
|
||||||
|
import java.sql.DriverManager;
|
||||||
|
import java.sql.PreparedStatement;
|
||||||
|
import java.sql.SQLException;
|
||||||
|
import java.sql.Statement;
|
||||||
|
|
||||||
import static org.junit.jupiter.api.Assertions.*;
|
import static org.junit.jupiter.api.Assertions.*;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* {@link OpenCodeSessionDiscovery} matches an opencode session record by the worker's cwd (its
|
* {@link OpenCodeSessionDiscovery} matches an opencode session row by the worker's cwd (its
|
||||||
* {@code directory}) against opencode's on-disk storage. These tests populate a TEMP storage root
|
* {@code directory}) against opencode's {@code opencode.db} SQLite database. These tests build a
|
||||||
* themselves — never the operator's real {@code ~/.local/share/opencode}.
|
* 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.
|
||||||
*/
|
*/
|
||||||
class OpenCodeSessionDiscoveryTest {
|
class OpenCodeSessionDiscoveryTest {
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Write a session record {@code {"id":..., "directory":...}} under
|
* Create {@code <root>/opencode.db} with a minimal {@code session} table (just the columns
|
||||||
* {@code <root>/session/<projectID>/<fileName>} and stamp it with a known last-modified time,
|
* {@link OpenCodeSessionDiscovery} reads: {@code id}, {@code directory}, {@code time_updated})
|
||||||
* so "most recently modified wins" is deterministic. Static so the launcher test can reuse it.
|
* and insert one row. Static so {@link OpenCodeLauncherTest} can reuse it.
|
||||||
*/
|
*/
|
||||||
static void writeRecord(Path root, String projectId, String fileName, String id,
|
static void writeRecord(Path root, String id, String directory, long timeUpdated) throws Exception {
|
||||||
String directory, long lastModifiedEpochMillis) throws Exception {
|
Path db = root.resolve("opencode.db");
|
||||||
Path dir = root.resolve("session").resolve(projectId);
|
try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + db)) {
|
||||||
Files.createDirectories(dir);
|
try (Statement statement = connection.createStatement()) {
|
||||||
Path file = dir.resolve(fileName);
|
statement.execute("CREATE TABLE IF NOT EXISTS session ("
|
||||||
Files.writeString(file, "{\"id\":\"" + id + "\",\"directory\":\"" + directory
|
+ "id TEXT PRIMARY KEY, directory TEXT, time_updated INTEGER)");
|
||||||
+ "\",\"projectID\":\"" + projectId + "\",\"version\":\"1.1.31\"}");
|
}
|
||||||
Files.setLastModifiedTime(file, FileTime.fromMillis(lastModifiedEpochMillis));
|
// 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();
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void findsTheRecordWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
|
void findsTheRowWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
|
||||||
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
|
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||||
writeRecord(root, "p2", "ses_b.json", "ses_bbb", "/w/b", 2000L);
|
writeRecord(root, "ses_bbb", "/w/b", 2000L);
|
||||||
|
|
||||||
assertEquals("ses_bbb", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/b"),
|
assertEquals("ses_bbb", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/b"),
|
||||||
"the record whose directory equals the cwd is the one found");
|
"the row whose directory equals the cwd is the one found");
|
||||||
assertEquals("ses_aaa", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
assertEquals("ses_aaa", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void aNonMatchingDirectoryYieldsNullRatherThanAMismatch(@TempDir Path root) throws Exception {
|
void aNonMatchingDirectoryYieldsNullRatherThanAMismatch(@TempDir Path root) throws Exception {
|
||||||
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
|
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||||
|
|
||||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/other"),
|
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/other"),
|
||||||
"no record for this cwd yet → null, not a wrong session");
|
"no row for this cwd yet → null, not a wrong session");
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void prefersTheMostRecentlyModifiedRecordWhenSeveralMatch(@TempDir Path root) throws Exception {
|
void prefersTheMostRecentlyUpdatedRowWhenSeveralMatch(@TempDir Path root) throws Exception {
|
||||||
writeRecord(root, "p1", "old.json", "ses_old", "/w/a", 1000L);
|
writeRecord(root, "ses_old", "/w/a", 1000L);
|
||||||
writeRecord(root, "p2", "new.json", "ses_new", "/w/a", 5000L);
|
writeRecord(root, "ses_new", "/w/a", 5000L);
|
||||||
|
|
||||||
assertEquals("ses_new", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
assertEquals("ses_new", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||||
"the freshest record for the cwd wins");
|
"the row with the highest time_updated for the cwd wins");
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void aMissingOrEmptyStorageRootYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
|
void aMissingDatabaseYieldsNullWithoutThrowing(@TempDir Path root) {
|
||||||
// Missing: no session dir at all under the root.
|
// No opencode.db at all under the root.
|
||||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
||||||
|
}
|
||||||
|
|
||||||
// Present but empty: a session dir with nothing in it produces no match, not a throw.
|
@Test
|
||||||
Path emptyRoot = root.resolve("empty");
|
void anEmptyDatabaseYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
|
||||||
Files.createDirectories(emptyRoot.resolve("session"));
|
Path db = root.resolve("opencode.db");
|
||||||
assertNull(new OpenCodeSessionDiscovery(emptyRoot).sessionIdForDirectory("/w/a"));
|
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"));
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
@@ -76,15 +97,41 @@ class OpenCodeSessionDiscoveryTest {
|
|||||||
assertNull(discovery.sessionIdForDirectory(" "));
|
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
|
@Test
|
||||||
void aMalformedRecordIsSkippedRatherThanFatal(@TempDir Path root) throws Exception {
|
void theDatabaseIsOpenedReadOnly(@TempDir Path root) throws Exception {
|
||||||
// A record that fails to parse must not abort the scan of its siblings.
|
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||||
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);
|
|
||||||
|
|
||||||
assertEquals("ses_good", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
try (Connection connection = new OpenCodeSessionDiscovery(root).openReadOnly();
|
||||||
"an unreadable record is skipped; a later valid one still matches");
|
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");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -710,6 +710,68 @@ class MessageServiceTest {
|
|||||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --- #137: a fleet_ask round-trip must not orphan the ticket's own reply -------------------
|
||||||
|
//
|
||||||
|
// The primary's fleet_send{turnId} answer call is itself bounded (a real MCP call, capped well
|
||||||
|
// under a minute) — far shorter than a resumed turn can genuinely take to finish real work. These
|
||||||
|
// drive the exact real delegation path (async send -> worker asks -> primary answers -> primary's
|
||||||
|
// own wait gives up -> worker's real fleet_reply arrives afterwards) rather than calling a reply
|
||||||
|
// sink directly, since the bug is specifically about which sink the resumed turn's reply reaches.
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket() throws Exception {
|
||||||
|
String ticket = messages.sendAsync(T, "task that asks");
|
||||||
|
awaitWaiting();
|
||||||
|
injectDelivery();
|
||||||
|
|
||||||
|
CompletableFuture<MessageService.AskResult> ask =
|
||||||
|
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||||
|
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||||
|
|
||||||
|
// The primary answers, but its own bounded wait for the worker's resumed turn is short and
|
||||||
|
// expires before the worker (still genuinely working) gets back to it.
|
||||||
|
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
|
||||||
|
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||||
|
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
|
||||||
|
"the primary's own bounded wait gives up before the worker finishes resuming");
|
||||||
|
|
||||||
|
// The worker keeps working past that window and only now calls fleet_reply.
|
||||||
|
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||||
|
|
||||||
|
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||||
|
assertEquals("PR opened: https://example/pulls/42", done.reply(),
|
||||||
|
"fleet_poll{ticket} must return the worker's real reply, not stay pending forever");
|
||||||
|
assertEquals("reply", done.replySource());
|
||||||
|
assertFalse(messages.hasStrandedReply(T),
|
||||||
|
"the reply completed its own ticket directly and never touched the inbox");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket() throws Exception {
|
||||||
|
String ticket = messages.sendAsync(T, "task that asks");
|
||||||
|
awaitWaiting();
|
||||||
|
injectDelivery();
|
||||||
|
|
||||||
|
CompletableFuture<MessageService.AskResult> ask =
|
||||||
|
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||||
|
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||||
|
|
||||||
|
MessageService.Reply answerReply = messages.answer(asking.turnId(), "config.yaml", 150);
|
||||||
|
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||||
|
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome());
|
||||||
|
|
||||||
|
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||||
|
|
||||||
|
// fleet_stop tears the worker's session down right after the reply landed — this must never
|
||||||
|
// report the misleading "the worker session was released before it replied": a reply is
|
||||||
|
// exactly what happened.
|
||||||
|
assertFalse(messages.abandon(T, "the worker session was released before it replied"),
|
||||||
|
"a reply already arrived, so nothing here is a genuine failure");
|
||||||
|
|
||||||
|
MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||||
|
assertEquals("PR opened: https://example/pulls/42", view.reply());
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
|
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
|
||||||
String ticket = messages.sendAsync(T, "task that asks");
|
String ticket = messages.sendAsync(T, "task that asks");
|
||||||
|
|||||||
Reference in New Issue
Block a user