Compare commits
21 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| fc2e26c0c5 | |||
| 4ac688b6d9 | |||
| a814d1ef00 | |||
| ea98856130 | |||
| a49e96835a | |||
| 63c19dcba7 | |||
| 2823349c8e | |||
| 6fc301d62c | |||
| bc99d64786 | |||
| 615af4ed0a | |||
| 045d229728 | |||
| d1fd5700f5 | |||
| 2d55b0b9a5 | |||
| d89ae94a2e | |||
| e6193c4098 | |||
| b66f0677ed | |||
| 23ada1981e | |||
| a237fbff9d | |||
| 31d5516991 | |||
| 5ba05d0bdb | |||
| 25726a5ae7 |
@@ -0,0 +1,52 @@
|
||||
# CB-175 report
|
||||
|
||||
## Change
|
||||
|
||||
`OpenCodeLauncher` reads the newest opencode session record for the worker cwd after spawn readiness.
|
||||
It compares the requested `provider/model` selector with `model.providerID/model.id` from the record.
|
||||
`variant` is not compared because a profile selector has no variant part.
|
||||
|
||||
An absent record, malformed record, or incomplete model object is unknown evidence. It does not log
|
||||
an error or quarantine the profile.
|
||||
|
||||
On a real mismatch, fleetd logs an ERROR with the requested and resolved selectors. The mismatch goes
|
||||
through `ExhaustionSink` into the existing `BackendQuarantine` and uses the profile's
|
||||
`effectiveCredentialId()`.
|
||||
|
||||
I chose a permanent, process-lifetime quarantine. A withdrawn selector cannot become correct after a
|
||||
cooldown. A timed retry could silently use the paid fallback again. `fleet_list` will show the usual
|
||||
quarantine state, with a very large remaining time, until fleetd restarts after an operator fixes the
|
||||
profile.
|
||||
|
||||
## Tests
|
||||
|
||||
Added tests for an exact match, a mismatch, and missing or unreadable session storage.
|
||||
|
||||
I proved the mismatch test fails without the quarantine call. I commented out the call and ran:
|
||||
|
||||
```text
|
||||
mvn -Dtest=OpenCodeLauncherTest#differentResolvedModelPermanentlyQuarantinesTheProfile test
|
||||
```
|
||||
|
||||
The result was:
|
||||
|
||||
```text
|
||||
[ERROR] Tests run: 1, Failures: 1, Errors: 0, Skipped: 0
|
||||
org.opentest4j.AssertionFailedError: a fallback model must block later spawns ==> expected: <true> but was: <false>
|
||||
[INFO] BUILD FAILURE
|
||||
```
|
||||
|
||||
I restored the call. I then ran `mvn clean install` in `fleetd/` without a pipe. Its result was:
|
||||
|
||||
```text
|
||||
[INFO] Tests run: 1040, Failures: 0, Errors: 0, Skipped: 0
|
||||
[INFO] BUILD SUCCESS
|
||||
```
|
||||
|
||||
## Limits and scope
|
||||
|
||||
I could not spawn a real opencode member or restart fleetd. I did not test this end to end against a
|
||||
live opencode session database.
|
||||
|
||||
I confirmed `ClaudeCodeLauncher` passes `--model` but does not read back the resolved model. I did
|
||||
not change it because it is outside this ticket's scope.
|
||||
@@ -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.
|
||||
@@ -168,6 +168,17 @@ public final class Fleetd {
|
||||
claudeProfiles.put(name, w);
|
||||
}
|
||||
});
|
||||
// One tracker covers timed backend exhaustion and permanent model-selector mismatches. The
|
||||
// latter cannot heal on a retry, so OpenCodeLauncher uses quarantinePermanently through the
|
||||
// sink below rather than letting a cooldown reopen a paid fallback.
|
||||
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
|
||||
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
|
||||
ExhaustionSink modelMismatchSink = (profileName, reason) -> {
|
||||
FleetConfig.Profile profile = config.get().profiles().get(profileName);
|
||||
if (profile != null) {
|
||||
quarantine.quarantinePermanently(profile.effectiveCredentialId());
|
||||
}
|
||||
};
|
||||
List<HerdrPeerLauncher> adapters = new ArrayList<>();
|
||||
// The claude-code adapter is the always-present default; keep it even with no profiles (so a
|
||||
// bridge configured with no workers, or opencode-only, still has a well-defined base adapter)
|
||||
@@ -184,15 +195,13 @@ public final class Fleetd {
|
||||
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
() -> config.get().memberCredentials()));
|
||||
() -> config.get().memberCredentials(), modelMismatchSink));
|
||||
}
|
||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
||||
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
||||
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
|
||||
// The cooldown is deferred (see FleetConfig#quarantineCooldownSeconds): it is read once
|
||||
// here, at startup, and a config reload only changes it for a daemon restart.
|
||||
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
|
||||
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
|
||||
PeerLauncher workers = new CompositePeerLauncher(
|
||||
adapters,
|
||||
cfg.effectiveDefaultProfile(),
|
||||
|
||||
@@ -38,6 +38,11 @@ public final class AgentControl {
|
||||
this.herdr = herdr;
|
||||
}
|
||||
|
||||
/** The herdr daemon this control object sends its agent calls to. */
|
||||
public HerdrClient herdr() {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
/** One agent-targeted call, translating a terminal id to its pane id (retrying once fresh). */
|
||||
private JsonNode agentCall(String method, String target, Map<String, Object> extra) {
|
||||
String resolved = resolveTarget(target);
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -2,6 +2,8 @@ package dev.ltms.fleet.member;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.peer.PeerHandle;
|
||||
@@ -21,6 +23,7 @@ import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.EnumSet;
|
||||
import java.util.HashSet;
|
||||
import java.util.IdentityHashMap;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -44,12 +47,17 @@ import java.util.stream.Collectors;
|
||||
* the single adapter that declares it. Profiles partition cleanly across adapters: the
|
||||
* constructor rejects a name claimed by two.</li>
|
||||
* <li><strong>By pane id</strong> — {@link #stop} routes to the adapter that spawned that pane
|
||||
* (recorded at spawn time). A pane the composite never spawned (only real for a caller that
|
||||
* hand-rolls an id) falls back to the first delegate; teardown is pane-id addressed and
|
||||
* tab cleanup is single-occupant guarded, so it is safe either way.</li>
|
||||
* (recorded at spawn time). A pane the composite never spawned, or one whose record was lost
|
||||
* to a daemon restart (CB-185 blocker 1 — {@link #spawnedBy} is in-memory only), can use the
|
||||
* fallback route in a one-daemon fleet. With more than one herdr daemon, {@link #probeOwner}
|
||||
* asks each configured daemon which one actually knows the pane: exactly one match routes
|
||||
* (and caches); no match is treated as already-gone; more than one match is a genuine
|
||||
* ambiguity (pane ids are per-daemon counters, so two daemons really can both hold, say,
|
||||
* {@code w1:p1}) and stop refuses rather than closing a pane on an arbitrary herdr daemon.</li>
|
||||
* <li><strong>Fleet-wide</strong> — {@link #reapOrphanWorkers} and {@link #capabilities} fan out
|
||||
* and combine. {@link #list} is deduplicated by pane id because every herdr-backed delegate
|
||||
* shares one herdr connection and so reports the same global agent set.</li>
|
||||
* and combine. {@link #list} is deduplicated by (owning daemon, pane id): delegates that share
|
||||
* one herdr connection report the same global agent set, but two daemons can each hold a pane
|
||||
* called {@code w1:p1}, so the daemon has to be part of the key.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>CB-518: an unqualified spawn is routed through a {@link PlacementPolicy}. The default
|
||||
@@ -435,12 +443,98 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
HerdrPeerLauncher d = spawnedBy.remove(id);
|
||||
HerdrPeerLauncher d = spawnedBy.get(id);
|
||||
if (d == null) {
|
||||
log.debug("stop({}) — no recorded owner, routing to the first adapter (pane-addressed)", id);
|
||||
d = delegates.getFirst();
|
||||
if (herdrDaemonCount() == 1) {
|
||||
log.debug("stop({}) — no recorded owner in a single-daemon fleet", id);
|
||||
d = delegates.getFirst();
|
||||
} else {
|
||||
d = probeOwner(id);
|
||||
if (d == null) {
|
||||
// No configured herdr daemon has ever heard of this pane. CB-185 blocker 1: this
|
||||
// is the normal case right after a daemon restart empties spawnedBy for a member
|
||||
// that has ALREADY been torn down since — the caller retried a stop that already
|
||||
// succeeded. Nothing to close and no owner to cache; matching the tolerance
|
||||
// HerdrPeerLauncher#stop already gives an already-gone pane (agent.close swallows
|
||||
// that as success), stop() here is a no-op rather than a refusal.
|
||||
log.debug("stop({}) — no configured herdr daemon knows this pane; "
|
||||
+ "treating as already stopped", id);
|
||||
return;
|
||||
}
|
||||
}
|
||||
}
|
||||
// Drop the owner record only after the delegate accepted the stop. Removing it first meant a
|
||||
// delegate that threw left the pane alive with its owner forgotten, so the retry fell into
|
||||
// the ambiguous branch above and refused the id for good.
|
||||
d.stop(id);
|
||||
spawnedBy.remove(id);
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-185 blocker 1: recover a spawnedBy cache miss by asking every distinct herdr daemon which
|
||||
* one actually knows {@code id} — the fix for "after a restart, every surviving member becomes
|
||||
* un-stoppable" (spawnedBy is in-memory only, so a restart empties it, and members intentionally
|
||||
* outlive the daemon).
|
||||
*
|
||||
* <p>Grouped by daemon identity, not by delegate, for the same reason {@link #list()} groups
|
||||
* that way: two adapters (claude-code, opencode) sharing one herdr connection would otherwise be
|
||||
* probed twice, and a pane on their shared daemon would look owned by two adapters instead of
|
||||
* one daemon.
|
||||
*
|
||||
* <p>A daemon that fails to answer {@code list()} (e.g. it is down) is treated as "does not know
|
||||
* this pane" rather than aborting the whole probe — one unreachable daemon must never make a
|
||||
* pane that a <em>different</em>, healthy daemon actually owns un-stoppable too, which would
|
||||
* resurrect the exact bug this method exists to fix.
|
||||
*
|
||||
* @return the owning delegate — cached into {@link #spawnedBy} so the next call is free — or
|
||||
* {@code null} when no daemon knows the pane
|
||||
* @throws IllegalArgumentException when more than one daemon claims the pane: pane ids are
|
||||
* per-daemon counters, so two daemons really can both hold, say, {@code w1:p1}, and there
|
||||
* is no way to tell which one the caller means
|
||||
*/
|
||||
private HerdrPeerLauncher probeOwner(String id) {
|
||||
Map<HerdrClient, HerdrPeerLauncher> byDaemon = new IdentityHashMap<>();
|
||||
for (HerdrPeerLauncher delegate : delegates) {
|
||||
byDaemon.putIfAbsent(delegate.herdr(), delegate);
|
||||
}
|
||||
List<HerdrPeerLauncher> owners = new ArrayList<>();
|
||||
for (HerdrPeerLauncher representative : byDaemon.values()) {
|
||||
List<Agent> agents;
|
||||
try {
|
||||
agents = representative.list();
|
||||
} catch (HerdrException e) {
|
||||
log.warn("stop({}) probe: a configured herdr daemon was unreachable ({}); "
|
||||
+ "treating it as not knowing this pane", id, e.getClass().getSimpleName());
|
||||
continue;
|
||||
}
|
||||
boolean knows = agents.stream().anyMatch(a -> id.equals(a.paneId()));
|
||||
if (knows) {
|
||||
owners.add(representative);
|
||||
}
|
||||
}
|
||||
if (owners.size() > 1) {
|
||||
throw new IllegalArgumentException("ambiguous paneId '" + id + "': "
|
||||
+ owners.size() + " configured herdr daemons report this pane — "
|
||||
+ "no way to tell which one the caller means");
|
||||
}
|
||||
if (owners.isEmpty()) {
|
||||
return null;
|
||||
}
|
||||
HerdrPeerLauncher owner = owners.get(0);
|
||||
spawnedBy.put(id, owner);
|
||||
return owner;
|
||||
}
|
||||
|
||||
/**
|
||||
* Count actual herdr daemons, not peer adapter kinds. Identity is intentional: separate client
|
||||
* objects may represent different daemons even if a client later implements value equality.
|
||||
*/
|
||||
private int herdrDaemonCount() {
|
||||
Set<HerdrClient> daemons = Collections.newSetFromMap(new IdentityHashMap<>());
|
||||
for (HerdrPeerLauncher delegate : delegates) {
|
||||
daemons.add(delegate.herdr());
|
||||
}
|
||||
return daemons.size();
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -475,14 +569,24 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
return route(profileName).capabilities();
|
||||
}
|
||||
|
||||
/** Every herdr agent, deduplicated by pane id (all delegates share one herdr and list globally). */
|
||||
/**
|
||||
* Every herdr agent, deduplicated by (owning daemon, pane id).
|
||||
*
|
||||
* <p>Delegates that share one {@link HerdrClient} see the same global agent set, so listing them
|
||||
* both would report every agent twice — that is what the dedupe is for. But pane ids are
|
||||
* per-daemon counters, so two daemons really can both hold {@code w1:p1} on different panes.
|
||||
* Keying on the pane id alone would silently drop one of them from {@code fleet_list} and from
|
||||
* every status view built on it. The daemon is part of the key for exactly that reason.
|
||||
*/
|
||||
@Override
|
||||
public List<Agent> list() {
|
||||
Map<HerdrClient, Integer> daemonIndex = new IdentityHashMap<>();
|
||||
Map<String, Agent> byPane = new LinkedHashMap<>();
|
||||
for (HerdrPeerLauncher d : delegates) {
|
||||
int daemon = daemonIndex.computeIfAbsent(d.herdr(), _ -> daemonIndex.size());
|
||||
for (Agent a : d.list()) {
|
||||
if (a.paneId() != null) {
|
||||
byPane.putIfAbsent(a.paneId(), a);
|
||||
byPane.putIfAbsent(daemon + "\u0000" + a.paneId(), a);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package dev.ltms.fleet.member;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.Tab;
|
||||
import dev.ltms.fleet.herdr.Workspace;
|
||||
@@ -520,6 +521,11 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
req.sessionName(), spawned.agentSessionId(), spawned.receipt());
|
||||
}
|
||||
|
||||
/** The herdr daemon that owns this launcher's pane coordinates. */
|
||||
public HerdrClient herdr() {
|
||||
return agents.herdr();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
return effectiveCwd(req.profileName(), req.requestedCwd(), req.callerCwd());
|
||||
@@ -1002,6 +1008,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();
|
||||
@@ -1010,8 +1023,8 @@ 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. */
|
||||
@@ -1061,12 +1074,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 "
|
||||
@@ -1171,18 +1187,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()
|
||||
@@ -1193,7 +1248,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 "
|
||||
|
||||
@@ -6,6 +6,7 @@ import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.CharterReceipt;
|
||||
import dev.ltms.fleet.peer.PeerHandle;
|
||||
@@ -71,6 +72,9 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
*/
|
||||
private final OpenCodeSessionDiscovery discovery;
|
||||
|
||||
/** Receives the profile when a resolved model differs from its requested selector. */
|
||||
private final ExhaustionSink modelMismatchSink;
|
||||
|
||||
/**
|
||||
* Production constructor — disables the spawn-ready gate ({@code spawnReadyTimeoutMs == 0}) so it
|
||||
* matches the legacy non-blocking spawn semantics. Config dirs are created under the JVM temp dir.
|
||||
@@ -121,7 +125,20 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials);
|
||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, ExhaustionSink.none());
|
||||
}
|
||||
|
||||
/** Production constructor with permanent-quarantine wiring for a model mismatch. */
|
||||
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,
|
||||
ExhaustionSink modelMismatchSink) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials, modelMismatchSink);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -163,12 +180,10 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
Function<String, String> env,
|
||||
long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper,
|
||||
Path configRoot, Path discoveryRoot,
|
||||
Supplier<FleetConfig.Fleet> fleet) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet);
|
||||
this.configRoot = configRoot;
|
||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||
Path configRoot, Path discoveryRoot,
|
||||
Supplier<FleetConfig.Fleet> fleet) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs, nowMillis, sleeper,
|
||||
configRoot, discoveryRoot, fleet, null, ExhaustionSink.none());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -179,13 +194,32 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
Function<String, String> env,
|
||||
long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper,
|
||||
Path configRoot, Path discoveryRoot,
|
||||
Supplier<FleetConfig.Fleet> fleet,
|
||||
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
|
||||
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, ExhaustionSink.none());
|
||||
}
|
||||
|
||||
/**
|
||||
* Full constructor with the model-mismatch quarantine callback. The callback is an
|
||||
* {@link ExhaustionSink} so model mismatches use the existing quarantine path rather than a
|
||||
* second state tracker.
|
||||
*/
|
||||
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,
|
||||
ExhaustionSink modelMismatchSink) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
|
||||
this.configRoot = configRoot;
|
||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||
this.modelMismatchSink = modelMismatchSink;
|
||||
}
|
||||
|
||||
private static Path defaultConfigRoot() {
|
||||
@@ -498,9 +532,34 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
PeerHandle inner = super.spawn(req);
|
||||
verifyResolvedModel(requireProfile(req.profileName()), effectiveCwd(req));
|
||||
return new SessionAwareHandle(inner, discovery, effectiveCwd(req));
|
||||
}
|
||||
|
||||
/**
|
||||
* Read the session record once after the spawn-readiness gate. opencode writes the model it
|
||||
* actually selected there. No record, incomplete model object, or a selector without one slash
|
||||
* is unknown evidence, so it must not quarantine a working profile.
|
||||
*/
|
||||
private void verifyResolvedModel(FleetConfig.Profile cfg, String cwd) {
|
||||
String[] requested = splitProviderModel(cfg.model());
|
||||
if (requested == null) {
|
||||
return;
|
||||
}
|
||||
OpenCodeSessionDiscovery.SessionRecord record = discovery.sessionForDirectory(cwd);
|
||||
if (record == null || record.providerId() == null || record.modelId() == null) {
|
||||
return;
|
||||
}
|
||||
String actual = record.providerId() + "/" + record.modelId();
|
||||
if (cfg.model().equals(actual)) {
|
||||
return;
|
||||
}
|
||||
log.error("opencode model mismatch for profile '{}': requested '{}' but resolved '{}'",
|
||||
cfg.profile(), cfg.model(), actual);
|
||||
modelMismatchSink.onExhausted(cfg.profile(), "opencode model mismatch: requested "
|
||||
+ cfg.model() + ", resolved " + actual);
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@link PeerHandle} that delegates everything to the base's worker handle but resolves
|
||||
* {@link #agentSessionId()} lazily through opencode session discovery. Delegate-only, so the
|
||||
|
||||
@@ -58,6 +58,15 @@ final class OpenCodeSessionDiscovery {
|
||||
* @return the matching session id, or {@code null} if none is known yet
|
||||
*/
|
||||
String sessionIdForDirectory(String directory) {
|
||||
SessionRecord record = sessionForDirectory(directory);
|
||||
return record == null ? null : record.id();
|
||||
}
|
||||
|
||||
/**
|
||||
* The newest session record for {@code directory}, or {@code null} when opencode has not written
|
||||
* one yet. This is the single storage seam for both session identity and resolved-model checks.
|
||||
*/
|
||||
SessionRecord sessionForDirectory(String directory) {
|
||||
if (directory == null || directory.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
@@ -65,20 +74,20 @@ final class OpenCodeSessionDiscovery {
|
||||
if (!Files.isDirectory(sessionRoot)) {
|
||||
return null;
|
||||
}
|
||||
String best = null;
|
||||
SessionRecord 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) {
|
||||
SessionRecord matched = matchRecord(record, directory);
|
||||
if (matched == null) {
|
||||
continue;
|
||||
}
|
||||
long mtime = lastModifiedEpochMillis(record);
|
||||
if (mtime > bestMtime) {
|
||||
bestMtime = mtime;
|
||||
best = id;
|
||||
best = matched;
|
||||
}
|
||||
}
|
||||
} catch (IOException ignored) {
|
||||
@@ -97,7 +106,7 @@ final class OpenCodeSessionDiscovery {
|
||||
* 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) {
|
||||
private SessionRecord matchRecord(Path record, String directory) {
|
||||
try {
|
||||
JsonNode node = json.readTree(record.toFile());
|
||||
JsonNode id = node == null ? null : node.get("id");
|
||||
@@ -105,12 +114,24 @@ final class OpenCodeSessionDiscovery {
|
||||
if (id == null || dir == null || !directory.equals(dir.asText())) {
|
||||
return null;
|
||||
}
|
||||
return id.asText();
|
||||
JsonNode model = node.path("model");
|
||||
String providerId = text(model, "providerID");
|
||||
String modelId = text(model, "id");
|
||||
return new SessionRecord(id.asText(), providerId, modelId);
|
||||
} catch (IOException e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private static String text(JsonNode node, String name) {
|
||||
JsonNode value = node.get(name);
|
||||
return value == null || value.isNull() || value.asText().isBlank() ? null : value.asText();
|
||||
}
|
||||
|
||||
/** The opencode fields fleetd reads from one session record. Null model fields mean unknown. */
|
||||
record SessionRecord(String id, String providerId, String modelId) {
|
||||
}
|
||||
|
||||
/** The record's last-modified epoch ms, or {@code Long.MIN_VALUE} if unreadable (never wins). */
|
||||
private static long lastModifiedEpochMillis(Path record) {
|
||||
try {
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -721,7 +730,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 +830,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());
|
||||
}
|
||||
|
||||
@@ -76,6 +76,20 @@ public final class BackendQuarantine {
|
||||
quarantinedUntilNanos.put(credentialId, nowNanos.getAsLong() + cooldownNanos);
|
||||
}
|
||||
|
||||
/**
|
||||
* Quarantine {@code credentialId} until the daemon restarts. This is for a configuration error
|
||||
* that cannot heal with time, unlike an exhausted backend. A model selector that opencode silently
|
||||
* resolves to another model stays wrong until an operator changes the profile, so a cooldown would
|
||||
* start the same unsafe work again.
|
||||
*/
|
||||
public void quarantinePermanently(String credentialId) {
|
||||
Objects.requireNonNull(credentialId, "credentialId");
|
||||
if (inert) {
|
||||
return;
|
||||
}
|
||||
quarantinedUntilNanos.put(credentialId, Long.MAX_VALUE);
|
||||
}
|
||||
|
||||
/** Whether {@code credentialId} is quarantined right now. */
|
||||
public boolean isQuarantined(String credentialId) {
|
||||
return remainingNanos(credentialId) > 0;
|
||||
@@ -110,6 +124,7 @@ public final class BackendQuarantine {
|
||||
}
|
||||
|
||||
private static long toSecondsRoundedUp(long nanos) {
|
||||
return (nanos + 999_999_999L) / 1_000_000_000L;
|
||||
long seconds = nanos / 1_000_000_000L;
|
||||
return seconds + (nanos % 1_000_000_000L == 0 ? 0 : 1);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -214,6 +214,17 @@ public final class FleetApp {
|
||||
* second daemon configured, a member daemon that is down must not be masked by a healthy lead
|
||||
* daemon — every spawn goes through the member daemon and would otherwise fail silently behind
|
||||
* a green {@code /healthz}.
|
||||
*
|
||||
* <p>CB-185 blocker 2: the {@code herdr} key always carries the <em>lead</em> daemon's
|
||||
* version/protocol, unchanged, because two consumers — {@code scripts/redeploy-fleetd.sh} and
|
||||
* {@code scripts/rename-checkout.sh} — read this endpoint already (both only check the HTTP
|
||||
* status code and print the body verbatim; neither parses a specific field, so adding a key
|
||||
* alongside {@code herdr} is safe). But it is the <em>member</em> daemon's protocol that decides
|
||||
* whether a spawn works, so when a second daemon is configured its version/protocol is reported
|
||||
* too, under a separate {@code member} key — never folded into {@code herdr}, which would make a
|
||||
* mismatch invisible to whichever consumer only reads that key. If the two protocol numbers
|
||||
* differ, {@code protocolMismatch: true} calls it out explicitly rather than leaving it to be
|
||||
* spotted by comparing two numbers by eye.
|
||||
*/
|
||||
private void healthz(Context ctx) {
|
||||
JsonNode pong;
|
||||
@@ -226,9 +237,15 @@ public final class FleetApp {
|
||||
"detail", e.getMessage()));
|
||||
return;
|
||||
}
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
body.put("status", "ok");
|
||||
body.put("herdr", Map.of(
|
||||
"version", pong.path("version").asText(""),
|
||||
"protocol", pong.path("protocol").asInt()));
|
||||
if (memberHerdr != herdr) {
|
||||
JsonNode memberPong;
|
||||
try {
|
||||
memberHerdr.call("ping");
|
||||
memberPong = memberHerdr.call("ping");
|
||||
} catch (HerdrException e) {
|
||||
ctx.status(503).json(Map.of(
|
||||
"status", "degraded",
|
||||
@@ -236,12 +253,16 @@ public final class FleetApp {
|
||||
"detail", e.getMessage()));
|
||||
return;
|
||||
}
|
||||
int leadProtocol = pong.path("protocol").asInt();
|
||||
int memberProtocol = memberPong.path("protocol").asInt();
|
||||
body.put("member", Map.of(
|
||||
"version", memberPong.path("version").asText(""),
|
||||
"protocol", memberProtocol));
|
||||
if (leadProtocol != memberProtocol) {
|
||||
body.put("protocolMismatch", true);
|
||||
}
|
||||
}
|
||||
ctx.status(200).json(Map.of(
|
||||
"status", "ok",
|
||||
"herdr", Map.of(
|
||||
"version", pong.path("version").asText(""),
|
||||
"protocol", pong.path("protocol").asInt())));
|
||||
ctx.status(200).json(body);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -31,6 +31,8 @@ public final class FakeHerdr implements HerdrClient {
|
||||
*/
|
||||
public final List<Call> calls = new CopyOnWriteArrayList<>();
|
||||
private boolean healthy = true;
|
||||
private String pingVersion = "0.8.0";
|
||||
private int pingProtocol = 19;
|
||||
private final List<String> extraWorkspaces = new ArrayList<>();
|
||||
private final List<String> extraAgents = new ArrayList<>();
|
||||
/** workspaceId → extra tabs that {@code tab.list} reports for it (CB-558 lead scans). */
|
||||
@@ -52,6 +54,16 @@ public final class FakeHerdr implements HerdrClient {
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code ping} report this version/protocol instead of the default 0.8.0/19 — CB-185
|
||||
* blocker 2's fixture for a lead and a member daemon running mismatched herdr versions.
|
||||
*/
|
||||
public FakeHerdr pingReports(String version, int protocol) {
|
||||
this.pingVersion = version;
|
||||
this.pingProtocol = protocol;
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Reject the first {@code n} {@code agent.start} calls with {@code agent_name_taken}. */
|
||||
public FakeHerdr agentNameTakenTimes(int n) {
|
||||
this.agentNameTakenFor = n;
|
||||
@@ -166,7 +178,8 @@ public final class FakeHerdr implements HerdrClient {
|
||||
try {
|
||||
return switch (method) {
|
||||
case "ping" -> mapper.readTree(
|
||||
"{\"type\":\"pong\",\"version\":\"0.8.0\",\"protocol\":19}");
|
||||
("{\"type\":\"pong\",\"version\":\"%s\",\"protocol\":%d}")
|
||||
.formatted(pingVersion, pingProtocol));
|
||||
case "workspace.list" -> mapper.readTree(("""
|
||||
{"type":"workspace_list","workspaces":[
|
||||
{"workspace_id":"w1","label":"dev-mgnl","focused":true,"pane_count":7,"agent_status":"unknown"},
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -8,6 +8,7 @@ import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.CharterReceipt;
|
||||
@@ -248,6 +249,193 @@ class CompositePeerLauncherTest {
|
||||
"stop routes to the spawning adapter and closes exactly that worker's pane");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopKeepsSameHerdrPaneIdSeparateByOwningAdapter() {
|
||||
// Separate herdr daemons can both issue w9:pRoot_1. The opaque handles identify their
|
||||
// spawning adapters, so each stop reaches only its recorded owner.
|
||||
FakeHerdr first = new FakeHerdr();
|
||||
FakeHerdr second = new FakeHerdr();
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
PeerHandle claude = composite.spawn(new SpawnRequest("claude", null, null));
|
||||
PeerHandle opencode = composite.spawn(new SpawnRequest("gemini", null, null));
|
||||
|
||||
assertNotEquals(claude.id(), opencode.id(), "each public paneId keeps its adapter owner");
|
||||
composite.stop(claude.id());
|
||||
assertTrue(first.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w9:pRoot_1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"the first daemon closes its own pane");
|
||||
assertFalse(second.called("pane.close"), "the matching pane on the second daemon stays live");
|
||||
|
||||
composite.stop(opencode.id());
|
||||
assertTrue(second.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w9:pRoot_1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"the second daemon then closes its own pane");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopAllowsALegacyBarePaneIdWithOneDaemon() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
PeerLauncher composite = new CompositePeerLauncher(List.of(claudeAdapter(herdr)), "claude");
|
||||
|
||||
composite.stop("w9:pRoot_1");
|
||||
|
||||
assertTrue(herdr.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w9:pRoot_1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"one daemon keeps the legacy bare-pane routing behaviour");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopAllowsAnUnownedPaneIdWithTwoAdaptersSharingOneDaemon() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
PeerLauncher composite = composite(herdr);
|
||||
|
||||
composite.stop("w9:pRoot_1");
|
||||
|
||||
assertTrue(herdr.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w9:pRoot_1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"two adapter kinds sharing one daemon keep the fallback route");
|
||||
}
|
||||
|
||||
@Test
|
||||
void listKeepsBothPanesWhenTwoDaemonsShareAPaneId() {
|
||||
// herdr pane ids are per-daemon counters, so two daemons really can both hold w1:p1 on
|
||||
// different panes. Deduplicating on the pane id alone dropped one of the two real agents.
|
||||
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
|
||||
FakeHerdr second = new FakeHerdr().withAgent("y", "term_y", "w1:p1", "w1:t1");
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
List<Agent> agents = composite.list();
|
||||
|
||||
assertEquals(2, agents.stream().filter(a -> "w1:p1".equals(a.paneId())).count(),
|
||||
"one w1:p1 per daemon survives — the pane id alone is not a unique key");
|
||||
assertTrue(agents.stream().anyMatch(a -> "term_x".equals(a.terminalId())));
|
||||
assertTrue(agents.stream().anyMatch(a -> "term_y".equals(a.terminalId())));
|
||||
}
|
||||
|
||||
@Test
|
||||
void listStillDeduplicatesTwoAdaptersSharingOneDaemon() {
|
||||
// Both adapters ask the SAME daemon, so both see the same agent set. Without the dedupe this
|
||||
// would report every agent twice; the daemon key must not break that.
|
||||
FakeHerdr herdr = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
|
||||
CompositePeerLauncher composite = composite(herdr);
|
||||
|
||||
List<Agent> agents = composite.list();
|
||||
|
||||
assertEquals(1, agents.stream().filter(a -> "w1:p1".equals(a.paneId())).count(),
|
||||
"one daemon still reports each of its agents once");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopKeepsTheOwnerRecordWhenTheDelegateRefusesTheStop() {
|
||||
// Removing the record before the delegate accepted the stop lost the owner on failure: the
|
||||
// pane was still alive, but the retry landed in the ambiguous branch and refused it for good.
|
||||
FakeHerdr first = new FakeHerdr().paneCloseFailsWith("pane_busy");
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(new FakeHerdr())), "claude");
|
||||
PeerHandle claude = composite.spawn(new SpawnRequest("claude", null, null));
|
||||
|
||||
assertThrows(HerdrException.class, () -> composite.stop(claude.id()));
|
||||
|
||||
// The retry must still know its owner — a HerdrException, never "ambiguous paneId".
|
||||
assertThrows(HerdrException.class, () -> composite.stop(claude.id()),
|
||||
"the owner record survives a failed stop, so the retry is not ambiguous");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopRejectsAnUnownedPaneIdWhenMultipleDaemonsCouldOwnIt() {
|
||||
// CB-185 blocker 1: genuine ambiguity — pane ids are per-daemon counters, so two daemons
|
||||
// can each really hold an agent at "w1:p1". Neither claims ownership through spawnedBy
|
||||
// (empty, as after a restart), so the probe must find BOTH and refuse rather than guess.
|
||||
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
|
||||
FakeHerdr second = new FakeHerdr().withAgent("y", "term_y", "w1:p1", "w1:t1");
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
IllegalArgumentException error = assertThrows(IllegalArgumentException.class,
|
||||
() -> composite.stop("w1:p1"));
|
||||
|
||||
assertEquals("ambiguous paneId 'w1:p1': 2 configured herdr daemons report this pane — "
|
||||
+ "no way to tell which one the caller means", error.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopOnAPaneNoConfiguredDaemonKnowsIsTreatedAsAlreadyStopped() {
|
||||
// CB-185 blocker 1, the zero-owner branch: spawnedBy is empty (as after a restart) and
|
||||
// neither daemon's agent.list mentions this pane at all — it is already gone. A retried
|
||||
// stop() on an already-gone pane must succeed quietly, not refuse forever.
|
||||
FakeHerdr first = new FakeHerdr();
|
||||
FakeHerdr second = new FakeHerdr();
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
assertDoesNotThrow(() -> composite.stop("w1:p1"));
|
||||
|
||||
assertFalse(first.called("pane.close"), "no owner was found, so no delegate is told to close anything");
|
||||
assertFalse(second.called("pane.close"), "no owner was found, so no delegate is told to close anything");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopWithEmptySpawnedByResolvesTheOwnerThroughAProbeAndSkipsTheOtherDaemon() {
|
||||
// CB-185 blocker 1, the main fix: after a restart spawnedBy is empty for every surviving
|
||||
// member. stop() must still find the one daemon that actually knows the pane and route
|
||||
// only to it — never touching the daemon that never held it.
|
||||
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
|
||||
FakeHerdr second = new FakeHerdr();
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
composite.stop("w1:p1");
|
||||
|
||||
assertTrue(first.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w1:p1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"the daemon that actually knows the pane closes it");
|
||||
assertFalse(second.called("pane.close"), "the daemon that never held the pane is never touched");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aProbeSurvivesOneUnreachableDaemonAndStillFindsTheOwnerOnTheOtherOne() {
|
||||
// CB-185 blocker 1 (lead review): a daemon that is DOWN while we probe must not abort the
|
||||
// whole probe — the pane the OPERATOR actually wants stopped can live on a different,
|
||||
// healthy daemon, and that pane must not become un-stoppable because a third one is down.
|
||||
FakeHerdr down = new FakeHerdr().healthy(false);
|
||||
FakeHerdr owner = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(down), opencodeAdapter(owner)), "claude");
|
||||
|
||||
assertDoesNotThrow(() -> composite.stop("w1:p1"),
|
||||
"the unreachable daemon must be skipped, not fail the whole stop");
|
||||
|
||||
assertTrue(owner.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w1:p1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"the healthy daemon that actually owns the pane still closes it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aProbedOwnerIsCachedSoARetryAfterAFailedStopNeedsNoSecondProbe() {
|
||||
// CB-185 blocker 1: the probe's whole point is to be cheap on repeat — a failed stop (e.g.
|
||||
// "pane_busy") must not force another agent.list() round trip on every retry.
|
||||
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1")
|
||||
.paneCloseFailsWith("pane_busy");
|
||||
FakeHerdr second = new FakeHerdr();
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
assertThrows(HerdrException.class, () -> composite.stop("w1:p1"));
|
||||
long listCallsAfterFirst = first.calls.stream().filter(c -> c.method().equals("agent.list")).count()
|
||||
+ second.calls.stream().filter(c -> c.method().equals("agent.list")).count();
|
||||
assertTrue(listCallsAfterFirst > 0, "the first stop needed a probe");
|
||||
|
||||
assertThrows(HerdrException.class, () -> composite.stop("w1:p1"),
|
||||
"still failing on the retry, but through the cached owner");
|
||||
long listCallsAfterSecond = first.calls.stream().filter(c -> c.method().equals("agent.list")).count()
|
||||
+ second.calls.stream().filter(c -> c.method().equals("agent.list")).count();
|
||||
assertEquals(listCallsAfterFirst, listCallsAfterSecond,
|
||||
"the retry is served from the cache — no additional agent.list probe");
|
||||
}
|
||||
|
||||
@Test
|
||||
void opencodeContextResetIsANoOpAndWarnsOnlyOnce() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -12,6 +12,7 @@ import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.peer.PeerHandle;
|
||||
import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
@@ -22,6 +23,7 @@ import java.util.Map;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Future;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
@@ -61,6 +63,23 @@ class OpenCodeLauncherTest {
|
||||
0, System::currentTimeMillis, () -> { }, configRoot, configRoot, null, () -> creds);
|
||||
}
|
||||
|
||||
private static OpenCodeLauncher serviceWithModelMismatchSink(FakeHerdr herdr, Path root,
|
||||
FleetConfig.Profile cfg,
|
||||
BackendQuarantine quarantine) {
|
||||
return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
|
||||
0, System::currentTimeMillis, () -> { }, root, root, null, null,
|
||||
(profile, _) -> quarantine.quarantinePermanently(profile));
|
||||
}
|
||||
|
||||
private static void writeResolvedModelRecord(Path root, String directory, String providerId,
|
||||
String modelId) throws Exception {
|
||||
Path record = Files.createDirectories(root.resolve("session").resolve("p1")).resolve("ses_a.json");
|
||||
Files.writeString(record, "{\"id\":\"ses_a\",\"directory\":\"" + directory
|
||||
+ "\",\"model\":{\"id\":\"" + modelId + "\",\"providerID\":\""
|
||||
+ providerId + "\",\"variant\":\"high\"}}");
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
private static Map<String, Object> lastStart(FakeHerdr herdr) {
|
||||
return (Map<String, Object>) herdr.lastCall("agent.start").params();
|
||||
@@ -211,6 +230,50 @@ class OpenCodeLauncherTest {
|
||||
"--auto is unconditional: a model-less worker still must never block on approval");
|
||||
}
|
||||
|
||||
// --- CB-175: verify opencode's recorded resolved model after spawn readiness ----------------
|
||||
|
||||
@Test
|
||||
void matchingResolvedModelDoesNotQuarantineTheProfile(@TempDir Path root) throws Exception {
|
||||
FleetConfig.Profile cfg = opencodeCfg("opencode/x-preview-f-free", null, null);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
|
||||
writeResolvedModelRecord(root, root.toString(), "opencode", "x-preview-f-free");
|
||||
|
||||
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, quarantine)
|
||||
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
|
||||
|
||||
assertFalse(quarantine.isQuarantined("gemini"), "the exact provider/model match is safe");
|
||||
}
|
||||
|
||||
@Test
|
||||
void differentResolvedModelPermanentlyQuarantinesTheProfile(@TempDir Path root) throws Exception {
|
||||
FleetConfig.Profile cfg = opencodeCfg("opencode/x-preview-f-free", null, null);
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
|
||||
writeResolvedModelRecord(root, root.toString(), "openai", "gpt-5.6-sol");
|
||||
|
||||
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, quarantine)
|
||||
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
|
||||
|
||||
assertTrue(quarantine.isQuarantined("gemini"), "a fallback model must block later spawns");
|
||||
assertEquals(Long.MAX_VALUE / 1_000_000_000L + 1, quarantine.remainingSeconds("gemini").orElseThrow(),
|
||||
"a withdrawn selector cannot become safe after the normal cooldown");
|
||||
}
|
||||
|
||||
@Test
|
||||
void missingOrUnreadableSessionDatabaseDoesNotQuarantine(@TempDir Path root) throws Exception {
|
||||
FleetConfig.Profile cfg = opencodeCfg("opencode/x-preview-f-free", null, null);
|
||||
BackendQuarantine missing = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
|
||||
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, missing)
|
||||
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
|
||||
assertFalse(missing.isQuarantined("gemini"), "a missing database is unknown evidence");
|
||||
|
||||
Path broken = Files.createDirectories(root.resolve("session").resolve("p1")).resolve("ses_a.json");
|
||||
Files.writeString(broken, "not JSON");
|
||||
BackendQuarantine unreadable = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(1));
|
||||
serviceWithModelMismatchSink(new FakeHerdr(), root, cfg, unreadable)
|
||||
.spawn(new SpawnRequest(null, root.toString(), null, null, null, MemberRole.DEV));
|
||||
assertFalse(unreadable.isQuarantined("gemini"), "an unreadable database is unknown evidence");
|
||||
}
|
||||
|
||||
// --- CB-617: --agent <role> when the role has an agent-definition file --------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -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();
|
||||
@@ -1057,6 +1079,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;
|
||||
|
||||
@@ -96,4 +96,55 @@ class FleetAppTwoDaemonTest {
|
||||
long calls = shared.calls.stream().filter(c -> c.method().equals("workspace.list")).count();
|
||||
assertEquals(1, calls, "single-daemon deployment must call workspace.list exactly once");
|
||||
}
|
||||
|
||||
// ── CB-185 blocker 2: /healthz must report the MEMBER daemon's protocol too ────────────────
|
||||
|
||||
@Test
|
||||
void healthzReportsBothDaemonsWhenTheirProtocolsDiffer() throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr().pingReports("0.8.0", 19);
|
||||
FakeHerdr member = new FakeHerdr().pingReports("0.7.0", 18);
|
||||
int port = start(lead, member);
|
||||
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
|
||||
assertEquals(200, res.statusCode(), res.body());
|
||||
assertTrue(res.body().contains("\"protocol\":19"),
|
||||
"the herdr key keeps reporting the LEAD's protocol, unchanged: " + res.body());
|
||||
assertTrue(res.body().contains("\"member\""), "a separate member key is present: " + res.body());
|
||||
assertTrue(res.body().contains("\"protocol\":18"),
|
||||
"the member key reports the member daemon's own protocol: " + res.body());
|
||||
assertTrue(res.body().contains("\"protocolMismatch\":true"),
|
||||
"a differing protocol is called out explicitly, not left to be spotted by eye: " + res.body());
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzReportsBothDaemonsWithNoMismatchWhenProtocolsMatch() throws Exception {
|
||||
int port = start(new FakeHerdr(), new FakeHerdr());
|
||||
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
|
||||
assertEquals(200, res.statusCode(), res.body());
|
||||
assertTrue(res.body().contains("\"member\""), "the member key is present whenever a second daemon "
|
||||
+ "is configured, even when the protocols happen to agree: " + res.body());
|
||||
assertFalse(res.body().contains("protocolMismatch"),
|
||||
"matching protocols must not raise a mismatch flag: " + res.body());
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzWithOneDaemonCarriesNoMemberOrMismatchKey() throws Exception {
|
||||
// The single-daemon deployment (no memberHerdrSocket) must see no change at all beyond the
|
||||
// historical body: no "member" key, no "protocolMismatch" key. (Map.of()'s own key order is
|
||||
// JVM-salted regardless of this fix, so this checks content, not exact key order.)
|
||||
FakeHerdr shared = new FakeHerdr();
|
||||
int port = start(shared, shared);
|
||||
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
|
||||
assertEquals(200, res.statusCode());
|
||||
assertTrue(res.body().contains("\"status\":\"ok\""), res.body());
|
||||
assertTrue(res.body().contains("\"protocol\":19"), res.body());
|
||||
assertTrue(res.body().contains("\"version\":\"0.8.0\""), res.body());
|
||||
assertFalse(res.body().contains("\"member\""), "no second daemon configured, so no member key: " + res.body());
|
||||
assertFalse(res.body().contains("protocolMismatch"), res.body());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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