Compare commits
18 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a052975420 | |||
| 3743789e8d | |||
| 388aba7632 | |||
| ad587eafa3 | |||
| 08ce9aef11 | |||
| 4ac688b6d9 | |||
| a814d1ef00 | |||
| ea98856130 | |||
| a49e96835a | |||
| 63c19dcba7 | |||
| 2823349c8e | |||
| 6fc301d62c | |||
| 615af4ed0a | |||
| d1fd5700f5 | |||
| 2d55b0b9a5 | |||
| d89ae94a2e | |||
| e6193c4098 | |||
| b66f0677ed |
@@ -0,0 +1,208 @@
|
||||
# Wiki audit for #168
|
||||
|
||||
**Source checked:** `.wiki-snapshot/` at `68e32c6` (2026-08-31). I did not use
|
||||
`wiki/`. Code references below are from the current `fleetd` source tree. A quoted
|
||||
line is a concrete claim that needs correction, unless the table says `KEEP`.
|
||||
|
||||
| Page | Verdict | One-line reason |
|
||||
|---|---|---|
|
||||
| `Home.md` | REVISE | Good overview, but it still names the retired product. |
|
||||
| `_Sidebar.md` | REVISE | The heading still says `claude-bridge`. |
|
||||
| `1-Architecture.md` | REBUILD | Its component contract mixes current names with removed tools, routes, and planned backends. |
|
||||
| `2-Message-Server.md` | REBUILD | The claimed MCP schema, mount command, REST/SSE surface, and fallback paths are pre-build design. |
|
||||
| `3-Approaches.md` | REVISE | Useful research history, but it presents unbuilt AgentAPI as a selectable fallback. |
|
||||
| `4-Setup.md` | RETIRE | It is an intentional stub that only redirects to chapter 13. |
|
||||
| `5-Operations.md` | RETIRE | It is an intentional stub that only redirects to chapter 13. |
|
||||
| `6-Team.md` | REBUILD | It teaches role-addressed sends and a Claude-only team model that the shipped API does not have. |
|
||||
| `7-Use-Cases.md` | REBUILD | Its flagship flow depends on removed `ccs` profiles and removed send parameters. |
|
||||
| `8-Roadmap.md` | REBUILD | It is a historical plan, but it presents old implementation choices and planned work as the current stack. |
|
||||
| `9-Implementation.md` | REBUILD | Its package, class, endpoint, and outcome map has drifted from the source. |
|
||||
| `10-Cross-Host-Messaging.md` | REVISE | It labels most federation work proposed, but misses the shipped `coordinator:` lead channel. |
|
||||
| `11-Features.md` | REVISE | It is the right catalogue, but code-path names are old and it misses the second-herdr-daemon capability. |
|
||||
| `12-Claude-to-OpenCode.md` | REVISE | The porting guide is mostly current, but calls the product and spawned-member path a bridge. |
|
||||
| `13-User-Guide.md` | REVISE | It is the best operator page, but needs the product rename and the second-herdr-daemon setup. |
|
||||
|
||||
## Pages needing work
|
||||
|
||||
### `Home.md` — REVISE
|
||||
|
||||
- Quote: `# claude-bridge` (line 1) and `` `claude-bridge` keeps`` (line 11).
|
||||
The product is `fleet` / `fleetd`. The MCP server identifies itself as `fleet` in
|
||||
`fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java:313-315`.
|
||||
- Quote: `AgentAPI ... swappable fallback injector` (lines 73-76).
|
||||
There is no AgentAPI implementation under `fleetd/src/main/java`; the actual
|
||||
launchers are selected by `Profile.kind` in
|
||||
`fleetd/src/main/java/dev/ltms/fleet/config/FleetConfig.java:265-270`.
|
||||
|
||||
### `_Sidebar.md` — REVISE
|
||||
|
||||
- Quote: `### 📖 claude-bridge` (line 1).
|
||||
Rename it to `fleet`. `FleetMcp` registers the current product-facing tool set at
|
||||
`fleetd/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java:301-326`.
|
||||
|
||||
### `1-Architecture.md` — REBUILD
|
||||
|
||||
- Quote: `` `claude-bridge` lets`` (line 3). The product was renamed; the MCP
|
||||
server name is `fleet` (`FleetMcp.java:313-315`).
|
||||
- Quote: ``fleet_read`` in the tool list (line 102). No such tool is registered.
|
||||
The complete registered list is `fleet_send` through `fleet_whoami` at
|
||||
`FleetMcp.java:301-326`; `fleet_read` is absent.
|
||||
- Quote: `SSE (GET /events)` (line 143). `FleetApp.build()` registers no `/events`
|
||||
route; its routes are listed at `FleetApp.java:143-159`.
|
||||
- Quote: `Redis Streams / NATS JetStream, or an embedded queue` (line 106).
|
||||
The shipped durable inbox is AMQP, configured by `broker`, at
|
||||
`FleetConfig.java:49-50` and `FleetConfig.java:655-714`.
|
||||
- Quote: `AgentAPI (fallback)` (line 107). No AgentAPI adapter exists; shipped
|
||||
launcher kinds are `claude-code` and `opencode` (`FleetConfig.java:265-270`).
|
||||
|
||||
### `2-Message-Server.md` — REBUILD
|
||||
|
||||
- Quote: `claude mcp add --transport http bridge http://127.0.0.1:8080/mcp`
|
||||
(line 67). The daemon defaults to port `8765` in `FleetConfig.java:183-187`,
|
||||
and identifies its server as `fleet` at `FleetMcp.java:313-315`.
|
||||
- Quote: ``fleet_send(message, target?, {block, timeout_seconds, auto_spawn,
|
||||
turn_id})`` (line 80). The real parameters are `sessionId`, `content`,
|
||||
`timeoutMs`, `wait`, `turnId`, and `coordId` (`FleetMcp.java:1096-1108`).
|
||||
- Quote: ``fleet_read(target, source)`` (line 85). It is not registered; see the
|
||||
complete registration at `FleetMcp.java:301-326`.
|
||||
- Quote: `docs/MCP-Contract.md ... normative` (lines 87-88). That is not a valid
|
||||
reference: only §6 is current, as the current operator guide itself says at
|
||||
`.wiki-snapshot/13-User-Guide.md:466`.
|
||||
- Quote: `SSE (GET /events)` (line 45). No route exists in the built REST surface,
|
||||
`FleetApp.java:143-159`.
|
||||
|
||||
### `3-Approaches.md` — REVISE
|
||||
|
||||
- Quote: `AgentAPI ... remains a swappable fallback injector` (lines 78-84).
|
||||
It was never built. The shipped adapter selection is only `claude-code` or
|
||||
`opencode` (`FleetConfig.java:265-270`). Keep it as discarded research, not an
|
||||
operational fallback.
|
||||
- Quote: `claude-bridge` (line 109). Rename the product to `fleet`; the runtime
|
||||
package is `dev.ltms.fleet`, for example `FleetMcp.java:1`.
|
||||
|
||||
### `4-Setup.md` — RETIRE
|
||||
|
||||
It is a 25-line redirect and says its procedure was never written (lines 3-9).
|
||||
Chapter 13 is the maintained install procedure. Keeping a second navigation page
|
||||
adds no working documentation.
|
||||
|
||||
### `5-Operations.md` — RETIRE
|
||||
|
||||
It is a 35-line redirect and says its runbook was never written (lines 3-14).
|
||||
Chapter 13 now owns run and recovery instructions.
|
||||
|
||||
### `6-Team.md` — REBUILD
|
||||
|
||||
- Quote: `fleet_send {role: w-claude, prompt: A}` (line 98). `fleet_send` accepts
|
||||
`sessionId` and `content`, not `role` or `prompt` (`FleetMcp.java:1096-1108`).
|
||||
- Quote: `some on Claude, some on the remote local LLM` (lines 3-5) and `Every
|
||||
worker is ... Claude Code` (line 25). `opencode` is a first-class launcher kind,
|
||||
not a Claude worker (`FleetConfig.java:265-270`).
|
||||
- Quote: `fleetd's concurrency policy` (line 121). The configured capacity control
|
||||
is per-profile `maxLoad` (`FleetConfig.java:251-264`), not the role routing model
|
||||
described here.
|
||||
|
||||
### `7-Use-Cases.md` — REBUILD
|
||||
|
||||
- Quote: `ccs profile` (line 10), `ccs + herdr` (line 22), and `ccs-spawn`
|
||||
(line 45). The configuration has `profiles` and `fleet`, not `ccs`:
|
||||
`FleetConfig.java:34-58` and `FleetConfig.java:81-101`.
|
||||
- Quote: `fleet_send({"to", "kind", "body", "block"})` (lines 55-62).
|
||||
None of those are the shipped send parameters. The schema is
|
||||
`FleetMcp.java:1096-1108`.
|
||||
- Quote: `fleet_list() → { "profiles": ... }` (lines 74-80). `fleet_list` is a
|
||||
roster view; `fleet_profiles` is the configured-backend view, as registered at
|
||||
`FleetMcp.java:307-311` and described at `FleetMcp.java:1176-1182`.
|
||||
|
||||
### `8-Roadmap.md` — REBUILD
|
||||
|
||||
- Quote: `Java 21+` (line 43). The current project guidance and source use Java 25;
|
||||
the `FleetConfig` source itself uses Java 25 unnamed lambda parameters, for
|
||||
example `FleetConfig.java:102`.
|
||||
- Quote: `herdr 0.7.0 / protocol 14` (line 46). The current REST health endpoint
|
||||
reports the live protocol returned by herdr (`FleetApp.java:240-244`), while the
|
||||
current operator guide records protocol 19 at
|
||||
`.wiki-snapshot/13-User-Guide.md:76-85`.
|
||||
- Quote: `ccs <profile> claude` and `ccs env <profile>` (lines 47-48). Shipped
|
||||
configuration uses `Profile` records and launcher `kind`,
|
||||
`FleetConfig.java:313-330` and `FleetConfig.java:265-270`.
|
||||
- Quote: `Redis Streams via Lettuce` (line 50). The actual durable inbox is AMQP
|
||||
`broker`, `FleetConfig.java:655-714`.
|
||||
|
||||
### `9-Implementation.md` — REBUILD
|
||||
|
||||
- Quote: `rest.FleetdApp` and `mcp.BridgeMcp` (lines 29-30). The classes are
|
||||
`rest.FleetApp` and `mcp.FleetMcp` (`FleetApp.java:46`; `FleetMcp.java:67`).
|
||||
- Quote: `dev.ltms.fleetd` (line 67). The source package is `dev.ltms.fleet`
|
||||
(`FleetMcp.java:1`).
|
||||
- Quote: `WorkerPresence` (line 110). The current class is `MemberPresence`, as
|
||||
imported and used by `FleetMcp` at `FleetMcp.java:12` and `465-469`.
|
||||
- Quote: the outcome list ending in `STALE_TURN` (lines 128-131). The code also
|
||||
has `BACKEND_EXHAUSTED` (`FleetMcp.java:550-554`) and async `ASKING` handling
|
||||
(`FleetMcp.java:664-668`).
|
||||
- Quote: `FleetdApp` (line 207) and `FleetdConfig` (line 211). These names do not
|
||||
resolve; current classes are `FleetApp` and `FleetConfig`.
|
||||
|
||||
### `10-Cross-Host-Messaging.md` — REVISE
|
||||
|
||||
- Quote: the chapter says the cross-host fabric is proposed except for the
|
||||
single-host inbox (lines 3-8). Cross-host **lead-to-lead** delivery shipped:
|
||||
`fleet_send` accepts `coordId` (`FleetMcp.java:1094-1107`) and publishes it at
|
||||
`FleetMcp.java:616-641`; configuration has `coordinator` at
|
||||
`FleetConfig.java:74-78` and `99-101`.
|
||||
- Quote: `bridge.dlx` (line 90). This product name is stale. The shipped lead path
|
||||
uses `LeadChannel`, not the proposed exchange flow (`FleetMcp.java:95-96` and
|
||||
`616-641`). Keep the proposed federation design, but add a clear shipped/proposed
|
||||
boundary for CB-637.
|
||||
|
||||
### `11-Features.md` — REVISE
|
||||
|
||||
- Quote: `mcp/BridgeMcp` (line 22), `config/FleetdConfig` (lines 25-27), and other
|
||||
index references. These paths no longer resolve; the source classes are
|
||||
`mcp/FleetMcp` (`FleetMcp.java:67`) and `config/FleetConfig`
|
||||
(`FleetConfig.java:81`).
|
||||
- Quote: `fleet_whoami` returns only `primary` or `worker` (lines 99-100).
|
||||
It also returns `architect` (`FleetMcp.java:1235-1244`).
|
||||
- The page needs the missing separate member-herdr-daemon feature listed below.
|
||||
|
||||
### `12-Claude-to-OpenCode.md` — REVISE
|
||||
|
||||
- Quote: `same bridge mount` (line 5) and `a bridge-spawned worker` (line 94).
|
||||
Rename the product path to `fleet`. The daemon exposes the MCP server as `fleet`
|
||||
(`FleetMcp.java:313-315`), and profiles select OpenCode with `kind: opencode`
|
||||
(`FleetConfig.java:332-335`).
|
||||
- Quote: the sample mount name is `fleetd` (line 67). The server name is `fleet`;
|
||||
update the sample to avoid teaching a second product name.
|
||||
|
||||
### `13-User-Guide.md` — REVISE
|
||||
|
||||
- Quote: `The bridge is the only channel` (line 63). The invariant is correct, but
|
||||
the product term needs the `fleet` rename. The daemon's MCP server name is
|
||||
`fleet` (`FleetMcp.java:313-315`).
|
||||
- Quote: it describes one herdr socket (lines 72-85). It needs the optional
|
||||
`memberHerdrSocket` setup and two-daemon health meaning. The config key is in
|
||||
`FleetConfig.java:34-37`, and `/healthz` checks both daemons when configured at
|
||||
`FleetApp.java:210-245`.
|
||||
|
||||
## MISSING
|
||||
|
||||
`11-Features.md` has a body section for **routing members through a separate herdr daemon**
|
||||
(`## memberHerdrSocket`, line 2174), but **no row in the index table** at the top of the page
|
||||
(lines 20-95). That table is how the page is meant to be read, so a capability absent from it is
|
||||
effectively undiscoverable. Lead note: this is my own omission — I added the section on 2026-08-31
|
||||
and did not add the matching row. Fixed in the wiki at `68e32c6`'s successor.
|
||||
|
||||
The original audit stated the feature had no entry at all. That was wrong: the section exists. The
|
||||
gap is the index row. Recorded here rather than silently corrected, because the difference matters —
|
||||
"undocumented" and "documented but unindexed" are different jobs.
|
||||
|
||||
Evidence for the feature itself: `FleetConfig.java:34-37` and `FleetApp.java:103-115`, `210-245`,
|
||||
and `247-263`.
|
||||
|
||||
## Audit method and coverage
|
||||
|
||||
I checked all 15 pages. I checked concrete tool, route, config, class, file, and
|
||||
product-name claims claim-by-claim on 11 pages: Home, Sidebar, 1, 2, 4, 5, 6, 7, 9,
|
||||
11, and 13. I skimmed the remaining four long historical or research pages (3, 8, 10,
|
||||
12), then checked their concrete claims that affect the verdict. This is an audit of
|
||||
the supplied snapshot, not a wiki rewrite.
|
||||
@@ -28,6 +28,7 @@
|
||||
<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>
|
||||
|
||||
<!--
|
||||
@@ -44,6 +45,12 @@
|
||||
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
|
||||
@@ -123,6 +130,17 @@
|
||||
<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()));
|
||||
() -> config.get().memberCredentials(), null, config::get));
|
||||
}
|
||||
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().memberCredentials(), config::get));
|
||||
}
|
||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
||||
// 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 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
|
||||
@@ -12,23 +15,46 @@ import java.util.Map;
|
||||
* 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. This scans agent panes; a spawn-time {@code pid→terminal} cache is
|
||||
* the obvious optimization once wired into {@code ClaudeCodeLauncher}.
|
||||
* 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}.
|
||||
*
|
||||
* <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).
|
||||
* 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 {
|
||||
|
||||
/**
|
||||
* 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;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -37,7 +63,13 @@ 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;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -49,8 +81,9 @@ public final class PaneLocator {
|
||||
if (pid <= 0) {
|
||||
return null;
|
||||
}
|
||||
Set<Long> ancestry = ancestorsOf(pid);
|
||||
for (HerdrClient herdr : herdrs) {
|
||||
String terminal = terminalForPid(herdr, pid);
|
||||
String terminal = terminalForPid(herdr, ancestry);
|
||||
if (terminal != null) {
|
||||
return terminal;
|
||||
}
|
||||
@@ -58,28 +91,54 @@ public final class PaneLocator {
|
||||
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")) {
|
||||
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 null;
|
||||
}
|
||||
|
||||
private static boolean paneOwnsPid(HerdrClient herdr, String paneId, long pid) {
|
||||
private static boolean paneOwnsAnyOf(HerdrClient herdr, String paneId, Set<Long> ancestry) {
|
||||
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 (info.path("shell_pid").asLong(-1) == pid) {
|
||||
if (ancestry.contains(info.path("shell_pid").asLong(-1))) {
|
||||
return true;
|
||||
}
|
||||
for (JsonNode p : info.path("foreground_processes")) {
|
||||
if (p.path("pid").asLong(-1) == pid) {
|
||||
if (ancestry.contains(p.path("pid").asLong(-1))) {
|
||||
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());
|
||||
}
|
||||
@@ -74,6 +74,14 @@ public final class CompletionResolver implements TurnListener {
|
||||
*/
|
||||
public static final long MIN_TURN_NANOS = Duration.ofSeconds(2).toNanos();
|
||||
|
||||
/**
|
||||
* fleetd#164 (part 2): one stable, explicit backend-failure marker seen on a Claude Code pane
|
||||
* when the backend itself rejected the turn (e.g. {@code "API Error: 400 invalid request body"}).
|
||||
* Kept deliberately narrow — a growing list of ad-hoc error strings rots as backends change their
|
||||
* wording; broader backend-error surfacing is out of scope here (fleetd#164 point 3).
|
||||
*/
|
||||
private static final Pattern BACKEND_ERROR = Pattern.compile("(?i)\\bAPI Error\\s*:");
|
||||
|
||||
private static final String CLIPPED_PANE_TAIL_MARKER =
|
||||
"[Pane tail clipped: member did not call fleet_reply.]";
|
||||
|
||||
@@ -278,6 +286,20 @@ public final class CompletionResolver implements TurnListener {
|
||||
}
|
||||
return;
|
||||
}
|
||||
// fleetd#164 (part 2): a scrape that read cleanly and produced content still isn't a real
|
||||
// reply when that content is the backend's own rejection (e.g. an HTTP 400 before the worker
|
||||
// did any work). Classify it as a failure naming the member, rather than handing the caller a
|
||||
// scrape that reads like a completed answer.
|
||||
String backendError = firstMatchingLine(assistantBlock, BACKEND_ERROR);
|
||||
if (backendError != null) {
|
||||
// Carry the whole scrape, not just the matched line. The pattern is a heuristic: a member
|
||||
// that forgot fleet_reply while reporting *about* a backend error matches it too. Failing
|
||||
// is still right — the caller must not read a scrape as an answer — but dropping the rest
|
||||
// of the pane would destroy the report, which is the same defect fleetd#164 is about.
|
||||
fail(target, turn, "member " + target + " ended on a backend error: " + backendError
|
||||
+ "\n--- pane tail ---\n" + tail);
|
||||
return;
|
||||
}
|
||||
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
|
||||
if (rendezvous.resolveCompletion(waiter, completion)) {
|
||||
inFlight.remove(target, turn);
|
||||
|
||||
@@ -94,13 +94,26 @@ 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) {
|
||||
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) {
|
||||
this(agents, spaces, guard, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||
fleet, memberCredentials);
|
||||
fleet, memberCredentials, hostEnvNames, config);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -153,10 +166,22 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
||||
Function<String, String> env,
|
||||
long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
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) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, config);
|
||||
this.guard = guard;
|
||||
}
|
||||
|
||||
@@ -174,7 +199,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);
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames, null);
|
||||
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. */
|
||||
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)
|
||||
@@ -238,8 +240,19 @@ 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<Set<String>> hostEnvNames, Supplier<FleetConfig> config) {
|
||||
this.fleet = fleet;
|
||||
this.namePrefix = namePrefix;
|
||||
this.agents = agents;
|
||||
@@ -252,6 +265,7 @@ 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 -------------------------------------------------------------------------
|
||||
@@ -1008,6 +1022,13 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* applied before the login shell runs and a sourced file can (and did) undo it. The control is
|
||||
* the ZDOTDIR scrub ({@link #applyEnvironmentAllowListPolicy}); {@code known}/{@code allow}
|
||||
* remain as reporting only via {@link #logCredentialGap}.
|
||||
*
|
||||
* <p>CB-633 follow-up (#192): under {@code allow-list} this method does NOT call {@link
|
||||
* #logCredentialGap} itself — at this point (called from {@link #baseEnv}, before {@link
|
||||
* #applyEnvironmentAllowListPolicy} runs) we do not yet know whether the pane's shell is zsh, so
|
||||
* we cannot yet pick correct wording. That decision, and the call, are deferred entirely to
|
||||
* {@link #applyEnvironmentAllowListPolicy}, which knows by then whether the scrub will actually
|
||||
* run.
|
||||
*/
|
||||
private void applyMemberCredentialPolicy(Map<String, String> workerEnv) {
|
||||
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
|
||||
@@ -1016,14 +1037,16 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
}
|
||||
if (!creds.isAllowList()) {
|
||||
overlayBlockedCredentials(workerEnv, creds);
|
||||
logCredentialGap(creds, null);
|
||||
}
|
||||
logCredentialGap(creds);
|
||||
}
|
||||
|
||||
/** Put {@link #BLOCKED_CREDENTIAL_SENTINEL} over every blocked name in the pane-creation env map. */
|
||||
private static void overlayBlockedCredentials(Map<String, String> workerEnv,
|
||||
FleetConfig.MemberCredentials creds) {
|
||||
for (String name : creds.blockedSet()) {
|
||||
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) {
|
||||
workerEnv.put(name, BLOCKED_CREDENTIAL_SENTINEL);
|
||||
}
|
||||
}
|
||||
@@ -1067,12 +1090,15 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
// logCredentialGap's WARN (below) is the only signal for this path.
|
||||
warnNonZsh(loginShell);
|
||||
overlayBlockedCredentials(launch.env(), creds);
|
||||
logCredentialGap(creds);
|
||||
logCredentialGap(creds, null);
|
||||
return null;
|
||||
}
|
||||
// Only reached when the scrub is actually about to run — the count below describes that
|
||||
// scrub, so it must not be logged before this gate (see the non-zsh branch above).
|
||||
// scrub, so it must not be logged before this gate (see the non-zsh branch above). Same
|
||||
// reasoning gates logCredentialGap's wording: passing the derived `allowed` set (non-null)
|
||||
// here, and ONLY here, is what tells it the scrub will really blank an unkept name — #192.
|
||||
logAllowListCoverage(allowed);
|
||||
logCredentialGap(creds, allowed);
|
||||
Path dir = EnvAllowListScrub.generate(Path.of(System.getProperty("java.io.tmpdir")), allowed);
|
||||
launch.env().put("ZDOTDIR", dir.toAbsolutePath().toString());
|
||||
log.info("memberCredentials policy=allow-list: profile={} generated ZDOTDIR {} — derived "
|
||||
@@ -1087,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.
|
||||
*/
|
||||
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()));
|
||||
MemberEnvAllowList.derive(profiles.values(), creds.allowSet(), brokerUriEnvNames));
|
||||
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
|
||||
@@ -1177,18 +1209,57 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
private static final Pattern CREDENTIAL_SHAPED_NAME =
|
||||
Pattern.compile("(?i).*(TOKEN|SECRET|_KEY|APIKEY|PASSWORD|CREDENTIAL|AUTH).*");
|
||||
|
||||
/** Guards {@link #logCredentialGap} to one WARN per launcher instance, not one per spawn. */
|
||||
private final AtomicBoolean credentialGapLogged = new AtomicBoolean();
|
||||
/**
|
||||
* Guards the {@code effectiveAllowed == null} branch of {@link #logCredentialGap} — the
|
||||
* genuinely-unprotected report (deny-by-default, and the allow-list non-zsh fallback) — to one
|
||||
* WARN per launcher instance, not one per spawn.
|
||||
*
|
||||
* <p>CB-633 follow-up (#192): kept SEPARATE from {@link #allowListGapLogged} on purpose.
|
||||
* {@code memberCredentials} is a live, re-read-per-spawn supplier, so the policy can change
|
||||
* between two spawns on the same launcher. A single shared flag would let a harmless allow-list
|
||||
* INFO on spawn 1 permanently suppress the real deny-by-default WARN a later spawn deserves —
|
||||
* the report that matters most getting hidden by the report that doesn't. Two flags mean each
|
||||
* report kind fires exactly once, independent of what the other kind already logged.
|
||||
*/
|
||||
private final AtomicBoolean unprotectedGapLogged = new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* Guards the {@code effectiveAllowed != null} branch of {@link #logCredentialGap} — the
|
||||
* allow-list-scrub-covered report — to one INFO per launcher instance. See {@link
|
||||
* #unprotectedGapLogged}'s javadoc for why this is a separate flag rather than a shared one.
|
||||
*/
|
||||
private final AtomicBoolean allowListGapLogged = new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* CB-596 criterion 4: a credential-shaped host env var name on neither {@code known} nor
|
||||
* {@code allow} is not silently allowed — it is reported. {@link #hostEnvNames} enumerates the
|
||||
* daemon's own environment (see that field's javadoc for why the daemon's env is read rather
|
||||
* than the spawned pane's, which the daemon has no channel to inspect at spawn time); this logs
|
||||
* every such NAME, at WARN, at most once per launcher instance — never a value, a prefix of a
|
||||
* value, or a hash of a value, so the log itself cannot leak anything.
|
||||
* every such NAME — never a value, a prefix of a value, or a hash of a value, so the log itself
|
||||
* cannot leak anything.
|
||||
*
|
||||
* <p>CB-633 follow-up (#192): {@code effectiveAllowed} picks the wording, and it must NOT be
|
||||
* picked from {@code creds.isAllowList()} — see {@link #applyEnvironmentAllowListPolicy}'s
|
||||
* javadoc for the reasoning this mirrors. {@code null} means no scrub-derived allow-list was
|
||||
* computed for this call — true on the deny-by-default path AND on the allow-list non-zsh
|
||||
* fallback, where nothing is ever scrubbed — so the whole gap is real and gets the WARN,
|
||||
* unchanged from before this fix. Non-null means this call came from the zsh branch of {@link
|
||||
* #applyEnvironmentAllowListPolicy}, reachable ONLY after that method's own zsh gate — but
|
||||
* {@code effectiveAllowed} is a SUPERSET of {@code known ∪ allow}: {@link MemberEnvAllowList#derive}
|
||||
* also unions in every profile's {@code gitTokenEnv}/{@code gitHostEnv}/{@code tokenEnv}/
|
||||
* {@code env:} keys, and {@link #derivedAllowedNames} further unions in this very spawn's own
|
||||
* env keys — so a name can be in the gap (uncovered by {@code known}/{@code allow}) AND still be
|
||||
* kept by the derived allow-list, in which case the scrub does NOT blank it and the member DOES
|
||||
* inherit it. Lead review on #192 caught this: the first cut of this fix reported the WHOLE gap
|
||||
* as scrub-blanked without checking that, which reported a real leak as safe — the exact
|
||||
* inversion #192 exists to remove. So on this path the gap is split with {@link
|
||||
* MemberEnvAllowList#keeps}, the SAME predicate the generated scrub itself evaluates, so this
|
||||
* split cannot drift from what the scrub actually does: the names it says are kept get the WARN
|
||||
* (same severity, and same guard, as the deny-by-default case — a name genuinely reaching a
|
||||
* member unprotected is equally serious whichever path put it there), and the names it says are
|
||||
* blanked keep the INFO.
|
||||
*/
|
||||
private void logCredentialGap(FleetConfig.MemberCredentials creds) {
|
||||
private void logCredentialGap(FleetConfig.MemberCredentials creds, Set<String> effectiveAllowed) {
|
||||
Set<String> covered = new HashSet<>(creds.known());
|
||||
covered.addAll(creds.allow());
|
||||
List<String> gap = hostEnvNames.get().stream()
|
||||
@@ -1199,7 +1270,37 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
if (gap.isEmpty()) {
|
||||
return;
|
||||
}
|
||||
if (credentialGapLogged.compareAndSet(false, true)) {
|
||||
if (effectiveAllowed == null) {
|
||||
warnGapUnprotected(gap);
|
||||
return;
|
||||
}
|
||||
List<String> keptByDerivedList = gap.stream()
|
||||
.filter(name -> MemberEnvAllowList.keeps(effectiveAllowed, name))
|
||||
.toList();
|
||||
List<String> blankedByScrub = gap.stream()
|
||||
.filter(name -> !MemberEnvAllowList.keeps(effectiveAllowed, name))
|
||||
.toList();
|
||||
if (!keptByDerivedList.isEmpty() && unprotectedGapLogged.compareAndSet(false, true)) {
|
||||
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||
+ "known: nor allow: — the derived allow-list keeps them anyway (a profile's "
|
||||
+ "gitTokenEnv/gitHostEnv/tokenEnv/env: names one, or this spawn injects it), "
|
||||
+ "so every member pane inherits them UNBLOCKED — {}. Add each to "
|
||||
+ "memberCredentials.known (or .allow if a member legitimately needs it), or "
|
||||
+ "remove it from whatever profile setting derives it in.",
|
||||
keptByDerivedList.size(), keptByDerivedList);
|
||||
}
|
||||
if (!blankedByScrub.isEmpty() && allowListGapLogged.compareAndSet(false, true)) {
|
||||
log.info("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||
+ "known: nor allow: — {}. The allow-list scrub blanks them anyway (they "
|
||||
+ "are not on the derived allow-list), so no member pane keeps them; add "
|
||||
+ "each to memberCredentials.known or .allow to make that explicit.",
|
||||
blankedByScrub.size(), blankedByScrub);
|
||||
}
|
||||
}
|
||||
|
||||
/** The deny-by-default (and allow-list non-zsh fallback) WARN — unchanged byte-for-byte by #192. */
|
||||
private void warnGapUnprotected(List<String> gap) {
|
||||
if (unprotectedGapLogged.compareAndSet(false, true)) {
|
||||
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||
+ "known: nor allow: — every member pane inherits them UNBLOCKED — {}. "
|
||||
+ "Add each to memberCredentials.known (blocked by default) or .allow "
|
||||
|
||||
@@ -40,8 +40,9 @@ 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} 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
|
||||
* 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
|
||||
* (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
|
||||
@@ -106,6 +107,16 @@ 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) {
|
||||
@@ -124,9 +135,37 @@ 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
|
||||
@@ -146,4 +185,16 @@ 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,12 +116,23 @@ 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) {
|
||||
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) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||
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,
|
||||
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);
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, null, config);
|
||||
this.configRoot = configRoot;
|
||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||
}
|
||||
|
||||
@@ -1,58 +1,87 @@
|
||||
package dev.ltms.fleet.member;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.sqlite.SQLiteConfig;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
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
|
||||
* 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 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><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>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
|
||||
* 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 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.
|
||||
* <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.
|
||||
*/
|
||||
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 ObjectMapper json;
|
||||
private final Path databasePath;
|
||||
private final AtomicBoolean warnedMissingDatabase = new AtomicBoolean(false);
|
||||
|
||||
OpenCodeSessionDiscovery(Path 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
|
||||
* {@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
|
||||
* 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 session the pane most likely corresponds to.
|
||||
*
|
||||
* <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.
|
||||
* <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.
|
||||
*
|
||||
* @param directory the worker's cwd, as resolved for this spawn
|
||||
* @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()) {
|
||||
return null;
|
||||
}
|
||||
Path sessionRoot = storageRoot.resolve("session");
|
||||
if (!Files.isDirectory(sessionRoot)) {
|
||||
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);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
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
|
||||
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");
|
||||
}
|
||||
}
|
||||
} catch (IOException ignored) {
|
||||
// storage root vanished or became unreadable — "no session known yet"
|
||||
} 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());
|
||||
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;
|
||||
}
|
||||
log.debug("no opencode session row for directory (root={}, directory={})",
|
||||
storageRoot, directory);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -168,14 +168,23 @@ public final class MessageService {
|
||||
private final String ticket;
|
||||
private final String target;
|
||||
private final CompletableFuture<Reply> future = new CompletableFuture<>();
|
||||
private final long createdNanos;
|
||||
/**
|
||||
* When {@link #future} resolved, or {@code null} while it is still pending — the clock
|
||||
* {@link #pruneTerminalTickets} measures the TTL from (#197). Deliberately a boxed
|
||||
* {@code Long} rather than a {@code long} with a sentinel: {@link System#nanoTime} may
|
||||
* legitimately return any value, zero and negatives included, so no numeric sentinel can mean
|
||||
* "not stamped yet". Stamped by a {@code whenComplete} hook registered in the constructor, so
|
||||
* every completion path stamps it — a reply, the completion fallback, a timeout, a failure,
|
||||
* or an abandon on teardown — without each of those having to remember to.
|
||||
*/
|
||||
private volatile Long completedNanos;
|
||||
private volatile Reply question;
|
||||
private volatile String turnId;
|
||||
|
||||
private Task(String ticket, String target, long createdNanos) {
|
||||
private Task(String ticket, String target, LongSupplier nowNanos) {
|
||||
this.ticket = ticket;
|
||||
this.target = target;
|
||||
this.createdNanos = createdNanos;
|
||||
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -357,21 +366,40 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* Route a worker's explicit {@code fleet_reply}: resolve an open send, or queue it in the
|
||||
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
|
||||
* result is <em>not</em> a failure — the reply is held for later drain.
|
||||
* Route a worker's explicit {@code fleet_reply}: resolve an open send, complete an async ticket
|
||||
* still parked waiting on this exact turn's answer, or — only once neither applies — queue it in
|
||||
* 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} /
|
||||
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
|
||||
* 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) {
|
||||
if (rendezvous.resolve(session, content)) {
|
||||
count(FleetMetrics.REPLIES, "path", "rendezvous");
|
||||
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);
|
||||
// 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.
|
||||
@@ -385,6 +413,24 @@ public final class MessageService {
|
||||
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. */
|
||||
private void count(String name, String... labels) {
|
||||
if (metrics != null) {
|
||||
@@ -426,19 +472,43 @@ public final class MessageService {
|
||||
* <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}.
|
||||
*
|
||||
* @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) {
|
||||
boolean hadStrandedReply = hasStrandedReply(target);
|
||||
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
|
||||
strandedReplies.remove(target);
|
||||
queuedDeliveries.remove(target);
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
||||
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
|
||||
boolean asyncFailed = false;
|
||||
Reply recovered = null; // lazily drained at most once, only if a task actually needs it
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null
|
||||
&& task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))) {
|
||||
asyncFailed = true;
|
||||
if (!target.equals(task.target) || task.question != null || task.future.isDone()) {
|
||||
continue;
|
||||
}
|
||||
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) {
|
||||
@@ -447,6 +517,21 @@ public final class MessageService {
|
||||
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
|
||||
* so that a subsequent drain or peek no longer returns it.
|
||||
@@ -721,7 +806,7 @@ public final class MessageService {
|
||||
*/
|
||||
public String sendAsync(String target, String content, Runnable onAccepted) {
|
||||
String ticket = "task-" + ticketSeq.incrementAndGet();
|
||||
Task task = new Task(ticket, target, nowNanos.getAsLong());
|
||||
Task task = new Task(ticket, target, nowNanos);
|
||||
tasks.put(ticket, task);
|
||||
if (pushLoop != null) {
|
||||
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
|
||||
@@ -821,12 +906,25 @@ public final class MessageService {
|
||||
* (or one the reminder cap already gave up on) is pruned here but never collected there, so it
|
||||
* lingers in {@code pendingTickets} forever and rides along on every later nudge to the same lead
|
||||
* — naming a ticket {@code fleet_poll} can no longer find (CB-588 follow-up).
|
||||
*
|
||||
* <p>The TTL runs from **completion**, not from creation (#197). It used to compare against
|
||||
* {@code createdNanos}, which made the real collection window {@code TTL minus however long the
|
||||
* task ran}: a delegation that took longer than the TTL had its reply destroyed on the first
|
||||
* sweep after it landed, every time. That is the normal case here — real work runs well past ten
|
||||
* minutes — and the reply lives only in {@code future}, so pruning it discards the worker's whole
|
||||
* report with nothing to fall back on. Measuring from completion gives every ticket the same full
|
||||
* window whatever its runtime, and still bounds {@code tasks}.
|
||||
*
|
||||
* <p>A task whose future is done but whose {@code completedNanos} is not stamped yet is left
|
||||
* alone. That window is the few instructions between {@code complete()} and the constructor's
|
||||
* {@code whenComplete} hook running; the next sweep collects it.
|
||||
*/
|
||||
private void pruneTerminalTickets() {
|
||||
long cutoff = nowNanos.getAsLong() - TICKET_TTL_NANOS;
|
||||
tasks.entrySet().removeIf(e -> {
|
||||
Task t = e.getValue();
|
||||
boolean expired = t.future.isDone() && t.createdNanos < cutoff;
|
||||
Long completed = t.completedNanos;
|
||||
boolean expired = t.future.isDone() && completed != null && completed < cutoff;
|
||||
if (expired && pushLoop != null) {
|
||||
pushLoop.ticketCollected(e.getKey());
|
||||
}
|
||||
|
||||
@@ -110,6 +110,7 @@ public final class GitWorktrees implements Worktrees {
|
||||
|
||||
@Override
|
||||
public String add(String repoRoot, String branch, String baseRef) {
|
||||
reportRemoteUrlsWithUserInfo(repoRoot);
|
||||
String base = (baseRef == null || baseRef.isBlank()) ? "HEAD" : baseRef;
|
||||
String nonce = nonce();
|
||||
Path root = resolveRoot(repoRoot);
|
||||
@@ -139,7 +140,7 @@ public final class GitWorktrees implements Worktrees {
|
||||
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
|
||||
return;
|
||||
}
|
||||
String origin = exec("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
|
||||
String origin = execRedacted("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
|
||||
URI uri;
|
||||
try {
|
||||
uri = new URI(origin);
|
||||
@@ -155,6 +156,9 @@ public final class GitWorktrees implements Worktrees {
|
||||
throw new WorktreeException("origin URL has invalid HTTPS user info; cannot provision safely");
|
||||
}
|
||||
String cleanOrigin = origin.substring(0, schemeEnd) + origin.substring(userInfoEnd + 1);
|
||||
// Plain exec, not execRedacted, is correct here: this call WRITES cleanOrigin (already
|
||||
// stripped of user-info above) rather than reading a URL back from stdout, so there is
|
||||
// nothing secret left in either its argv or its stdout to redact.
|
||||
exec("git", "-C", repoRoot, "remote", "set-url", "origin", cleanOrigin);
|
||||
log.info("removed HTTPS user info from forge origin before provisioning worktree");
|
||||
}
|
||||
@@ -164,7 +168,7 @@ public final class GitWorktrees implements Worktrees {
|
||||
if (exitCode("git", "-C", worktreePath, "remote", "get-url", "--all", "origin") != 0) {
|
||||
return;
|
||||
}
|
||||
String origins = exec("git", "-C", worktreePath, "remote", "get-url", "--all", "origin");
|
||||
String origins = execRedacted("git", "-C", worktreePath, "remote", "get-url", "--all", "origin");
|
||||
for (String origin : origins.split("\\R")) {
|
||||
try {
|
||||
URI uri = new URI(origin);
|
||||
@@ -177,6 +181,99 @@ public final class GitWorktrees implements Worktrees {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Report — never refuse — every remote whose fetch or push URL carries user-info outside SSH.
|
||||
* A linked worktree shares its parent repository's git config, so a credential on ANY remote
|
||||
* (not only {@code origin}) or in a {@code pushurl} is just as readable by a member as one on
|
||||
* {@code origin}'s HTTPS fetch URL — the one case {@link #removeUserInfoFromHttpsOrigin} and
|
||||
* {@link #requireCredentialFreeHttpsOrigin} already strip and refuse. This check is additive: it
|
||||
* only logs a warning, it never mutates config and never refuses the provision.
|
||||
*
|
||||
* <p>A reporting-only check must never be able to abort a provision — PR #173 shipped one that
|
||||
* ran unguarded at the top of {@link #add}, and every git call inside it can throw ({@link #exec}
|
||||
* turns a non-zero exit or its 30-second timeout into a {@link WorktreeException}). Every git call
|
||||
* here is therefore wrapped, and on failure only the exception's <em>class</em> is logged, never
|
||||
* its message: the enumerating {@code git remote} call is not redacted, and its stderr is read
|
||||
* from the very config that may hold the URL this check exists to find.
|
||||
*/
|
||||
private void reportRemoteUrlsWithUserInfo(String repoRoot) {
|
||||
List<String> remotes;
|
||||
try {
|
||||
remotes = exec("git", "-C", repoRoot, "remote").lines()
|
||||
.map(String::trim)
|
||||
.filter(r -> !r.isBlank())
|
||||
.toList();
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("could not enumerate remotes to check for credentialed URLs in {}: {}",
|
||||
Path.of(repoRoot).toAbsolutePath().normalize(), e.getClass().getName());
|
||||
return;
|
||||
}
|
||||
for (String remote : remotes) {
|
||||
reportOneRemoteUrlsWithUserInfo(repoRoot, remote);
|
||||
}
|
||||
}
|
||||
|
||||
private void reportOneRemoteUrlsWithUserInfo(String repoRoot, String remote) {
|
||||
boolean leaks = remoteUrlsLeakUserInfo(repoRoot, remote, false)
|
||||
|| remoteUrlsLeakUserInfo(repoRoot, remote, true);
|
||||
if (leaks) {
|
||||
log.warn("member worktree shares a remote URL containing user-info: remote={} repository={}; "
|
||||
+ "remove credentials from the repository's git config",
|
||||
remote, Path.of(repoRoot).toAbsolutePath().normalize());
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* True when any resolved fetch (or, if {@code push}, push) URL for {@code remote} carries
|
||||
* non-empty user-info outside SSH. Never throws — a git failure here is caught, logged (its
|
||||
* class only, per the javadoc above), and treated as "nothing found", so it cannot abort or
|
||||
* otherwise affect provisioning. Uses {@link #execRedacted} because the command's stdout is
|
||||
* itself the URL this check exists to find.
|
||||
*/
|
||||
private boolean remoteUrlsLeakUserInfo(String repoRoot, String remote, boolean push) {
|
||||
try {
|
||||
String out = push
|
||||
? execRedacted("git", "-C", repoRoot, "remote", "get-url", "--push", "--all", remote)
|
||||
: execRedacted("git", "-C", repoRoot, "remote", "get-url", "--all", remote);
|
||||
return out.lines().anyMatch(url -> !url.isBlank() && urlLeaksUserInfo(url.trim()));
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("could not read the {} URL for remote {} to check for credentials: {}",
|
||||
push ? "push" : "fetch", remote, e.getClass().getName());
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* True when {@code rawUrl} parses as an absolute URI with a non-SSH-family scheme and non-empty
|
||||
* user-info. An unparsable or scheme-less URL — including the ssh scp-like shorthand
|
||||
* ({@code user@host:path}) — is not this check's concern and is treated as "no finding", the
|
||||
* same way {@link #configureHttpsUrlRewriteForSshOrigin} leaves that shorthand untouched.
|
||||
*/
|
||||
private static boolean urlLeaksUserInfo(String rawUrl) {
|
||||
URI uri;
|
||||
try {
|
||||
uri = new URI(rawUrl);
|
||||
} catch (URISyntaxException e) {
|
||||
return false;
|
||||
}
|
||||
String scheme = uri.getScheme();
|
||||
if (scheme == null || isSshLikeScheme(scheme)) {
|
||||
return false;
|
||||
}
|
||||
String userInfo = uri.getUserInfo();
|
||||
return userInfo != null && !userInfo.isEmpty();
|
||||
}
|
||||
|
||||
/**
|
||||
* SSH-family schemes deliberately excluded from {@link #urlLeaksUserInfo}: there, the user part
|
||||
* selects an account and authentication itself happens over SSH, so it is not a credential the
|
||||
* way HTTPS/HTTP user-info is.
|
||||
*/
|
||||
private static boolean isSshLikeScheme(String scheme) {
|
||||
return "ssh".equalsIgnoreCase(scheme) || "git+ssh".equalsIgnoreCase(scheme)
|
||||
|| "ssh+git".equalsIgnoreCase(scheme);
|
||||
}
|
||||
|
||||
/** Configure a per-worktree helper that supplies a token from the member environment at call time. */
|
||||
private void configureEnvironmentCredentialHelper(String repoRoot, String worktreePath) {
|
||||
exec("git", "-C", repoRoot, "config", "extensions.worktreeConfig", "true");
|
||||
@@ -234,7 +331,7 @@ public final class GitWorktrees implements Worktrees {
|
||||
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
|
||||
return;
|
||||
}
|
||||
String origin = exec("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
|
||||
String origin = execRedacted("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
|
||||
URI uri;
|
||||
try {
|
||||
uri = new URI(origin);
|
||||
@@ -603,6 +700,31 @@ public final class GitWorktrees implements Worktrees {
|
||||
|
||||
/** Same as {@link #exec(String...)}, with extra environment variables set on the child process. */
|
||||
private String exec(Map<String, String> extraEnv, String... command) {
|
||||
return exec(extraEnv, false, command);
|
||||
}
|
||||
|
||||
/**
|
||||
* Same as {@link #exec(String...)}, for a command whose stdout may itself carry a credential
|
||||
* (e.g. {@code git remote get-url}, whose output is a URL). Stdout is still returned normally on
|
||||
* success — callers still get the URL to inspect — but it is suppressed from BOTH the timeout
|
||||
* message and the non-zero-exit message, so a failing call here can never copy it into a
|
||||
* {@link WorktreeException}, and from there into a caller's log.
|
||||
*/
|
||||
private String execRedacted(String... command) {
|
||||
return exec(Map.of(), true, command);
|
||||
}
|
||||
|
||||
/**
|
||||
* Shared implementation for {@link #exec(Map, String...)} and {@link #execRedacted(String...)}.
|
||||
* {@code redactOutput} suppresses captured stdout from both failure messages below.
|
||||
*
|
||||
* <p>Package-private, not {@code private}: also a test seam, the same way the
|
||||
* {@code afterWorktreeAdded} constructor parameter is. It lets a test drive the redaction
|
||||
* guarantee directly — a synthetic failing command whose stdout carries a test marker passed
|
||||
* through {@code extraEnv} rather than argv — without depending on finding a real git failure
|
||||
* mode that happens to echo a URL onto stdout before exiting non-zero.
|
||||
*/
|
||||
String exec(Map<String, String> extraEnv, boolean redactOutput, String... command) {
|
||||
String out;
|
||||
int code;
|
||||
Process p;
|
||||
@@ -624,7 +746,8 @@ public final class GitWorktrees implements Worktrees {
|
||||
try {
|
||||
if (!p.waitFor(30, TimeUnit.SECONDS)) {
|
||||
p.destroyForcibly();
|
||||
throw new WorktreeException("command timed out: " + String.join(" ", command) + "\n" + out);
|
||||
throw new WorktreeException("command timed out: " + String.join(" ", command)
|
||||
+ (redactOutput || out.isBlank() ? "" : "\n" + out));
|
||||
}
|
||||
code = p.exitValue();
|
||||
} catch (InterruptedException e) {
|
||||
@@ -634,7 +757,7 @@ public final class GitWorktrees implements Worktrees {
|
||||
}
|
||||
if (code != 0) {
|
||||
throw new WorktreeException("exit " + code + " for: " + String.join(" ", command)
|
||||
+ (out.isBlank() ? "" : "\n" + out));
|
||||
+ (redactOutput || out.isBlank() ? "" : "\n" + out));
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
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). */
|
||||
@@ -63,4 +67,123 @@ 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() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -566,6 +566,73 @@ class CompletionResolverTest {
|
||||
assertEquals("The usage limit has been reached.", waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
// --- fleetd#164 (part 2 addendum): narrow BACKEND_ERROR pattern classification ---------
|
||||
|
||||
@Test
|
||||
void classifiesABackendErrorLineAsAFailureInsteadOfACompletedReply() {
|
||||
String block = "⏺ API Error: 400 invalid request body\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertTrue(waiter.isDone(), "a backend-error scrape still resolves the blocked send");
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"a backend rejection is a failure, not a completed reply");
|
||||
}
|
||||
|
||||
@Test
|
||||
void theBackendErrorReasonNamesTheMemberAndCarriesTheMatchedLine() {
|
||||
String block = "⏺ API Error: 400 invalid request body\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
String reason = waiter.getNow(null).text();
|
||||
assertTrue(reason.contains("term_a"), "the failure names the member: " + reason);
|
||||
assertTrue(reason.contains("API Error: 400 invalid request body"),
|
||||
"the failure carries the matched backend-error line: " + reason);
|
||||
}
|
||||
|
||||
@Test
|
||||
void aCaseInsensitiveApiErrorLineIsStillClassifiedAsABackendError() {
|
||||
String block = "⏺ api error: rate limited\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(), "the pattern is case-insensitive");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBackendErrorFailureStillCarriesTheRestOfTheScrape() {
|
||||
// The pattern is a heuristic: a member that forgot fleet_reply while *reporting on* a backend
|
||||
// error matches it too. Failing is still correct, but the report itself must survive — losing
|
||||
// it would be the same information-destroying defect fleetd#164 exists to fix.
|
||||
String block = "\u23fa I looked into the gateway problem.\n"
|
||||
+ "The log line was: API Error: 400 invalid request body\n"
|
||||
+ "The cause is a missing content-type header.\n\u276f ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
String reason = waiter.getNow(null).text();
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
|
||||
assertTrue(reason.contains("The cause is a missing content-type header."),
|
||||
"the failure carries the rest of the pane, not only the matched line: " + reason);
|
||||
}
|
||||
|
||||
@Test
|
||||
void coverageIsOffWhenNoProfileHasAPatternConfigured() {
|
||||
assertEquals("off (no profile has an exhaustedPattern configured; profiles: [terra])",
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package dev.ltms.fleet.member;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
@@ -1055,6 +1056,288 @@ class ClaudeCodeLauncherTest {
|
||||
"a name already on allow: is covered, not a gap");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up (#192): the deny-by-default WARN wording is a promise an operator relies on —
|
||||
* pinned byte-for-byte so a future edit cannot drift it (e.g. while picking wording for the
|
||||
* allow-list path) without a test noticing.
|
||||
*/
|
||||
@Test
|
||||
void denyByDefaultKeepsTheExactCredentialGapWarn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
_ -> null, 0, System::currentTimeMillis, () -> {}, null, () -> TEST_MEMBER_CREDENTIALS,
|
||||
() -> Set.of("A_BRAND_NEW_SECRET_TOKEN"));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
svc.spawn();
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
("memberCredentials gap: 1 credential-shaped env var name(s) are on neither "
|
||||
+ "known: nor allow: — every member pane inherits them UNBLOCKED — "
|
||||
+ "[A_BRAND_NEW_SECRET_TOKEN]. Add each to memberCredentials.known "
|
||||
+ "(blocked by default) or .allow (if a member legitimately needs it).")
|
||||
.equals(e.getFormattedMessage())),
|
||||
"the deny-by-default WARN text must not drift — got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up (#192), defect 1: under {@code policy: allow-list} on a zsh login shell the
|
||||
* generated ZDOTDIR scrub genuinely blanks an unkept credential-shaped name, so the report must
|
||||
* not say the pane inherits it UNBLOCKED — that claim is exactly what PR #174 got wrong. This
|
||||
* goes through the real spawn path (not {@code HerdrPeerLauncherAllowListWiringTest}'s fixture,
|
||||
* which overrides {@code buildLaunch} and bypasses none of the logic under test here — the
|
||||
* shell-dependent branch lives in {@code applyEnvironmentAllowListPolicy}, which every spawn
|
||||
* still passes through).
|
||||
*/
|
||||
@Test
|
||||
void allowListPolicyOnZshReportsTheGapWithoutClaimingItIsUnblocked() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
FleetConfig.MemberCredentials creds = new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST,
|
||||
List.of("AI_GATEWAY_TOKEN"), List.of());
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
name -> "SHELL".equals(name) ? "/bin/zsh" : null,
|
||||
0, System::currentTimeMillis, () -> {}, null, () -> creds,
|
||||
() -> Set.of("AI_GATEWAY_TOKEN", "A_BRAND_NEW_SECRET_TOKEN"));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
svc.spawn();
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getFormattedMessage().contains("memberCredentials gap")
|
||||
&& e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN")),
|
||||
"the unkept name must still be reported — got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
assertFalse(appender.list.stream().anyMatch(e ->
|
||||
e.getFormattedMessage().contains("memberCredentials gap")
|
||||
&& e.getFormattedMessage().contains("UNBLOCKED")),
|
||||
"on zsh the scrub genuinely blanks the name, so the report must not claim it is "
|
||||
+ "inherited unblocked — got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up (#192): the mirror of the zsh test above. On a non-zsh login shell {@code
|
||||
* ZDOTDIR} is ignored, so no scrub ever runs — the report must keep the WARN wording (a name here
|
||||
* really is inherited unblocked) rather than claiming a scrub protects it. This is the trap PR
|
||||
* #174 fell into the other direction: keying the wording on the shell, not on {@code
|
||||
* creds.isAllowList()}, is what keeps this branch correct.
|
||||
*/
|
||||
@Test
|
||||
void allowListPolicyOnNonZshKeepsTheWarnWording() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
FleetConfig.MemberCredentials creds = new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST,
|
||||
List.of("AI_GATEWAY_TOKEN"), List.of());
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
name -> "SHELL".equals(name) ? "/bin/bash" : null,
|
||||
0, System::currentTimeMillis, () -> {}, null, () -> creds,
|
||||
() -> Set.of("AI_GATEWAY_TOKEN", "A_BRAND_NEW_SECRET_TOKEN"));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
svc.spawn();
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getFormattedMessage().contains("memberCredentials gap")
|
||||
&& e.getFormattedMessage().contains("UNBLOCKED")
|
||||
&& e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN")),
|
||||
"no scrub runs on a non-zsh shell, so the WARN wording must be kept — got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
assertFalse(appender.list.stream().anyMatch(e ->
|
||||
e.getFormattedMessage().contains("memberCredentials gap")
|
||||
&& e.getFormattedMessage().toLowerCase(java.util.Locale.ROOT).contains("scrub")),
|
||||
"nothing is scrubbed on this path, so the report must not claim otherwise — got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up (#192), defect 2: {@code memberCredentials} is a live, re-read-per-spawn
|
||||
* supplier, so the policy can change between two spawns on the same launcher. Before this fix a
|
||||
* single {@code AtomicBoolean} guarded both report kinds, so the harmless allow-list INFO on the
|
||||
* first spawn would permanently suppress the real deny-by-default WARN a later spawn deserves.
|
||||
* This goes through {@link ClaudeCodeLauncher#buildLaunch}'s real {@code baseEnv()} path — the
|
||||
* WARN this test pins fires from {@code applyMemberCredentialPolicy}, which {@code
|
||||
* HerdrPeerLauncherAllowListWiringTest}'s fixture never reaches at all (see its class javadoc).
|
||||
*/
|
||||
@Test
|
||||
void secondSpawnStillWarnsAfterPolicyChangesFromAllowListToDenyByDefault() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
AtomicReference<FleetConfig.MemberCredentials> creds = new AtomicReference<>(
|
||||
new FleetConfig.MemberCredentials(FleetConfig.MemberCredentials.POLICY_ALLOW_LIST,
|
||||
List.of("AI_GATEWAY_TOKEN"), List.of()));
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
name -> "SHELL".equals(name) ? "/bin/zsh" : null,
|
||||
0, System::currentTimeMillis, () -> {}, null, creds::get,
|
||||
() -> Set.of("AI_GATEWAY_TOKEN", "A_BRAND_NEW_SECRET_TOKEN"));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
try {
|
||||
logger.addAppender(appender);
|
||||
svc.spawn(); // allow-list + zsh: harmless INFO, sets the allow-list guard only
|
||||
|
||||
creds.set(new FleetConfig.MemberCredentials(FleetConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT,
|
||||
List.of("AI_GATEWAY_TOKEN"), List.of()));
|
||||
svc.spawn(); // policy reloaded to deny-by-default: this WARN must NOT be suppressed
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getFormattedMessage().contains("memberCredentials gap")
|
||||
&& e.getFormattedMessage().contains("UNBLOCKED")),
|
||||
"the second spawn's deny-by-default WARN must still fire even though the first "
|
||||
+ "spawn's allow-list INFO already logged the same underlying gap — got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up (#192), lead-review fix: {@code effectiveAllowed} is a SUPERSET of
|
||||
* {@code known ∪ allow} — {@code MemberEnvAllowList.derive} also unions in every profile's
|
||||
* {@code tokenEnv} (among other fields), so a credential-shaped name can be uncovered by
|
||||
* {@code known:}/{@code allow:} and STILL survive the scrub because a profile's own
|
||||
* {@code tokenEnv} names it. Here {@code tokenEnv} is deliberately set to a credential-shaped
|
||||
* name the operator forgot to list — the misconfiguration this report exists to catch. The scrub
|
||||
* genuinely keeps it, so the report must WARN, not claim (as the pre-lead-review cut of this fix
|
||||
* did) that "no member pane keeps them".
|
||||
*/
|
||||
@Test
|
||||
void allowListWarnsWhenTheDerivedAllowListKeepsAnUncoveredName() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "SOME_LEAKY_TOKEN",
|
||||
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
FleetConfig.MemberCredentials creds = new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of());
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
name -> "SHELL".equals(name) ? "/bin/zsh" : null,
|
||||
0, System::currentTimeMillis, () -> {}, null, () -> creds,
|
||||
() -> Set.of("SOME_LEAKY_TOKEN"));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
svc.spawn();
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.WARN
|
||||
&& e.getFormattedMessage().contains("memberCredentials gap")
|
||||
&& e.getFormattedMessage().contains("SOME_LEAKY_TOKEN")
|
||||
&& e.getFormattedMessage().contains("UNBLOCKED")),
|
||||
"a name kept by the derived allow-list (via this profile's tokenEnv) must still WARN "
|
||||
+ "— got: " + appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
assertFalse(appender.list.stream().anyMatch(e ->
|
||||
e.getFormattedMessage().contains("memberCredentials gap")
|
||||
&& e.getFormattedMessage().contains("SOME_LEAKY_TOKEN")
|
||||
&& e.getFormattedMessage().contains("blanks them anyway")),
|
||||
"the scrub does NOT blank this name, so the INFO wording must not claim it does — got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up (#192), lead-review fix: a mixed gap — one name the derived allow-list keeps
|
||||
* (this profile's {@code tokenEnv}), one it does not — must split cleanly: the WARN names only
|
||||
* the kept one, the INFO names only the blanked one. Proves the split uses {@code
|
||||
* MemberEnvAllowList.keeps} per-name rather than an all-or-nothing decision for the whole gap.
|
||||
*/
|
||||
@Test
|
||||
void allowListSplitsAMixedGapBetweenTheWarnAndTheInfo() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "SOME_LEAKY_TOKEN",
|
||||
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
FleetConfig.MemberCredentials creds = new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of());
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
name -> "SHELL".equals(name) ? "/bin/zsh" : null,
|
||||
0, System::currentTimeMillis, () -> {}, null, () -> creds,
|
||||
() -> Set.of("SOME_LEAKY_TOKEN", "A_BRAND_NEW_SECRET_TOKEN"));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
svc.spawn();
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
boolean warnNamesOnlyKept = appender.list.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.WARN
|
||||
&& e.getFormattedMessage().contains("memberCredentials gap")
|
||||
&& e.getFormattedMessage().contains("SOME_LEAKY_TOKEN")
|
||||
&& !e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN"));
|
||||
boolean infoNamesOnlyBlanked = appender.list.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.INFO
|
||||
&& e.getFormattedMessage().contains("memberCredentials gap")
|
||||
&& e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN")
|
||||
&& !e.getFormattedMessage().contains("SOME_LEAKY_TOKEN"));
|
||||
|
||||
assertTrue(warnNamesOnlyKept, "the WARN must name the derived-list-kept variable and only it "
|
||||
+ "— got: " + appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
assertTrue(infoNamesOnlyBlanked, "the INFO must name the scrub-blanked variable and only it "
|
||||
+ "— got: " + appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* The half of CB-592 that can actually survive the pane's login shell. BRIDGED_MEMBER is a name
|
||||
* secrets.sh never exports, so nothing overwrites it — measured: GITEA_TOKEN is injected the
|
||||
|
||||
+45
-2
@@ -161,6 +161,34 @@ 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
|
||||
@@ -260,15 +288,24 @@ 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) {
|
||||
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),
|
||||
Map.of("test", profile()), "test",
|
||||
name -> "SHELL".equals(name) ? shell : null,
|
||||
0, () -> 0L, () -> { }, null, creds, hostEnvNames);
|
||||
0, () -> 0L, () -> { }, null, creds, hostEnvNames, config);
|
||||
}
|
||||
|
||||
@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"));
|
||||
}
|
||||
|
||||
@@ -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. */
|
||||
@Test
|
||||
void theGeneratedDirectoryIsRemovedWhenThePaneIsStopped() {
|
||||
|
||||
@@ -121,6 +121,34 @@ 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,8 +334,7 @@ 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, "p1", "ses_a.json",
|
||||
"ses_resolved", "/work/dir", 1000L);
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "ses_resolved", "/work/dir", 1000L);
|
||||
assertEquals("ses_resolved", handle.agentSessionId(),
|
||||
"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.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.*;
|
||||
|
||||
/**
|
||||
* {@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}.
|
||||
* {@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.
|
||||
*/
|
||||
class OpenCodeSessionDiscoveryTest {
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* 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.
|
||||
*/
|
||||
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));
|
||||
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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
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);
|
||||
void findsTheRowWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||
writeRecord(root, "ses_bbb", "/w/b", 2000L);
|
||||
|
||||
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"));
|
||||
}
|
||||
|
||||
@Test
|
||||
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"),
|
||||
"no record for this cwd yet → null, not a wrong session");
|
||||
"no row for this cwd yet → null, not a wrong session");
|
||||
}
|
||||
|
||||
@Test
|
||||
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);
|
||||
void prefersTheMostRecentlyUpdatedRowWhenSeveralMatch(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_old", "/w/a", 1000L);
|
||||
writeRecord(root, "ses_new", "/w/a", 5000L);
|
||||
|
||||
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
|
||||
void aMissingOrEmptyStorageRootYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
|
||||
// Missing: no session dir at all under the root.
|
||||
void aMissingDatabaseYieldsNullWithoutThrowing(@TempDir Path root) {
|
||||
// No opencode.db at all under the root.
|
||||
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
|
||||
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"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -76,15 +97,41 @@ 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 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);
|
||||
void theDatabaseIsOpenedReadOnly(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "ses_aaa", "/w/a", 1000L);
|
||||
|
||||
assertEquals("ses_good", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||
"an unreadable record is skipped; a later valid one still matches");
|
||||
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");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -86,6 +86,28 @@ class MessageServiceTest {
|
||||
assertTrue(reply.completed(), "a scraped completion still counts as completed");
|
||||
}
|
||||
|
||||
@Test
|
||||
void backendErrorScrapeThroughMessageServiceFailsInsteadOfBecomingReplyText() throws Exception {
|
||||
// fleetd#164 (part 2 addendum): a scrape that reads cleanly but is only the backend's own
|
||||
// rejection (e.g. an HTTP 400) must reach the caller as WORKER_FAILED, not as a completed
|
||||
// reply whose text happens to be the error line.
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
awaitWaiting();
|
||||
|
||||
herdr.readText("$ prompt");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
herdr.readText("⏺ API Error: 400 invalid request body");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
MessageService.Reply reply = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome(),
|
||||
"a backend rejection must use the caller's failure outcome, not a completed reply");
|
||||
assertFalse(reply.completed(), "plain backend errors are never fallback reply content");
|
||||
assertTrue(reply.text().contains("API Error: 400 invalid request body"),
|
||||
"the visible backend error is carried as the failure reason: " + reply.text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void explicitFleetReplyResolvesAsReplied() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
@@ -688,6 +710,68 @@ class MessageServiceTest {
|
||||
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
|
||||
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
@@ -1057,6 +1141,74 @@ class MessageServiceTest {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* #197: the ticket TTL must run from COMPLETION, not from creation.
|
||||
*
|
||||
* <p>It used to compare the cutoff against {@code createdNanos}, so the real window to collect a
|
||||
* reply was {@code TTL minus however long the task ran}. A delegation that ran longer than the
|
||||
* TTL was already past the cutoff the moment it finished, so the very next prune destroyed its
|
||||
* reply — and the reply lives only in the task's future, so nothing could get it back. That is
|
||||
* the normal case for real work here, not an edge case: three workers in one session ran well
|
||||
* past ten minutes and two of their complete reports were lost this way.
|
||||
*
|
||||
* <p>The task below runs for longer than the whole TTL before it replies, which is exactly the
|
||||
* shape that used to lose everything. Remove the fix and this fails: {@code poll} returns
|
||||
* {@code null} because the ticket was pruned on arrival.
|
||||
*/
|
||||
@Test
|
||||
void aTaskRunningLongerThanTheTtlStillKeepsItsReport() throws Exception {
|
||||
java.util.concurrent.atomic.AtomicLong clock = new java.util.concurrent.atomic.AtomicLong(1_000_000_000L);
|
||||
try (var wiring = wireWithPushLoop(1, 50, clock::get)) {
|
||||
String slow = wiring.service().sendAsync(T, "a task that takes longer than the TTL");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
// The worker is still working, and has been for longer than the entire TTL. Nothing may
|
||||
// be pruned yet — the ticket has not finished, so there is no report to keep or lose.
|
||||
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(30));
|
||||
|
||||
// Only now does it reply. Under the old clock this reply was born already expired.
|
||||
assertTrue(rendezvous.resolve(T, "the long report"));
|
||||
awaitTicketPhaseOn(wiring.service(), slow, MessageService.Phase.DONE);
|
||||
|
||||
// A second delegation runs pruneTerminalTickets before it returns.
|
||||
wiring.service().sendAsync(T, "an unrelated second task");
|
||||
|
||||
MessageService.TaskView view = wiring.service().poll(slow);
|
||||
assertNotNull(view, "a ticket that completed just now must survive the prune, however "
|
||||
+ "long its task ran — the TTL is the window to COLLECT the report, not the "
|
||||
+ "budget for producing it");
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
assertEquals("the long report", view.reply(),
|
||||
"the worker's actual report must still be there, not just the ticket");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The other half of #197: the TTL must still bound {@code tasks}. Measuring from completion
|
||||
* would be a leak if a finished ticket were then kept forever, so this pins the eviction that
|
||||
* still has to happen — the same ticket, left uncollected for longer than the TTL AFTER it
|
||||
* finished, is gone.
|
||||
*/
|
||||
@Test
|
||||
void aFinishedTicketIsStillPrunedOnceTheTtlPassesSinceItFinished() throws Exception {
|
||||
java.util.concurrent.atomic.AtomicLong clock = new java.util.concurrent.atomic.AtomicLong(1_000_000_000L);
|
||||
try (var wiring = wireWithPushLoop(1, 50, clock::get)) {
|
||||
String done = wiring.service().sendAsync(T, "a quick task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "quick result"));
|
||||
awaitTicketPhaseOn(wiring.service(), done, MessageService.Phase.DONE);
|
||||
|
||||
// Nobody collected it, and the TTL has now passed since it FINISHED.
|
||||
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
|
||||
wiring.service().sendAsync(T, "an unrelated second task");
|
||||
|
||||
assertNull(wiring.service().poll(done),
|
||||
"the TTL must still evict an uncollected finished ticket, or tasks grows forever");
|
||||
}
|
||||
}
|
||||
|
||||
private MessageService.TaskView awaitTicketPhaseOn(MessageService svc, String ticket,
|
||||
MessageService.Phase phase) throws Exception {
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
|
||||
@@ -1,13 +1,23 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.IThrowableProxy;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
@@ -352,6 +362,151 @@ class GitWorktreesTest {
|
||||
assertEquals("worktree origin contains HTTPS user info; refusing provision", error.getMessage());
|
||||
}
|
||||
|
||||
// ---- CB-189: broader remote-URL coverage — every remote, both fetch and push URLs, any
|
||||
// non-SSH scheme. Reporting only, additive to the origin/https strip-and-refuse tests above. ----
|
||||
|
||||
private Logger reportingLogger;
|
||||
private ListAppender<ILoggingEvent> reportingAppender;
|
||||
|
||||
/** {@link GitWorktrees}'s own logger, captured fresh for each test so assertions never see a
|
||||
* message left over from a previous test. */
|
||||
@BeforeEach
|
||||
void attachReportingLogCapture() {
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
reportingLogger = ctx.getLogger(GitWorktrees.class);
|
||||
reportingLogger.setLevel(Level.WARN);
|
||||
reportingAppender = new ListAppender<>();
|
||||
reportingAppender.setContext(ctx);
|
||||
reportingAppender.start();
|
||||
reportingLogger.addAppender(reportingAppender);
|
||||
}
|
||||
|
||||
@AfterEach
|
||||
void detachReportingLogCapture() {
|
||||
reportingLogger.detachAppender(reportingAppender);
|
||||
}
|
||||
|
||||
private List<String> capturedMessages() {
|
||||
return reportingAppender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
}
|
||||
|
||||
/** Asserts {@code secret} appears in no captured message, and in no attached exception's
|
||||
* message either — the constraint is that a credential must never reach a log, however it
|
||||
* would have gotten there. */
|
||||
private void assertNoLeak(String secret) {
|
||||
for (ILoggingEvent event : reportingAppender.list) {
|
||||
assertFalse(event.getFormattedMessage().contains(secret),
|
||||
"log message leaked a credential (" + secret + "): " + event.getFormattedMessage());
|
||||
IThrowableProxy thrown = event.getThrowableProxy();
|
||||
if (thrown != null && thrown.getMessage() != null) {
|
||||
assertFalse(thrown.getMessage().contains(secret),
|
||||
"logged exception leaked a credential (" + secret + "): " + thrown.getMessage());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Gap 1: only {@code origin} was ever inspected. A credential on any other remote's fetch URL
|
||||
* must now be reported. */
|
||||
@Test
|
||||
void aCredentialedUrlOnANonOriginRemoteIsReported(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
|
||||
git(repo, "remote", "add", "upstream", "https://leaky-upstream-token@git.ltms.dev/akb/kb.git");
|
||||
|
||||
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-189-a", "HEAD");
|
||||
|
||||
List<String> messages = capturedMessages();
|
||||
assertTrue(messages.stream().anyMatch(m -> m.contains("remote=upstream")),
|
||||
"expected a report naming the leaking non-origin remote:\n" + messages);
|
||||
assertNoLeak("leaky-upstream-token");
|
||||
assertNoLeak("https://leaky-upstream-token@git.ltms.dev/akb/kb.git");
|
||||
assertNoLeak("git.ltms.dev");
|
||||
}
|
||||
|
||||
/** Gap 2: push URLs were never inspected. A credential visible only on {@code pushurl} — the
|
||||
* fetch URL for the same remote stays clean — must now be reported. */
|
||||
@Test
|
||||
void aCredentialedPushUrlIsReported(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
|
||||
git(repo, "remote", "add", "mirror", "https://git.ltms.dev/akb/mirror.git");
|
||||
git(repo, "remote", "set-url", "--push", "mirror",
|
||||
"https://leaky-push-token@git.ltms.dev/akb/mirror.git");
|
||||
|
||||
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-189-b", "HEAD");
|
||||
|
||||
List<String> messages = capturedMessages();
|
||||
assertTrue(messages.stream().anyMatch(m -> m.contains("remote=mirror")),
|
||||
"expected a report naming the remote with the leaking pushurl:\n" + messages);
|
||||
assertNoLeak("leaky-push-token");
|
||||
assertNoLeak("https://leaky-push-token@git.ltms.dev/akb/mirror.git");
|
||||
assertNoLeak("git.ltms.dev");
|
||||
}
|
||||
|
||||
/** Gap 3: only {@code https} was handled. A plain {@code http://user:pass@…} remote — worse
|
||||
* than https, not better — must now be reported. */
|
||||
@Test
|
||||
void anHttpUrlWithCredentialsIsReported(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
|
||||
git(repo, "remote", "add", "insecure", "http://plainuser:plainpass@git.ltms.dev/akb/kb.git");
|
||||
|
||||
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-189-c", "HEAD");
|
||||
|
||||
List<String> messages = capturedMessages();
|
||||
assertTrue(messages.stream().anyMatch(m -> m.contains("remote=insecure")),
|
||||
"expected a report for the credentialed plain-http remote:\n" + messages);
|
||||
assertNoLeak("plainuser");
|
||||
assertNoLeak("plainpass");
|
||||
assertNoLeak("plainuser:plainpass");
|
||||
assertNoLeak("git.ltms.dev");
|
||||
}
|
||||
|
||||
/** A normal {@code ssh://} remote and a credential-free {@code https://} remote must produce no
|
||||
* report at all — the check must not cry wolf on ordinary, safe configuration. */
|
||||
@Test
|
||||
void anSshRemoteAndACleanHttpsRemoteProduceNoReport(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "ssh://git@git.ltms.dev:2224/akb/kb.git");
|
||||
git(repo, "remote", "add", "clean", "https://git.ltms.dev/akb/kb.git");
|
||||
|
||||
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-189-d", "HEAD");
|
||||
|
||||
assertTrue(reportingAppender.list.isEmpty(),
|
||||
"expected no report for an ssh remote and a credential-free https remote, got:\n"
|
||||
+ capturedMessages());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-189 review fix. Even a FAILING command that read a credential onto its stdout must never
|
||||
* let that value reach the thrown {@link WorktreeException}'s message — this is gap 4 from the
|
||||
* CB-189 issue, and the reason {@link GitWorktrees#execRedacted} exists at all. Drives the
|
||||
* shared {@code exec}/{@code execRedacted} seam directly (it is package-private for exactly this,
|
||||
* the same way the {@code afterWorktreeAdded} constructor parameter is a test seam) with a
|
||||
* synthetic, non-git command whose stdout carries a marker — passed through the environment,
|
||||
* never through argv, so the marker cannot leak via the command line that IS always printed
|
||||
* unconditionally in the exception message — and which exits non-zero. This isolates the
|
||||
* redaction guarantee itself rather than depending on a specific git failure mode that happens to
|
||||
* echo a URL onto stdout before failing: none of the git subcommands this class actually runs was
|
||||
* found to have one (a corrupted config makes {@code git config --get} fail before it ever reads
|
||||
* the target key, so its output never carries the URL either). The marker is generated per-test
|
||||
* run and injected only by the test, never a real-looking credential, so even a failing assertion
|
||||
* could not itself print a secret.
|
||||
*/
|
||||
@Test
|
||||
void execRedactedNeverCopiesFailingCommandOutputIntoTheExceptionMessage(@TempDir Path tmp) {
|
||||
String marker = "cb189-marker-" + System.nanoTime();
|
||||
GitWorktrees worktrees = new GitWorktrees(tmp.toString());
|
||||
|
||||
WorktreeException thrown = assertThrows(WorktreeException.class, () -> worktrees.exec(
|
||||
Map.of("MARKER", marker), true, "sh", "-c", "echo \"$MARKER\"; exit 7"));
|
||||
|
||||
assertNotNull(thrown.getMessage());
|
||||
assertFalse(thrown.getMessage().contains(marker),
|
||||
"a failing command's captured stdout leaked into the exception message: "
|
||||
+ thrown.getMessage());
|
||||
}
|
||||
|
||||
/** Neutralizing must not look like work in progress, or a worker would commit it into its PR. */
|
||||
@Test
|
||||
void theNeutralizedConfigIsNotAPendingLocalModification(@TempDir Path tmp) throws Exception {
|
||||
|
||||
Reference in New Issue
Block a user