Compare commits
30 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a22480c117 | |||
| 24f404f989 | |||
| 17af61e8dd | |||
| 6af87b6ad6 | |||
| fc655e78c2 | |||
| 11c3ff67b6 | |||
| d867c87100 | |||
| b0c4cedfab | |||
| 85417d5215 | |||
| 4accc746bd | |||
| 21c4c8cbef | |||
| 430f5b0dae | |||
| 42731833d0 | |||
| 65acf066ad | |||
| ee8f570fd7 | |||
| fa3f910d44 | |||
| 46ac6e4e38 | |||
| 3bfa82839b | |||
| 65ccf2e4ad | |||
| 7a3b27f76f | |||
| 82e7be564c | |||
| c5e24197bf | |||
| bcb402b688 | |||
| 7c4170ff6d | |||
| a6095743f0 | |||
| 66e776d178 | |||
| 26f64cba45 | |||
| 51047848f1 | |||
| d8c0b657e8 | |||
| 312c0584ce |
@@ -27,9 +27,14 @@ can make every broker probe look empty.
|
||||
reaches the report:
|
||||
|
||||
```bash
|
||||
sed -E 's#://[^@]*@#://<redacted>@#'
|
||||
sed -E 's#://[^@]*@#://<redacted>@#g'
|
||||
```
|
||||
|
||||
**The `g` flag is not optional.** Without it `sed` replaces only the first match on each line, so a
|
||||
line carrying two URIs leaks the second one. `scripts/redeploy-fleetd.sh --check` prints lines like
|
||||
that. Checked on 2026-08-27: without `g`, `amqp://u1:p1@h1/mac and http://u2:p2@h2:15672/api`
|
||||
redacts the first pair and prints `u2:p2` in the clear.
|
||||
|
||||
Keep `pipefail` on when applying that filter. Otherwise the filter can hide a failed probe. Apply
|
||||
the same no-print rule to the management password below, even though it is not in an AMQP URI.
|
||||
|
||||
@@ -44,7 +49,7 @@ Do not copy those checks into new shell code. The script reads `LAVINMQ_URI`, so
|
||||
```bash
|
||||
set -o pipefail
|
||||
scripts/redeploy-fleetd.sh --check 2>&1 \
|
||||
| sed -E 's#://[^@]*@#://<redacted>@#'
|
||||
| sed -E 's#://[^@]*@#://<redacted>@#g'
|
||||
git rev-parse HEAD
|
||||
```
|
||||
|
||||
@@ -108,8 +113,19 @@ PY
|
||||
|
||||
**What this tier cannot see:** it proves facts only about the Mac daemon at `127.0.0.1:8765`.
|
||||
It cannot show the fleet01 daemon, broker queue depth, or broker consumers. The fleet01 REST service
|
||||
at `10.10.20.13:8765` is not reachable from the Mac, and SSH as `dai.ha@10.10.20.13` is denied.
|
||||
Say this in the report rather than omitting fleet01.
|
||||
at `10.10.20.13:8765` is not reachable from the Mac. Say this in the report rather than omitting
|
||||
fleet01.
|
||||
|
||||
**But fleet01 IS reachable over SSH — checked 2026-08-28.** An older version of this line said SSH
|
||||
was denied. That is true only for the user `dai.ha`. The host alias `fleet01` maps to user `ltms`,
|
||||
and `ssh fleet01` works with key auth:
|
||||
|
||||
```bash
|
||||
ssh -o BatchMode=yes -o ConnectTimeout=6 fleet01 'echo $(id -un)@$(hostname)'
|
||||
```
|
||||
|
||||
So fleet01's daemon PID, uptime, jar and `/healthz` **can** be reported — over SSH, not over REST.
|
||||
Do that rather than writing `not reachable`. `ltms` also has passwordless sudo there.
|
||||
|
||||
## 3. Tier 2 — the shared broker (run when management access exists)
|
||||
|
||||
|
||||
@@ -93,21 +93,20 @@ bind:
|
||||
# intervalSeconds → how often a tick runs (default 30). ENFORCED floor of 15: the code computes
|
||||
# Math.max(15, intervalSeconds), so a lower value is silently raised, not
|
||||
# rejected.
|
||||
# workingSuspectAfterSeconds, paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by
|
||||
# anything — the dormant monitor only consumes intervalSeconds today (CB-573
|
||||
# shipped ahead of the evidence publishers these two knobs are for). Setting
|
||||
# them changes nothing right now, and no minimum is enforced on either, because
|
||||
# nothing reads them to enforce one. They exist so a later build can start
|
||||
# honouring them without another config-shape change.
|
||||
# workingSuspectAfterSeconds → age before a BUSY member is suspected of a stall (default 600).
|
||||
# ENFORCED floor of 300: a lower value is silently raised.
|
||||
# paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by anything. Setting it changes
|
||||
# nothing right now. It exists so a later build can start honouring it without
|
||||
# another config-shape change.
|
||||
# notifications.mode → "webhook" flips what fleet_list REPORTS (healthCoverage: "full" instead
|
||||
# of "detection-only") — it does NOT make fleetd send any webhook call; no
|
||||
# delivery mechanism is implemented yet. Any other value, or omitting the
|
||||
# block, reports "detection-only".
|
||||
# health:
|
||||
# enabled: true
|
||||
# intervalSeconds: 30
|
||||
# workingSuspectAfterSeconds: 600
|
||||
# paneProbeIntervalSeconds: 60
|
||||
# intervalSeconds: 30 # floor 15
|
||||
# workingSuspectAfterSeconds: 600 # floor 300 — how long BUSY with no activity means STALL_SUSPECTED
|
||||
# paneProbeIntervalSeconds: 60 # parsed, but nothing reads it yet — changing it changes nothing
|
||||
# notifications:
|
||||
# mode: disabled
|
||||
|
||||
@@ -115,6 +114,9 @@ bind:
|
||||
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
|
||||
# Optional socket for member panes. Omit this to use herdrSocket for both leads and members.
|
||||
# memberHerdrSocket: /Users/member/.config/herdr/herdr.sock
|
||||
|
||||
# How member sessions are spawned. Define one or more named profiles (backends) under
|
||||
# `profiles`; each key is the profile name (also the ccs profile). A profile says only WHICH
|
||||
# BACKEND — model, CLI adapter, credentials, cost. It says nothing about what a member spawned on
|
||||
|
||||
@@ -7,6 +7,7 @@ import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
import dev.ltms.fleet.herdr.LeadTabScanner;
|
||||
import dev.ltms.fleet.lead.LeadLauncher;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
@@ -149,9 +150,12 @@ public final class Fleetd {
|
||||
: UnixSocketHerdrClient.defaultSocketPath();
|
||||
|
||||
UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect(socket, new com.fasterxml.jackson.databind.ObjectMapper());
|
||||
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
WorkspaceControl spaces = new WorkspaceControl(herdr);
|
||||
UnixSocketHerdrClient memberHerdr = cfg.memberHerdrSocket() != null && !cfg.memberHerdrSocket().isBlank()
|
||||
? UnixSocketHerdrClient.connect(Path.of(cfg.memberHerdrSocket()), new com.fasterxml.jackson.databind.ObjectMapper())
|
||||
: herdr;
|
||||
AtomicReference<Supplier<Map<String, String>>> leadsRef = new AtomicReference<>(Map::of);
|
||||
HerdrRouter router = new HerdrRouter(herdr, memberHerdr,
|
||||
target -> leadsRef.get().get().containsKey(target));
|
||||
// CB-402: one adapter per configured peer kind, fronted by a composite router. A profile's
|
||||
// `kind:` selects its adapter — claude-code (the default) and opencode partition the profile
|
||||
// set — and the composite dispatches each SPI call to the adapter that owns the profile/pane.
|
||||
@@ -169,14 +173,14 @@ public final class Fleetd {
|
||||
// bridge configured with no workers, or opencode-only, still has a well-defined base adapter)
|
||||
// unless opencode is the only kind configured.
|
||||
if (!claudeProfiles.isEmpty() || opencodeProfiles.isEmpty()) {
|
||||
adapters.add(new ClaudeCodeLauncher(agents, spaces, guard,
|
||||
adapters.add(new ClaudeCodeLauncher(router.memberAgents(), router.memberSpaces(), guard,
|
||||
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
() -> config.get().memberCredentials()));
|
||||
}
|
||||
if (!opencodeProfiles.isEmpty()) {
|
||||
adapters.add(new OpenCodeLauncher(agents, spaces,
|
||||
adapters.add(new OpenCodeLauncher(router.memberAgents(), router.memberSpaces(),
|
||||
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||
() -> config.get().fleet(),
|
||||
@@ -271,6 +275,7 @@ public final class Fleetd {
|
||||
// operational cadence, not identity, so there is no correctness reason to give every
|
||||
// lead its own scanner.
|
||||
int scanIntervalSeconds = leaders.values().iterator().next().scanIntervalSeconds();
|
||||
// This must use the lead daemon: scanning member tabs would demote the lead to a worker.
|
||||
leads = new LeadTabScanner(herdr, tabToName, Set.of(),
|
||||
TimeUnit.SECONDS.toNanos(scanIntervalSeconds), System::nanoTime);
|
||||
log.info("lead scan: tabs {} host a lead (rescan every {}s, shared fleet space)",
|
||||
@@ -278,13 +283,14 @@ public final class Fleetd {
|
||||
} else {
|
||||
leads = () -> leadTerminals;
|
||||
}
|
||||
leadsRef.set(leads);
|
||||
|
||||
// CB-558: start any declared lead that is not already running. After the scanner is built,
|
||||
// because both read the same tab labels and the ordering makes that dependency visible; and
|
||||
// only when herdr answered, because the launcher's whole safety property is that it can
|
||||
// count live leads first — it must never guess and risk a second orchestrator.
|
||||
if (herdrUp && !leaders.isEmpty()) {
|
||||
int launched = new LeadLauncher(agents, spaces, cfg).ensureLeads();
|
||||
int launched = new LeadLauncher(router.leadAgents(), router.leadSpaces(), cfg).ensureLeads();
|
||||
if (launched > 0) {
|
||||
log.info("lead auto-launch: {} lead(s) started", launched);
|
||||
}
|
||||
@@ -340,6 +346,7 @@ public final class Fleetd {
|
||||
+ "BACKEND_EXHAUSTED): {}", credentialId,
|
||||
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
|
||||
});
|
||||
AgentControl agents = router.memberAgents();
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink);
|
||||
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
|
||||
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
|
||||
@@ -380,9 +387,10 @@ public final class Fleetd {
|
||||
sessions.onTurnFailed(target);
|
||||
}
|
||||
};
|
||||
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
|
||||
Predicate<String> deliverable = deliverableTo(presence, leads);
|
||||
Injector injector = new Injector(router, turnListener, deliverable,
|
||||
presence::forget);
|
||||
StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS);
|
||||
StatusPoller poller = new StatusPoller(router, injector, Injector.POLL_INTERVAL_MILLIS);
|
||||
poller.start();
|
||||
|
||||
// CB-307: reply inbox. A broker: block selects the AMQP-backed durable adapter; absent (or
|
||||
@@ -419,7 +427,7 @@ public final class Fleetd {
|
||||
// are counted at their single funnel rather than at each of the two caller-facing surfaces.
|
||||
// CB-512: the push loop takes it too, so nudge outcomes (delivered|exhausted) are counted.
|
||||
Metrics metrics = FleetMetrics.create(sessions, replyInbox);
|
||||
var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox,
|
||||
var pushLoop = new ReplyPushLoop(primaryRegistry, router.leadAgents(), replyInbox,
|
||||
pushScheduler, maxReminders, backoffMs, metrics);
|
||||
// CB-551: idle-lead heartbeat. Opt-in; absent `leadHeartbeat:` this is never constructed, so
|
||||
// an upgraded daemon cannot silently start spending subscription on nudging an idle lead.
|
||||
@@ -429,7 +437,7 @@ public final class Fleetd {
|
||||
Thread.ofVirtual().name("bridge-heartbeat-").unstarted(r));
|
||||
if (cfg.leadHeartbeat() != null) {
|
||||
var hb = cfg.leadHeartbeat();
|
||||
heartbeat = new LeadHeartbeatLoop(primaryRegistry, agents, replyInbox, sessions::roster,
|
||||
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
|
||||
pushLoop, heartbeatScheduler, System::nanoTime,
|
||||
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
|
||||
metrics);
|
||||
@@ -438,7 +446,7 @@ public final class Fleetd {
|
||||
heartbeat = null;
|
||||
heartbeatScheduler.shutdownNow();
|
||||
}
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
|
||||
MessageService messages = new MessageService(router, injector, rendezvous, replyInbox,
|
||||
pushLoop, metrics);
|
||||
|
||||
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
|
||||
@@ -449,7 +457,8 @@ public final class Fleetd {
|
||||
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it,
|
||||
// through the same idempotent target-wide operation CB-516 already uses on release.
|
||||
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
|
||||
System::nanoTime, cfg.health().intervalOrDefault(), messages::abandon);
|
||||
System::nanoTime, cfg.health().intervalOrDefault(),
|
||||
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
|
||||
String coverage = FleetHealthMonitor.coverage(true,
|
||||
cfg.health().notifications() != null && cfg.health().notifications().configured());
|
||||
if ("detection-only".equals(coverage)) {
|
||||
@@ -489,8 +498,11 @@ public final class Fleetd {
|
||||
|
||||
// MCP server face (CB-105): fleet_send/fleet_reply/fleet_status, mounted at /mcp.
|
||||
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
|
||||
// CB-185: a caller's pane can live on either daemon (a lead's on the lead daemon, a
|
||||
// member's on the member daemon) — search both, lead first. Collapses to one scan when
|
||||
// memberHerdrSocket is unset (herdr == memberHerdr).
|
||||
ConnectionIdentity identity = new ConnectionIdentity(
|
||||
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
||||
new PaneLocator(herdr, memberHerdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
||||
|
||||
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
|
||||
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
|
||||
@@ -536,7 +548,7 @@ public final class Fleetd {
|
||||
if (leadMailbox != null) {
|
||||
var leadCoordScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
||||
Thread.ofVirtual().name("bridge-leadcoord-").unstarted(r));
|
||||
leadCoordLoop = new LeadCoordLoop(leadMailbox, agents, leads, leadCoordScheduler,
|
||||
leadCoordLoop = new LeadCoordLoop(leadMailbox, router.leadAgents(), leads, leadCoordScheduler,
|
||||
LEAD_COORD_INTERVAL_MS);
|
||||
leadCoordLoop.start();
|
||||
leadCoordSchedulerRef = leadCoordScheduler;
|
||||
@@ -587,11 +599,13 @@ public final class Fleetd {
|
||||
log.debug("lead mailbox close: {}", e.toString());
|
||||
}
|
||||
}
|
||||
herdr.close();
|
||||
router.close();
|
||||
}));
|
||||
|
||||
Javalin app = new FleetApp(herdr, workers, sessions, messages, presence, mcp.servlet(),
|
||||
callers, metrics).build();
|
||||
// CB-185: give FleetApp both daemons — /healthz must require both to answer and
|
||||
// GET /sessions must merge across both, or a down/unpolled member daemon is invisible.
|
||||
Javalin app = new FleetApp(herdr, memberHerdr, workers, sessions, messages, presence, mcp.servlet(),
|
||||
callers, metrics, deliverable).build();
|
||||
app.start(cfg.bind().host(), cfg.bind().port());
|
||||
log.info("fleetd listening on {}:{}, herdr socket {}",
|
||||
cfg.bind().host(), cfg.bind().port(), socket);
|
||||
|
||||
@@ -71,7 +71,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
|
||||
/** Keys that cannot change under a running daemon — see the class doc. */
|
||||
private static final Set<String> COLD_KEYS =
|
||||
Set.of("bind", "herdrSocket", "broker", "auth");
|
||||
Set.of("bind", "herdrSocket", "memberHerdrSocket", "broker", "auth");
|
||||
|
||||
private final Path path;
|
||||
private final AtomicReference<FleetConfig> current;
|
||||
@@ -190,6 +190,9 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
if (!Objects.equals(old.herdrSocket(), fresh.herdrSocket())) {
|
||||
changed.add("herdrSocket");
|
||||
}
|
||||
if (!Objects.equals(old.memberHerdrSocket(), fresh.memberHerdrSocket())) {
|
||||
changed.add("memberHerdrSocket");
|
||||
}
|
||||
if (!Objects.equals(old.broker(), fresh.broker())) {
|
||||
changed.add("broker");
|
||||
}
|
||||
|
||||
@@ -32,7 +32,8 @@ import java.util.Set;
|
||||
* silently dropping a whole block is indistinguishable from honouring it.
|
||||
*
|
||||
* @param bind REST/MCP listen host:port
|
||||
* @param herdrSocket path to herdr's Unix socket ({@code null} → client default)
|
||||
* @param herdrSocket path to the lead herdr Unix socket ({@code null} → client default)
|
||||
* @param memberHerdrSocket optional member herdr Unix socket ({@code null}/blank → lead socket)
|
||||
* @param profiles named backend profiles, keyed by profile name (multi-backend fleet). A
|
||||
* profile answers <em>which backend</em> — model, CLI adapter, credentials,
|
||||
* cost. It says nothing about what the member spawned on it is for; that is
|
||||
@@ -80,6 +81,7 @@ import java.util.Set;
|
||||
public record FleetConfig(
|
||||
Bind bind,
|
||||
String herdrSocket,
|
||||
String memberHerdrSocket,
|
||||
Map<String, Profile> profiles,
|
||||
Guard guard,
|
||||
String worktreeRoot,
|
||||
@@ -105,7 +107,7 @@ public record FleetConfig(
|
||||
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||
ConfigReload configReload, Integer quarantineCooldownSeconds,
|
||||
MemberCredentials memberCredentials) {
|
||||
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, null);
|
||||
}
|
||||
@@ -116,9 +118,9 @@ public record FleetConfig(
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||
ConfigReload configReload, Integer quarantineCooldownSeconds) {
|
||||
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, quarantineCooldownSeconds, null);
|
||||
configReload, quarantineCooldownSeconds, null, null);
|
||||
}
|
||||
|
||||
/** Default cooldown (CB-578 stage B) when {@code quarantineCooldownSeconds} is absent/non-positive. */
|
||||
@@ -129,8 +131,8 @@ public record FleetConfig(
|
||||
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, String placement, Auth auth) {
|
||||
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null, null);
|
||||
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null, null, null, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the optional {@code health:} block was added. */
|
||||
@@ -138,8 +140,8 @@ public record FleetConfig(
|
||||
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, String placement, Auth auth, ConfigReload configReload) {
|
||||
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload, null);
|
||||
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload, null, null, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the CB-578 stage B {@code quarantineCooldownSeconds} field was added. */
|
||||
@@ -148,9 +150,9 @@ public record FleetConfig(
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||
ConfigReload configReload) {
|
||||
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
this(bind, herdrSocket, null, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, null);
|
||||
configReload, null, null, null);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -642,6 +644,9 @@ public record FleetConfig(
|
||||
Integer paneProbeIntervalSeconds, Notifications notifications) {
|
||||
public boolean isEnabled() { return Boolean.TRUE.equals(enabled); }
|
||||
public int intervalOrDefault() { return Math.max(15, intervalSeconds == null ? 30 : intervalSeconds); }
|
||||
public int workingSuspectAfterOrDefault() {
|
||||
return Math.max(300, workingSuspectAfterSeconds == null ? 600 : workingSuspectAfterSeconds);
|
||||
}
|
||||
public record Notifications(String mode) {
|
||||
public boolean configured() { return "webhook".equalsIgnoreCase(mode); }
|
||||
}
|
||||
@@ -1317,7 +1322,7 @@ public record FleetConfig(
|
||||
* {@code fleetd.yaml} itself is gitignored.
|
||||
*/
|
||||
static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
|
||||
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
|
||||
"bind", "herdrSocket", "memberHerdrSocket", "profiles", "guard", "worktreeRoot",
|
||||
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
|
||||
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
|
||||
"memberCredentials", "coordinator");
|
||||
@@ -1938,7 +1943,7 @@ public record FleetConfig(
|
||||
: new MemberCredentials(null, List.of(), List.of());
|
||||
// coordinator is left as-is, like broker/primary above: null keeps no LeadMailbox opened,
|
||||
// and this ticket's Coordinator is config-only anyway (nothing yet reads it at startup).
|
||||
return new FleetConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
|
||||
return new FleetConfig(b, herdrSocket, memberHerdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
|
||||
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
|
||||
quarantineCooldown, mc, coordinator);
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package dev.ltms.fleet.health;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import org.slf4j.Logger;
|
||||
@@ -25,6 +26,8 @@ public final class FleetHealthMonitor {
|
||||
|
||||
/** Bounded attempts to run {@link #failTarget} for one transition. Never retried tick-to-tick (CB-580). */
|
||||
static final int MAX_FAIL_TARGET_ATTEMPTS = 3;
|
||||
// CB-641: Match the injector's 60s readiness gate so health allows a full first boot.
|
||||
static final long READINESS_GRACE_NANOS = TimeUnit.SECONDS.toNanos(60);
|
||||
|
||||
private final AgentControl agents;
|
||||
private final Supplier<List<MemberSession>> roster;
|
||||
@@ -32,12 +35,30 @@ public final class FleetHealthMonitor {
|
||||
private final ScheduledExecutorService scheduler;
|
||||
private final LongSupplier clock;
|
||||
private final long intervalSeconds;
|
||||
private final long workingSuspectAfterNanos;
|
||||
private final BiConsumer<String, String> failTarget;
|
||||
private final Map<String, HealthPrior> priors = new HashMap<>();
|
||||
private final Map<String, HealthState> states = new HashMap<>();
|
||||
/**
|
||||
* CB-643: consecutive ticks on which a target looked like an orphaned delegation. The fact
|
||||
* {@link MessageService#hasOrphanedDelegation} reports is a true snapshot, but it can read true
|
||||
* for one tick during an ordinary race — an async ticket exists before its virtual thread has
|
||||
* reached {@code rendezvous.open()}, so for that instant nothing is accepted or queued behind
|
||||
* it. {@code decide} maps the field straight to {@code DELEGATION_ORPHANED} with no cross-tick
|
||||
* smoothing of its own, so a single racy read would log a fault that clears on the next tick.
|
||||
* Requiring two consecutive observations costs one interval of latency on a real orphan and
|
||||
* removes that false positive entirely.
|
||||
*/
|
||||
private final Map<String, Integer> orphanStreaks = new HashMap<>();
|
||||
|
||||
// These facts need the evidence publishers introduced by later M4 units. They are not negatives.
|
||||
private static final boolean NOT_YET_OBSERVED = false;
|
||||
/** How many consecutive ticks a target must look orphaned before health reports it (CB-643). */
|
||||
static final int ORPHAN_CONFIRM_TICKS = 2;
|
||||
|
||||
// CB-643: every HealthSnapshot field now carries real evidence. The NOT_YET_OBSERVED placeholder
|
||||
// that stood in for 7 of the 12 is gone, and with it the reason 8 of the 9 fault states were
|
||||
// unreachable — GONE and NEVER_READY included, which is what kept CB-580's failTarget from ever
|
||||
// firing. Do not reintroduce a constant here: a field with no publisher is a dead state, and the
|
||||
// tests pass either way, so nothing else will tell you.
|
||||
|
||||
/**
|
||||
* @param failTarget CB-568's idempotent target-wide failure operation (e.g. {@code messages::abandon}),
|
||||
@@ -48,13 +69,14 @@ public final class FleetHealthMonitor {
|
||||
*/
|
||||
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
|
||||
BiConsumer<String, String> failTarget) {
|
||||
long workingSuspectAfterSeconds, BiConsumer<String, String> failTarget) {
|
||||
this.agents = agents;
|
||||
this.roster = roster;
|
||||
this.messages = messages;
|
||||
this.scheduler = scheduler;
|
||||
this.clock = clock;
|
||||
this.intervalSeconds = intervalSeconds;
|
||||
this.workingSuspectAfterNanos = TimeUnit.SECONDS.toNanos(workingSuspectAfterSeconds);
|
||||
this.failTarget = Objects.requireNonNull(failTarget, "failTarget");
|
||||
}
|
||||
|
||||
@@ -69,28 +91,50 @@ public final class FleetHealthMonitor {
|
||||
// Package-private so tests can run one tick without waiting.
|
||||
void tick() {
|
||||
try {
|
||||
List<Agent> agentsNow = agents.list(); // Exactly one list call for this complete observation.
|
||||
List<MemberSession> rosterNow = roster.get(); // One in-memory roster snapshot for this tick.
|
||||
List<Agent> agentsNow;
|
||||
boolean controlLinkDown = false;
|
||||
try {
|
||||
agentsNow = agents.list(); // Exactly one list call for this complete observation.
|
||||
} catch (HerdrException error) {
|
||||
agentsNow = List.of();
|
||||
controlLinkDown = true;
|
||||
log.warn("fleet health control link unavailable; classifying roster", error);
|
||||
}
|
||||
Map<String, Agent> live = new HashMap<>();
|
||||
for (Agent agent : agentsNow) live.put(agent.terminalId(), agent);
|
||||
HashSet<String> current = new HashSet<>();
|
||||
long nowNanos = clock.getAsLong();
|
||||
for (MemberSession session : rosterNow) {
|
||||
current.add(session.terminalId());
|
||||
Agent agent = live.get(session.terminalId());
|
||||
AgentStatus status = agent == null ? AgentStatus.UNKNOWN : agent.status();
|
||||
boolean accepted = messages.hasAcceptedDelivery(session.terminalId());
|
||||
HealthSnapshot snapshot = new HealthSnapshot(session.state(), status, accepted, NOT_YET_OBSERVED,
|
||||
messages.hasInboxMessage(session.terminalId()), agent != null, NOT_YET_OBSERVED,
|
||||
NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED);
|
||||
boolean present = agent != null;
|
||||
boolean targetNotFound = !controlLinkDown && !present
|
||||
&& session.state() != MemberSession.State.SPAWNING;
|
||||
boolean readinessGraceElapsed = nowNanos - session.spawnedAtNanos() >= READINESS_GRACE_NANOS;
|
||||
boolean stalled = session.state() == MemberSession.State.BUSY
|
||||
&& nowNanos - session.lastActivityAtNanos() >= workingSuspectAfterNanos;
|
||||
// CB-643: the three message-layer facts CB-640 published. Read them here rather than
|
||||
// leaving them false — that constant is what made 8 of the 9 fault states dead.
|
||||
boolean queuedDelivery = messages.hasQueuedDelivery(session.terminalId());
|
||||
boolean replyStranded = messages.hasStrandedReply(session.terminalId());
|
||||
boolean orphanedDelegation = confirmOrphan(session.terminalId(),
|
||||
messages.hasOrphanedDelegation(session.terminalId()));
|
||||
HealthSnapshot snapshot = new HealthSnapshot(session.state(), status, accepted, queuedDelivery,
|
||||
messages.hasInboxMessage(session.terminalId()), present, targetNotFound, controlLinkDown,
|
||||
readinessGraceElapsed, orphanedDelegation, replyStranded, stalled);
|
||||
HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE),
|
||||
clock.getAsLong());
|
||||
nowNanos);
|
||||
priors.put(session.terminalId(), decision.prior());
|
||||
reportTransition(session.terminalId(), decision.state());
|
||||
}
|
||||
priors.keySet().retainAll(current);
|
||||
states.keySet().retainAll(current);
|
||||
orphanStreaks.keySet().retainAll(current);
|
||||
} catch (Throwable error) {
|
||||
// A list failure is health evidence, and must never kill the monitor's only scheduler task.
|
||||
// Any unclassified collection failure must never kill the monitor's only scheduler task.
|
||||
log.warn("fleet health collection failed; will retry next tick", error);
|
||||
} finally {
|
||||
if (!scheduler.isShutdown()) {
|
||||
@@ -99,6 +143,20 @@ public final class FleetHealthMonitor {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Debounce {@link MessageService#hasOrphanedDelegation} across ticks (CB-643). Returns true only
|
||||
* once {@code observed} has held for {@link #ORPHAN_CONFIRM_TICKS} consecutive ticks; a single
|
||||
* false reading resets the streak, so a transient race never reaches the classifier.
|
||||
*/
|
||||
private boolean confirmOrphan(String target, boolean observed) {
|
||||
if (!observed) {
|
||||
orphanStreaks.remove(target);
|
||||
return false;
|
||||
}
|
||||
int streak = orphanStreaks.merge(target, 1, Integer::sum);
|
||||
return streak >= ORPHAN_CONFIRM_TICKS;
|
||||
}
|
||||
|
||||
void reportTransition(String target, HealthState next) {
|
||||
HealthState previous = states.put(target, next);
|
||||
if (previous == next) return;
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
package dev.ltms.fleet.herdr;
|
||||
|
||||
import java.util.Objects;
|
||||
import java.util.function.Predicate;
|
||||
|
||||
/** Routes lead operations and member operations to their owning herdr daemon. */
|
||||
public final class HerdrRouter implements AutoCloseable {
|
||||
private final HerdrClient lead;
|
||||
private final HerdrClient member;
|
||||
private final AgentControl leadAgents;
|
||||
private final AgentControl memberAgents;
|
||||
private final WorkspaceControl leadSpaces;
|
||||
private final WorkspaceControl memberSpaces;
|
||||
private final Predicate<String> isLead;
|
||||
|
||||
public HerdrRouter(HerdrClient lead, HerdrClient member, Predicate<String> isLead) {
|
||||
this.lead = Objects.requireNonNull(lead, "lead");
|
||||
this.member = member != null ? member : lead;
|
||||
this.isLead = Objects.requireNonNull(isLead, "isLead");
|
||||
leadAgents = new AgentControl(this.lead);
|
||||
memberAgents = this.member == this.lead ? leadAgents : new AgentControl(this.member);
|
||||
leadSpaces = new WorkspaceControl(this.lead);
|
||||
memberSpaces = this.member == this.lead ? leadSpaces : new WorkspaceControl(this.member);
|
||||
}
|
||||
|
||||
public AgentControl leadAgents() { return leadAgents; }
|
||||
public WorkspaceControl leadSpaces() { return leadSpaces; }
|
||||
public AgentControl memberAgents() { return memberAgents; }
|
||||
public WorkspaceControl memberSpaces() { return memberSpaces; }
|
||||
public AgentControl agentsFor(String targetId) { return isLead.test(targetId) ? leadAgents : memberAgents; }
|
||||
|
||||
HerdrClient leadClient() { return lead; }
|
||||
HerdrClient memberClient() { return member; }
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
lead.close();
|
||||
if (member != lead) member.close();
|
||||
}
|
||||
}
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.fleet.herdr;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
@@ -13,33 +14,61 @@ import java.util.Map;
|
||||
* <p>herdr owns the PID→pane truth: {@code pane.process_info} reports each pane's {@code shell_pid}
|
||||
* and foreground process PIDs. This scans agent panes; a spawn-time {@code pid→terminal} cache is
|
||||
* the obvious optimization once wired into {@code ClaudeCodeLauncher}.
|
||||
*
|
||||
* <p>CB-185 split the fleet across two herdr daemons — lead operations on one, members on the
|
||||
* other ({@code memberHerdrSocket}). A caller's pane can live on <em>either</em> daemon (a lead's
|
||||
* MCP connection resolves against the lead daemon; a member's against the member daemon), so this
|
||||
* must be able to search more than one client. {@link #PaneLocator(HerdrClient, HerdrClient)}
|
||||
* searches the lead client first, then the member client, and collapses to a single scan when the
|
||||
* two are the same object (the historical single-daemon deployment).
|
||||
*/
|
||||
public final class PaneLocator {
|
||||
|
||||
private final HerdrClient herdr;
|
||||
private final List<HerdrClient> herdrs;
|
||||
|
||||
/** Search only this client — the single-daemon deployment. */
|
||||
public PaneLocator(HerdrClient herdr) {
|
||||
this.herdr = herdr;
|
||||
this.herdrs = List.of(herdr);
|
||||
}
|
||||
|
||||
/**
|
||||
* Search {@code lead} first, then {@code member} — the two-daemon deployment (CB-185). When
|
||||
* the caller passes the same client for both (no {@code memberHerdrSocket} configured), this
|
||||
* collapses to one client and one scan, exactly {@link #PaneLocator(HerdrClient)}'s behaviour.
|
||||
*/
|
||||
public PaneLocator(HerdrClient lead, HerdrClient member) {
|
||||
this.herdrs = lead == member ? List.of(lead) : List.of(lead, member);
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@code terminal_id} of the agent pane whose process tree contains {@code pid}, or
|
||||
* {@code null} if no agent pane owns it (e.g. the caller is the primary, or off-host).
|
||||
* {@code null} if no agent pane on any searched daemon owns it (e.g. the caller is the
|
||||
* primary, or off-host).
|
||||
*/
|
||||
public String terminalForPid(long pid) {
|
||||
if (pid <= 0) {
|
||||
return null;
|
||||
}
|
||||
for (HerdrClient herdr : herdrs) {
|
||||
String terminal = terminalForPid(herdr, pid);
|
||||
if (terminal != null) {
|
||||
return terminal;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
private static String terminalForPid(HerdrClient herdr, long pid) {
|
||||
for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) {
|
||||
String paneId = pane.path("pane_id").asText(null);
|
||||
if (paneId != null && paneOwnsPid(paneId, pid)) {
|
||||
if (paneId != null && paneOwnsPid(herdr, paneId, pid)) {
|
||||
return pane.path("terminal_id").asText(null);
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
private boolean paneOwnsPid(String paneId, long pid) {
|
||||
private static boolean paneOwnsPid(HerdrClient herdr, String paneId, long pid) {
|
||||
JsonNode info;
|
||||
try {
|
||||
info = herdr.call("pane.process_info", Map.of("pane_id", paneId)).path("process_info");
|
||||
|
||||
@@ -6,12 +6,14 @@ import dev.ltms.fleet.msg.TurnToken;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
/**
|
||||
@@ -61,6 +63,17 @@ public final class CompletionResolver implements TurnListener {
|
||||
/** Cap the scraped tail so a long transcript can't return an unbounded blob. */
|
||||
static final int MAX_SCRAPE_CHARS = 4000;
|
||||
|
||||
/**
|
||||
* fleetd#164: the floor below which a {@code BUSY -> DONE} transition cannot be real work. A
|
||||
* backend that rejects a turn outright (e.g. an HTTP 400 from the model, before the worker read
|
||||
* a single file or produced a token) drives the exact same confirmed {@code working -> idle}
|
||||
* transition a genuine completion does — just in about a second instead of the many seconds a
|
||||
* real turn costs. {@link #onTurnComplete} cannot tell those two cases apart from the transition
|
||||
* alone, so a turn that settles inside this floor is treated as a crash signature and resolved
|
||||
* as a failure, never as a (possibly empty) success.
|
||||
*/
|
||||
public static final long MIN_TURN_NANOS = Duration.ofSeconds(2).toNanos();
|
||||
|
||||
private static final String CLIPPED_PANE_TAIL_MARKER =
|
||||
"[Pane tail clipped: member did not call fleet_reply.]";
|
||||
|
||||
@@ -68,6 +81,7 @@ public final class CompletionResolver implements TurnListener {
|
||||
private final Rendezvous rendezvous;
|
||||
private final ExhaustedPatternLookup exhaustedPatterns;
|
||||
private final ExhaustionSink exhaustionSink;
|
||||
private final LongSupplier nowNanos;
|
||||
|
||||
/**
|
||||
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its
|
||||
@@ -79,8 +93,21 @@ public final class CompletionResolver implements TurnListener {
|
||||
* reference: a completion scrape equal to it means the worker produced no new output (the previous
|
||||
* turn's wind-down sampled as this boundary), so it is suppressed. Overwritten on each delivery;
|
||||
* cleared when the turn resolves. Package-private so tests can capture and replay a specific turn.
|
||||
*
|
||||
* <p>{@code deliveredAtNanos} (fleetd#164) is the {@link #nowNanos} reading taken at delivery —
|
||||
* the other half of the {@link #MIN_TURN_NANOS} floor check, compared against a fresh reading at
|
||||
* resolution time.
|
||||
*/
|
||||
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline) {
|
||||
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos) {
|
||||
|
||||
/**
|
||||
* Convenience for tests exercising scrape/suppression logic that don't care about turn
|
||||
* timing: back-dates the delivery far enough that {@link #MIN_TURN_NANOS} can never fire.
|
||||
* Not used by production code — {@link #captureBaseline} always records a real reading.
|
||||
*/
|
||||
InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline) {
|
||||
this(waiter, baseline, Long.MIN_VALUE / 2);
|
||||
}
|
||||
}
|
||||
|
||||
private final ConcurrentHashMap<String, InFlight> inFlight = new ConcurrentHashMap<>();
|
||||
@@ -97,10 +124,25 @@ public final class CompletionResolver implements TurnListener {
|
||||
*/
|
||||
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
|
||||
ExhaustionSink exhaustionSink) {
|
||||
this(agents, rendezvous, exhaustedPatterns, exhaustionSink, System::nanoTime);
|
||||
}
|
||||
|
||||
/**
|
||||
* Test constructor with an injectable clock (fleetd#164), matching the {@code LongSupplier}
|
||||
* pattern {@link dev.ltms.fleet.session.SessionManager} and {@link dev.ltms.fleet.msg.MessageService}
|
||||
* already use: lets a test place a turn's delivery and its resolution at an exact, controllable
|
||||
* distance apart around the {@link #MIN_TURN_NANOS} floor, without a real sleep. Public (rather
|
||||
* than package-private like those two) because callers that wire a full {@code MessageService}
|
||||
* fixture — e.g. {@code MessageServiceTest} — construct this resolver directly from another
|
||||
* package.
|
||||
*/
|
||||
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
|
||||
ExhaustionSink exhaustionSink, LongSupplier nowNanos) {
|
||||
this.agents = agents;
|
||||
this.rendezvous = rendezvous;
|
||||
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
|
||||
this.exhaustionSink = Objects.requireNonNull(exhaustionSink, "exhaustionSink");
|
||||
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -130,7 +172,7 @@ public final class CompletionResolver implements TurnListener {
|
||||
baseline = null; // fail open: no baseline ⇒ no suppression
|
||||
log.debug("delivery baseline for {} failed: {}", target, e.getMessage());
|
||||
}
|
||||
inFlight.put(target, new InFlight(waiter, baseline));
|
||||
inFlight.put(target, new InFlight(waiter, baseline, nowNanos.getAsLong()));
|
||||
}
|
||||
|
||||
/** The turn currently baselined for {@code target}, or {@code null} — a test hook for the captureBaseline path. */
|
||||
@@ -176,6 +218,15 @@ public final class CompletionResolver implements TurnListener {
|
||||
inFlight.remove(target, turn);
|
||||
return;
|
||||
}
|
||||
// fleetd#164: a BUSY -> DONE transition inside the floor cannot be real work — it's a crash
|
||||
// signature (e.g. a backend HTTP 400 before the worker did anything), not a fast answer. Fail
|
||||
// it before spending a scrape on the ordinary path; the reason still carries whatever is on
|
||||
// screen, since that is usually the backend's own error.
|
||||
long elapsedNanos = nowNanos.getAsLong() - turn.deliveredAtNanos();
|
||||
if (elapsedNanos < MIN_TURN_NANOS) {
|
||||
fail(target, turn, tooFastReason(target, elapsedNanos));
|
||||
return;
|
||||
}
|
||||
String tail;
|
||||
String assistantBlock = null;
|
||||
int originalLength = 0;
|
||||
@@ -187,20 +238,25 @@ public final class CompletionResolver implements TurnListener {
|
||||
clipped = originalLength > MAX_SCRAPE_CHARS;
|
||||
tail = clip(assistantBlock);
|
||||
} catch (RuntimeException e) {
|
||||
// The worker finished but we couldn't read its screen — still resolve the send so the
|
||||
// caller unblocks; an empty tail beats hanging until the caller's timeout.
|
||||
log.warn("completion scrape for {} failed; resolving with an empty tail: {}",
|
||||
target, e.getMessage());
|
||||
log.warn("completion scrape for {} failed: {}", target, e.getMessage());
|
||||
tail = "";
|
||||
scrapeFailed = true;
|
||||
}
|
||||
// fleetd#164: a scrape nobody could read, and a scrape that read cleanly but produced nothing,
|
||||
// both used to resolve the send as a SUCCESS carrying "" — indistinguishable from a worker that
|
||||
// genuinely finished with nothing to say. That is the defect: fail loudly instead, naming the
|
||||
// member, so a caller (including a lead deciding whether to delegate again) can tell a lost
|
||||
// turn from a real empty answer.
|
||||
if (scrapeFailed || tail.isEmpty()) {
|
||||
fail(target, turn, emptyScrapeReason(target, scrapeFailed));
|
||||
return;
|
||||
}
|
||||
// Misattribution guard (CB-115): if the scrape is byte-identical to the pane content at
|
||||
// delivery, this turn produced no new output — the boundary belongs to the previous turn's
|
||||
// wind-down (common on rapid back-to-back sends). Suppress rather than resolve the send with
|
||||
// a stale answer; the real fleet_reply (or a later genuine completion) resolves it instead.
|
||||
// A scrape that failed to read is exempt — an empty tail there is "couldn't see", not "no change".
|
||||
String baseline = turn.baseline();
|
||||
if (!scrapeFailed && baseline != null && baseline.equals(tail)) {
|
||||
if (baseline != null && baseline.equals(tail)) {
|
||||
log.debug("suppressing misattributed completion for {} (no output change since delivery)",
|
||||
target);
|
||||
return; // keep the in-flight record: a later genuine completion still needs it
|
||||
@@ -208,21 +264,19 @@ public final class CompletionResolver implements TurnListener {
|
||||
// CB-578 stage A: a turn that ended with no fleet_reply AND whose scrape matches the
|
||||
// backend's configured usage-limit pattern is a refusal, not an answer. Classify it as
|
||||
// BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply.
|
||||
if (!scrapeFailed) {
|
||||
Pattern exhausted = exhaustedPatterns.patternFor(target);
|
||||
String matchedLine = exhausted == null ? null : firstMatchingLine(assistantBlock, exhausted);
|
||||
if (matchedLine != null) {
|
||||
String reason = "backend exhausted (usage limit): " + matchedLine;
|
||||
if (rendezvous.resolveExhausted(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("completion for {} classified BACKEND_EXHAUSTED (no fleet_reply; scrape "
|
||||
+ "matched the profile's exhausted pattern): {}", target, reason);
|
||||
// CB-578 stage B: only on the resolution that actually won the race — a late
|
||||
// duplicate must never quarantine a credential twice for one refusal.
|
||||
exhaustionSink.onExhausted(target, reason);
|
||||
}
|
||||
return;
|
||||
Pattern exhausted = exhaustedPatterns.patternFor(target);
|
||||
String matchedLine = exhausted == null ? null : firstMatchingLine(assistantBlock, exhausted);
|
||||
if (matchedLine != null) {
|
||||
String reason = "backend exhausted (usage limit): " + matchedLine;
|
||||
if (rendezvous.resolveExhausted(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("completion for {} classified BACKEND_EXHAUSTED (no fleet_reply; scrape "
|
||||
+ "matched the profile's exhausted pattern): {}", target, reason);
|
||||
// CB-578 stage B: only on the resolution that actually won the race — a late
|
||||
// duplicate must never quarantine a credential twice for one refusal.
|
||||
exhaustionSink.onExhausted(target, reason);
|
||||
}
|
||||
return;
|
||||
}
|
||||
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
|
||||
if (rendezvous.resolveCompletion(waiter, completion)) {
|
||||
@@ -272,6 +326,36 @@ public final class CompletionResolver implements TurnListener {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd#164: the failure reason for a turn that settled inside {@link #MIN_TURN_NANOS} — names
|
||||
* the member and both timings, and appends whatever the pane shows (usually the backend's own
|
||||
* error) so the caller sees the cause, not just "it failed".
|
||||
*/
|
||||
private String tooFastReason(String target, long elapsedNanos) {
|
||||
String scrape;
|
||||
try {
|
||||
scrape = clip(agents.read(target, SCRAPE_SOURCE));
|
||||
} catch (RuntimeException e) {
|
||||
scrape = "";
|
||||
}
|
||||
String reason = String.format(
|
||||
"member %s went BUSY -> DONE in %dms (floor %dms) — too fast to be real work, most "
|
||||
+ "likely a backend error before any work started",
|
||||
target, elapsedNanos / 1_000_000, MIN_TURN_NANOS / 1_000_000);
|
||||
return scrape.isBlank() ? reason : reason + ": " + scrape;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd#164: the failure reason for a scrape that produced zero characters — names the member
|
||||
* and says plainly that the turn produced nothing, so a caller (a lead deciding whether to
|
||||
* delegate again included) never mistakes a lost turn for a genuinely empty reply.
|
||||
*/
|
||||
private static String emptyScrapeReason(String target, boolean scrapeFailed) {
|
||||
return "member " + target + " turn completed with an empty scrape (0 chars) — "
|
||||
+ (scrapeFailed ? "its pane could not be read; " : "")
|
||||
+ "treating as a lost turn, not a real answer";
|
||||
}
|
||||
|
||||
/**
|
||||
* The first line of {@code text} matching {@code pattern}, stripped — the CB-578 stage A
|
||||
* evidence carried in a {@code BACKEND_EXHAUSTED} reason so the operator sees the real refusal
|
||||
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.fleet.inject;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
@@ -89,6 +90,7 @@ public final class Injector {
|
||||
public static final long POLL_INTERVAL_MILLIS = 250;
|
||||
|
||||
private final AgentControl agents;
|
||||
private final HerdrRouter router;
|
||||
private final TurnListener turnListener;
|
||||
private final Predicate<String> ready; // CB-113: a target is deliverable only when available
|
||||
private final Consumer<String> forget; // CB-114: clear a gone worker's readiness/presence
|
||||
@@ -124,11 +126,25 @@ public final class Injector {
|
||||
public Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
|
||||
Consumer<String> forget) {
|
||||
this.agents = agents;
|
||||
this.router = null;
|
||||
this.turnListener = turnListener;
|
||||
this.ready = ready;
|
||||
this.forget = forget;
|
||||
}
|
||||
|
||||
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
|
||||
Consumer<String> forget) {
|
||||
this.agents = null;
|
||||
this.router = router;
|
||||
this.turnListener = turnListener;
|
||||
this.ready = ready;
|
||||
this.forget = forget;
|
||||
}
|
||||
|
||||
private AgentControl agentsFor(String target) {
|
||||
return router != null ? router.agentsFor(target) : agents;
|
||||
}
|
||||
|
||||
/** A pending message and the future that completes when it has been delivered. */
|
||||
private record Pending(String text, TurnToken token, CompletableFuture<Void> delivered) {
|
||||
}
|
||||
@@ -253,7 +269,7 @@ public final class Injector {
|
||||
if (p != null && ready.test(target)) {
|
||||
t.notReadySincePoll = 0;
|
||||
try {
|
||||
agents.send(target, p.text());
|
||||
agentsFor(target).send(target, p.text());
|
||||
t.queue.poll();
|
||||
t.awaitingPickup = true;
|
||||
t.awaitingCompletion = true;
|
||||
@@ -313,7 +329,7 @@ public final class Injector {
|
||||
// thread while it holds the target lock.
|
||||
if (resubmit) {
|
||||
try {
|
||||
agents.submit(target); // nudge a raced Enter so the pending paste submits
|
||||
agentsFor(target).submit(target); // nudge a raced Enter so the pending paste submits
|
||||
} catch (RuntimeException e) {
|
||||
log.debug("resubmit to {} failed (will retry next poll): {}", target, e.getMessage());
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package dev.ltms.fleet.inject;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -22,6 +23,7 @@ public final class StatusPoller {
|
||||
private static final Logger log = LoggerFactory.getLogger(StatusPoller.class);
|
||||
|
||||
private final AgentControl agents;
|
||||
private final HerdrRouter router;
|
||||
private final Injector injector;
|
||||
private final StatusRefiner refiner;
|
||||
private final long intervalMillis;
|
||||
@@ -35,11 +37,24 @@ public final class StatusPoller {
|
||||
public StatusPoller(AgentControl agents, Injector injector, StatusRefiner refiner,
|
||||
long intervalMillis) {
|
||||
this.agents = agents;
|
||||
this.router = null;
|
||||
this.injector = injector;
|
||||
this.refiner = refiner;
|
||||
this.intervalMillis = intervalMillis;
|
||||
}
|
||||
|
||||
public StatusPoller(HerdrRouter router, Injector injector, long intervalMillis) {
|
||||
this.agents = null;
|
||||
this.router = router;
|
||||
this.injector = injector;
|
||||
// CB-185: this refiner's own AgentControl (member) is only a default for the legacy 2-arg
|
||||
// refine() overload — the loop below always calls the 3-arg refine(target, raw, control)
|
||||
// with the per-target control from router.agentsFor(target), so a lead target is refined
|
||||
// against the LEAD daemon even though this field points at the member one.
|
||||
this.refiner = new StatusRefiner(router.memberAgents());
|
||||
this.intervalMillis = intervalMillis;
|
||||
}
|
||||
|
||||
/** Start the polling loop on a virtual thread. Idempotent. */
|
||||
public synchronized void start() {
|
||||
if (running) return;
|
||||
@@ -56,7 +71,11 @@ public final class StatusPoller {
|
||||
try {
|
||||
// herdr's agent_status can misreport a settled worker as `unknown`; refine it
|
||||
// against the pane content before it drives delivery/completion (CB-115).
|
||||
AgentStatus status = refiner.refine(target, agents.status(target));
|
||||
// CB-185: refine THROUGH the same control the raw status came from — a router
|
||||
// splits lead/member targets across two herdr daemons, and reading a lead's pane
|
||||
// through the (fixed) member refiner never finds it, wedging that lead at UNKNOWN.
|
||||
AgentControl control = router != null ? router.agentsFor(target) : agents;
|
||||
AgentStatus status = refiner.refine(target, control.status(target), control);
|
||||
injector.onStatus(target, status);
|
||||
} catch (HerdrException e) {
|
||||
// The worker's agent is gone — stop trying and unblock its waiters.
|
||||
|
||||
@@ -41,15 +41,32 @@ public final class StatusRefiner {
|
||||
}
|
||||
|
||||
/**
|
||||
* Return a trustworthy status for {@code target}. Any non-{@code UNKNOWN} {@code raw} is returned
|
||||
* unchanged; an {@code UNKNOWN} triggers a pane read and content classification. A read failure
|
||||
* leaves it {@code UNKNOWN} (the safe default: no delivery, and the stall path still applies).
|
||||
* Return a trustworthy status for {@code target}, reading its pane through this refiner's own
|
||||
* {@link AgentControl}. Equivalent to {@link #refine(String, AgentStatus, AgentControl)} with
|
||||
* that control — kept for callers that only ever talk to one herdr daemon.
|
||||
*/
|
||||
public AgentStatus refine(String target, AgentStatus raw) {
|
||||
return refine(target, raw, agents);
|
||||
}
|
||||
|
||||
/**
|
||||
* Return a trustworthy status for {@code target}. Any non-{@code UNKNOWN} {@code raw} is returned
|
||||
* unchanged; an {@code UNKNOWN} triggers a pane read (through {@code control}) and content
|
||||
* classification. A read failure leaves it {@code UNKNOWN} (the safe default: no delivery, and
|
||||
* the stall path still applies).
|
||||
*
|
||||
* <p>CB-185: {@code control} must be the {@link AgentControl} for the <em>same</em> daemon the
|
||||
* raw status was sampled from — a router splits lead and member targets across two herdr
|
||||
* daemons, and reading a lead's pane through the member client (or vice versa) fails to find
|
||||
* the pane and leaves the target wedged at {@code UNKNOWN} forever. Callers that route per
|
||||
* target (e.g. {@code StatusPoller}) must pass that target's control explicitly rather than
|
||||
* relying on the control fixed at construction.
|
||||
*/
|
||||
public AgentStatus refine(String target, AgentStatus raw, AgentControl control) {
|
||||
if (raw != AgentStatus.UNKNOWN) return raw;
|
||||
String pane;
|
||||
try {
|
||||
pane = agents.read(target, PROBE_SOURCE);
|
||||
pane = control.read(target, PROBE_SOURCE);
|
||||
} catch (RuntimeException e) {
|
||||
log.debug("status refine read for {} failed; leaving UNKNOWN: {}", target, e.getMessage());
|
||||
return AgentStatus.UNKNOWN;
|
||||
|
||||
@@ -1033,34 +1033,40 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* pass it through {@code tab.create}/{@code pane.split}. Returns the directory for teardown
|
||||
* registration, or {@code null} when the policy does not apply.
|
||||
*
|
||||
* <p>The allow-list handed to the generator is the derived profile set ({@link
|
||||
* MemberEnvAllowList#derive}) UNIONed with the exact keys of THIS launch's env map — names the
|
||||
* daemon itself injects must survive its own control. {@code SSH_AUTH_SOCK} is added ONLY when
|
||||
* the config explicitly allows it; by default it is absent, so the scrub blanks it like any
|
||||
* other non-derived name.
|
||||
* <p>The allow-list handed to the generator is the derived profile set UNIONed with the
|
||||
* operator's own {@code memberCredentials.allow:} names ({@link MemberEnvAllowList#derive(
|
||||
* Collection, Set)} — CB-633 follow-up) and with the exact keys of THIS launch's env map —
|
||||
* names the daemon itself injects must survive its own control. {@code SSH_AUTH_SOCK} is added
|
||||
* ONLY when the config explicitly allows it, EVEN IF the operator also listed it under
|
||||
* {@code allow:}; by default it is absent, so the scrub blanks it like any other non-derived
|
||||
* name. It stays a one-off decision because it is a live handle to the operator's ssh-agent, not
|
||||
* a value — a member holding it can sign with every key the agent holds, so letting it ride in
|
||||
* on the generic {@code allow:} list would hand that out for an unrelated reason.
|
||||
*/
|
||||
private Path applyEnvironmentAllowListPolicy(FleetConfig.Profile cfg, Launch launch) {
|
||||
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
|
||||
if (creds == null || !creds.isAllowList()) {
|
||||
return null;
|
||||
}
|
||||
Set<String> allowed = derivedAllowedNames(creds, launch);
|
||||
String loginShell = resolveEnv("SHELL");
|
||||
boolean zsh = loginShell != null && (loginShell.endsWith("/zsh") || loginShell.equals("zsh"));
|
||||
if (!zsh) {
|
||||
// A non-zsh login shell ignores ZDOTDIR entirely: NO scrub would run, so pretending
|
||||
// otherwise would be worse than saying so. Warn loudly and fall back to the CB-596
|
||||
// sentinel overlay over the enumerated known: names — weaker (a sourced file can undo
|
||||
// it), but strictly better than nothing.
|
||||
// it), but strictly better than nothing. Deliberately no "allowed N of M" line here: the
|
||||
// scrub this count describes does not run on this path, so printing it would tell an
|
||||
// operator that a fraction of names were blocked when the real number blocked is zero.
|
||||
// logCredentialGap's WARN (below) is the only signal for this path.
|
||||
warnNonZsh(loginShell);
|
||||
overlayBlockedCredentials(launch.env(), creds);
|
||||
logCredentialGap(creds);
|
||||
return null;
|
||||
}
|
||||
Set<String> allowed = new java.util.TreeSet<>(MemberEnvAllowList.derive(profiles.values()));
|
||||
if (creds.sshAuthSockAllowed()) {
|
||||
allowed.add(SSH_AUTH_SOCK);
|
||||
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
|
||||
allowed.addAll(launch.env().keySet());
|
||||
// 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).
|
||||
logAllowListCoverage(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 "
|
||||
@@ -1069,8 +1075,42 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
return dir;
|
||||
}
|
||||
|
||||
/**
|
||||
* The full kept-name set for this spawn: the profile-derived names, unioned with {@code
|
||||
* memberCredentials.allow:} (CB-633 follow-up — previously ignored by this whole policy), the
|
||||
* ssh-agent handle when explicitly allowed, and the exact keys of THIS launch's own env map.
|
||||
*/
|
||||
private Set<String> derivedAllowedNames(FleetConfig.MemberCredentials creds, Launch launch) {
|
||||
Set<String> allowed = new java.util.TreeSet<>(
|
||||
MemberEnvAllowList.derive(profiles.values(), creds.allowSet()));
|
||||
if (creds.sshAuthSockAllowed()) {
|
||||
allowed.add(SSH_AUTH_SOCK);
|
||||
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
|
||||
allowed.addAll(launch.env().keySet());
|
||||
return allowed;
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up: one INFO line per allow-list spawn WHOSE SCRUB ACTUALLY RUNS, so an operator
|
||||
* can read a single log line and know the scrub ran and how much of the visible environment it
|
||||
* will keep. Callable ONLY from the zsh branch of {@link #applyEnvironmentAllowListPolicy}, after
|
||||
* the shell gate — logging it before that gate (or on the non-zsh fallback, where nothing is
|
||||
* scrubbed) would tell an operator a fraction of names were blocked when the real number blocked
|
||||
* is zero, which is worse than not logging at all. {@code M} is {@link #hostEnvNames}' size (the
|
||||
* daemon's own environment — see that field's javadoc for why it stands in for the pane's, which
|
||||
* the daemon has no channel to inspect at spawn time) and {@code N} is how many of those names
|
||||
* survive {@code allowed} (including the {@code LC_*} prefix rule). Neither number is a constant:
|
||||
* both come from the actual derived set and the actual environment this spawn sees. Never logs a
|
||||
* variable NAME or VALUE — only the counts.
|
||||
*/
|
||||
private void logAllowListCoverage(Set<String> allowed) {
|
||||
Set<String> hostNames = hostEnvNames.get();
|
||||
long kept = hostNames.stream().filter(name -> MemberEnvAllowList.keeps(allowed, name)).count();
|
||||
log.info("member credentials: allowed {} of {}", kept, hostNames.size());
|
||||
}
|
||||
|
||||
/** The operator ssh-agent handle — kept ONLY by explicit config decision, never by default. */
|
||||
private static final String SSH_AUTH_SOCK = "SSH_AUTH_SOCK";
|
||||
private static final String SSH_AUTH_SOCK = MemberEnvAllowList.SSH_AUTH_SOCK;
|
||||
|
||||
/**
|
||||
* CB-633: a non-zsh login shell means the allow-list control CANNOT run — say so once per
|
||||
|
||||
@@ -35,9 +35,27 @@ import java.util.TreeSet;
|
||||
* <p>{@code SSH_AUTH_SOCK} is deliberately NOT here. It is a handle to the operator's ssh-agent — a
|
||||
* member holding it can sign with the operator's keys — so keeping it is a config decision
|
||||
* ({@code memberCredentials.sshAuthSock: allow}), not a derivation default.
|
||||
*
|
||||
* <p><b>CB-633 follow-up:</b> the union also includes {@code memberCredentials.allow:} — the
|
||||
* operator's own explicit list. Before this, {@code policy: allow-list} silently ignored every name
|
||||
* an operator wrote under {@code allow:} unless a profile happened to carry it too, which meant
|
||||
* turning the policy on could blank credentials working members already depended on. {@code
|
||||
* SSH_AUTH_SOCK} is the one exception: even when the operator lists it under {@code allow:}, it is
|
||||
* excluded here and added back ONLY by the caller when {@code sshAuthSock: allow} is explicitly set
|
||||
* (see {@link #SSH_AUTH_SOCK}'s javadoc) — it is a live handle to the operator's own ssh-agent, not
|
||||
* a value, so treating it like any other allow-listed name would hand a member every key the
|
||||
* operator's agent holds the moment they typed the name under {@code allow:} for an unrelated
|
||||
* reason.
|
||||
*/
|
||||
public final class MemberEnvAllowList {
|
||||
|
||||
/**
|
||||
* The operator's ssh-agent socket path. Deliberately excluded from {@link #derive}'s union of
|
||||
* {@code memberCredentials.allow:} — see the class javadoc's CB-633 follow-up note. Governed
|
||||
* ONLY by {@code memberCredentials.sshAuthSock}, never by appearing in {@code allow:}.
|
||||
*/
|
||||
public static final String SSH_AUTH_SOCK = "SSH_AUTH_SOCK";
|
||||
|
||||
/**
|
||||
* Names that are not credentials and that a login shell or agent binary genuinely needs.
|
||||
*
|
||||
@@ -73,9 +91,21 @@ public final class MemberEnvAllowList {
|
||||
|
||||
/**
|
||||
* Derive the allowed NAME set from the given profiles plus {@link #INFRASTRUCTURE_PASSTHROUGH}.
|
||||
* Deterministic (sorted) so generated scrub files are diffable run-to-run.
|
||||
* Equivalent to {@link #derive(Collection, Set)} with no operator-configured names — kept for
|
||||
* callers (and existing tests) that only care about the profile-derived half.
|
||||
*/
|
||||
public static Set<String> derive(Collection<FleetConfig.Profile> profiles) {
|
||||
return derive(profiles, Set.of());
|
||||
}
|
||||
|
||||
/**
|
||||
* Derive the allowed NAME set: the profile-derived union above, PLUS {@code configuredAllow} —
|
||||
* the operator's own {@code memberCredentials.allow:} list (CB-633 follow-up). {@code
|
||||
* SSH_AUTH_SOCK} is dropped from {@code configuredAllow} even if the operator listed it there;
|
||||
* see the class javadoc for why. Deterministic (sorted) so generated scrub files are diffable
|
||||
* run-to-run.
|
||||
*/
|
||||
public static Set<String> derive(Collection<FleetConfig.Profile> profiles, Set<String> configuredAllow) {
|
||||
Set<String> derived = new TreeSet<>(INFRASTRUCTURE_PASSTHROUGH);
|
||||
if (profiles != null) {
|
||||
for (FleetConfig.Profile p : profiles) {
|
||||
@@ -87,6 +117,13 @@ public final class MemberEnvAllowList {
|
||||
}
|
||||
}
|
||||
}
|
||||
if (configuredAllow != null) {
|
||||
for (String name : configuredAllow) {
|
||||
if (name != null && !name.isBlank() && !SSH_AUTH_SOCK.equals(name)) {
|
||||
derived.add(name);
|
||||
}
|
||||
}
|
||||
}
|
||||
return Set.copyOf(derived);
|
||||
}
|
||||
|
||||
|
||||
@@ -29,8 +29,9 @@ import java.util.concurrent.TimeoutException;
|
||||
*
|
||||
* <p><strong>Mapping — consume-and-hold with deferred manual ack.</strong> Each target has a durable
|
||||
* queue {@code agent.<target>.inbox}. The gateway that owns the target starts a manual-ack consumer
|
||||
* ({@link #own}) that pulls persistent messages off that queue into an in-memory <em>held</em> map
|
||||
* (keyed by {@code msgId}) but does <em>not</em> ack them. {@link #peek} returns that snapshot;
|
||||
* ({@link #own}) that pulls persistent messages, up to its prefetch window, off that queue into an
|
||||
* in-memory <em>held</em> map (keyed by {@code msgId}) but does <em>not</em> ack them.
|
||||
* {@link #peek} returns that snapshot;
|
||||
* {@link #ack} acks the broker delivery-tag and drops the entry. Because messages stay unacked until
|
||||
* the owning gateway actually drains them, a crash (or a {@code java -jar} bounce) before caller-ack
|
||||
* leaves them on the broker — it redelivers on reconnect. That is the durability the in-memory
|
||||
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.fleet.msg;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.metrics.FleetMetrics;
|
||||
import dev.ltms.fleet.metrics.Metrics;
|
||||
@@ -189,6 +190,7 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
private final AgentControl agents;
|
||||
private final HerdrRouter router;
|
||||
private final Injector injector;
|
||||
private final Rendezvous rendezvous;
|
||||
private final ReplyInbox inbox;
|
||||
@@ -204,6 +206,25 @@ public final class MessageService {
|
||||
new ConcurrentHashMap<>();
|
||||
/** Async tickets paused on a specific {@code fleet_ask} turn. */
|
||||
private final ConcurrentHashMap<String, Task> asyncTasksByTurn = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* Targets whose most recent {@code fleet_reply} arrived with no send awaiting it (CB-640) —
|
||||
* {@link Rendezvous#resolve} returned {@code false} and the reply was queued into the inbox
|
||||
* instead (see {@link #reply}). The reply itself is not lost (it sits in the inbox for a
|
||||
* later drain), but the stranding is a fact the health layer needs to see. Bounded by the
|
||||
* target's own lifecycle rather than a TTL: an entry is cleared the next time this target's
|
||||
* delivery is accepted ({@link #send}) or the target is torn down ({@link #abandon}), so the
|
||||
* map holds at most one entry per session with an unresolved stranding right now.
|
||||
*/
|
||||
private final ConcurrentHashMap<String, Boolean> strandedReplies = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* Targets whose last send timed out with {@link Outcome#TIMED_OUT_QUEUED} (CB-640) — the
|
||||
* message never reached the {@link Injector} delivery window before the caller's deadline, so
|
||||
* it is still sitting in the injector's own per-target queue. Set where {@link #send} already
|
||||
* computes {@code wasDelivered} for that outcome; no new queue is kept here, only the fact.
|
||||
* Cleared the same way as {@link #strandedReplies}: the next accepted delivery for the target
|
||||
* ({@link #send} opening a fresh waiter) or a teardown ({@link #abandon}).
|
||||
*/
|
||||
private final ConcurrentHashMap<String, Boolean> queuedDeliveries = new ConcurrentHashMap<>();
|
||||
private final AtomicLong ticketSeq = new AtomicLong();
|
||||
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
|
||||
Thread.ofVirtual().name("bridge-async-", 0).factory());
|
||||
@@ -238,6 +259,7 @@ public final class MessageService {
|
||||
MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox,
|
||||
ReplyPushLoop pushLoop, Metrics metrics, LongSupplier nowNanos) {
|
||||
this.agents = agents;
|
||||
this.router = null;
|
||||
this.injector = injector;
|
||||
this.rendezvous = rendezvous;
|
||||
this.inbox = inbox;
|
||||
@@ -246,6 +268,20 @@ public final class MessageService {
|
||||
this.nowNanos = nowNanos;
|
||||
}
|
||||
|
||||
public MessageService(HerdrRouter router, Injector injector, Rendezvous rendezvous, ReplyInbox inbox,
|
||||
ReplyPushLoop pushLoop, Metrics metrics) {
|
||||
this.agents = null;
|
||||
this.router = router;
|
||||
this.injector = injector;
|
||||
this.rendezvous = rendezvous;
|
||||
this.inbox = inbox;
|
||||
this.pushLoop = pushLoop;
|
||||
this.metrics = metrics;
|
||||
this.nowNanos = System::nanoTime;
|
||||
}
|
||||
|
||||
private AgentControl agentsFor(String target) { return router != null ? router.agentsFor(target) : agents; }
|
||||
|
||||
/** Create with an explicit {@link ReplyInbox} and no push loop. */
|
||||
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox) {
|
||||
this(agents, injector, rendezvous, inbox, null);
|
||||
@@ -258,7 +294,7 @@ public final class MessageService {
|
||||
|
||||
/** Current lifecycle status of a worker (the {@code GET /sessions/{id}/status} surface). */
|
||||
public AgentStatus status(String target) {
|
||||
return agents.status(target);
|
||||
return agentsFor(target).status(target);
|
||||
}
|
||||
|
||||
/** Read-only delegation fact for fleet views. */
|
||||
@@ -271,6 +307,55 @@ public final class MessageService {
|
||||
return !inbox.peek(target).isEmpty();
|
||||
}
|
||||
|
||||
/**
|
||||
* Read-only delegation fact for fleet views (CB-640): {@code target}'s last send timed out
|
||||
* before the {@link Injector} ever delivered it — the caller saw
|
||||
* {@link Outcome#TIMED_OUT_QUEUED} (see the {@code TimeoutException} branch of {@link #send}),
|
||||
* and the message is still sitting in the injector's per-target queue waiting for the worker
|
||||
* to go idle. Distinct from {@link Outcome#TIMED_OUT_WORKING}, where delivery already happened
|
||||
* and only the reply is outstanding. Cleared the next time this target's delivery is accepted
|
||||
* or the target is abandoned — see {@link #queuedDeliveries}.
|
||||
*/
|
||||
public boolean hasQueuedDelivery(String target) {
|
||||
return target != null && queuedDeliveries.containsKey(target);
|
||||
}
|
||||
|
||||
/**
|
||||
* Read-only delegation fact for fleet views (CB-640): {@code target}'s last {@code fleet_reply}
|
||||
* arrived while no send was waiting for it, so {@link Rendezvous#resolve} returned
|
||||
* {@code false} and {@link #reply} fell back to queueing it in the inbox (see the CB-307
|
||||
* javadoc there and {@code docs/CB-307-Reliable-Delivery.md} §1). Cleared the next time this
|
||||
* target's delivery is accepted or the target is abandoned — see {@link #strandedReplies}.
|
||||
*/
|
||||
public boolean hasStrandedReply(String target) {
|
||||
return target != null && strandedReplies.containsKey(target);
|
||||
}
|
||||
|
||||
/**
|
||||
* Read-only delegation fact for fleet views (CB-640): an async ticket is still
|
||||
* {@link Phase#PENDING} against {@code target}, yet nothing is actually in flight for it — no
|
||||
* open rendezvous waiter ({@link #hasAcceptedDelivery}) and no message still sitting in the
|
||||
* injector's queue ({@link #hasQueuedDelivery}). A healthy PENDING ticket can briefly look this
|
||||
* way while its virtual thread has not yet been scheduled or is blocked on the session lock
|
||||
* behind another send to the same target, so this is a snapshot fact for the health classifier
|
||||
* to weigh across ticks, not proof on its own that the ticket is stuck. It also genuinely
|
||||
* persists — not just as a passing race — once an async {@code fleet_ask} lapses unanswered:
|
||||
* {@link #ask} clears the ticket's question and returns it to {@code PENDING}, but {@link #send}
|
||||
* already closed the forward waiter the instant the question surfaced, so the target has
|
||||
* neither an accepted nor a queued delivery left to show for it.
|
||||
*/
|
||||
public boolean hasOrphanedDelegation(String target) {
|
||||
if (target == null || hasAcceptedDelivery(target) || hasQueuedDelivery(target)) {
|
||||
return false;
|
||||
}
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null && !task.future.isDone()) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Route a worker's explicit {@code fleet_reply}: resolve an open send, or queue it in the
|
||||
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
|
||||
@@ -288,6 +373,9 @@ public final class MessageService {
|
||||
return true; // a live send took it — unchanged fast path
|
||||
}
|
||||
inbox.publish(session, UUID.randomUUID().toString(), content);
|
||||
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
|
||||
// worker whose replies keep missing their waiter, not only the queue depth this leaves behind.
|
||||
strandedReplies.put(session, Boolean.TRUE);
|
||||
// A rising inbox share is the signal CB-307 exists to make visible: the worker finished but
|
||||
// nobody was waiting, so delivery now depends on the push loop and a drain.
|
||||
count(FleetMetrics.REPLIES, "path", "inbox");
|
||||
@@ -341,6 +429,9 @@ public final class MessageService {
|
||||
* @return true if a live waiter was failed
|
||||
*/
|
||||
public boolean abandon(String target, String reason) {
|
||||
// CB-640: the session is gone — nothing will ever accept or deliver into it now.
|
||||
strandedReplies.remove(target);
|
||||
queuedDeliveries.remove(target);
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
||||
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
|
||||
boolean asyncFailed = false;
|
||||
@@ -446,6 +537,11 @@ public final class MessageService {
|
||||
// enqueue failure is safely closed by the finally below: nothing is left queued, and the
|
||||
// failed send leaves no stale waiter behind.
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
|
||||
// CB-640: this send now owns target's delivery, so any earlier stranded-reply or
|
||||
// still-queued fact no longer describes the live state — clear both rather than let
|
||||
// them outlive the send that supersedes them.
|
||||
strandedReplies.remove(target);
|
||||
queuedDeliveries.remove(target);
|
||||
try {
|
||||
if (task != null) {
|
||||
asyncTasksByWaiter.put(reply, task);
|
||||
@@ -464,6 +560,11 @@ public final class MessageService {
|
||||
} catch (TimeoutException e) {
|
||||
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
|
||||
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
|
||||
if (!wasDelivered) {
|
||||
// CB-640: still sitting in the injector's queue, waiting for the member to
|
||||
// go idle — record the fact for fleet health (see queuedDeliveries).
|
||||
queuedDeliveries.put(target, Boolean.TRUE);
|
||||
}
|
||||
return recorded(new Reply(
|
||||
wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
|
||||
} catch (ExecutionException e) {
|
||||
@@ -704,7 +805,7 @@ public final class MessageService {
|
||||
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
|
||||
private String liveStatus(String target) {
|
||||
try {
|
||||
return agents.status(target).name().toLowerCase();
|
||||
return agentsFor(target).status(target).name().toLowerCase();
|
||||
} catch (RuntimeException e) {
|
||||
return "unknown";
|
||||
}
|
||||
|
||||
@@ -31,6 +31,7 @@ import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Predicate;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
@@ -54,11 +55,12 @@ public final class FleetApp {
|
||||
/** Context attribute under which the resolved caller is stashed by the auth filter. */
|
||||
private static final String CALLER = "fleetd.caller";
|
||||
|
||||
private final HerdrClient herdr;
|
||||
private final HerdrClient herdr; // lead daemon
|
||||
private final HerdrClient memberHerdr; // CB-185: member daemon (same object when unconfigured)
|
||||
private final PeerLauncher workers;
|
||||
private final SessionManager sessions; // CB-301: authoritative session registry
|
||||
private final MessageService messages;
|
||||
private final MemberPresence presence; // CB-113: which workers are MCP-connected (available)
|
||||
private final Predicate<String> deliverable;
|
||||
private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable)
|
||||
private final CallerResolver auth; // CB-501: null → authz not enforced (legacy behaviour)
|
||||
private final Metrics metrics; // CB-502: null → /metrics not exposed
|
||||
@@ -70,9 +72,9 @@ public final class FleetApp {
|
||||
* behaviour without each needing an auth fixture.
|
||||
*/
|
||||
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet) {
|
||||
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null);
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet) {
|
||||
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null, presence::isPresent);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -82,13 +84,39 @@ public final class FleetApp {
|
||||
* the endpoint
|
||||
*/
|
||||
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) {
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) {
|
||||
this(herdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, presence::isPresent);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param deliverable the injector's readiness gate, shared so status reports its real result
|
||||
*/
|
||||
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
|
||||
Predicate<String> deliverable) {
|
||||
this(herdr, herdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, deliverable);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param herdr the lead daemon's client
|
||||
* @param memberHerdr the member daemon's client (CB-185); pass the same instance as
|
||||
* {@code herdr} for a single-daemon deployment — {@code healthz}/{@code
|
||||
* sessions} then make exactly one herdr call each, unchanged from before
|
||||
* the two-daemon router existed
|
||||
* @param deliverable the injector's readiness gate, shared so status reports its real result
|
||||
*/
|
||||
public FleetApp(HerdrClient herdr, HerdrClient memberHerdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
|
||||
Predicate<String> deliverable) {
|
||||
this.herdr = herdr;
|
||||
this.memberHerdr = memberHerdr != null ? memberHerdr : herdr;
|
||||
this.workers = workers;
|
||||
this.sessions = sessions;
|
||||
this.messages = messages;
|
||||
this.presence = presence;
|
||||
this.deliverable = deliverable;
|
||||
this.mcpServlet = mcpServlet;
|
||||
this.auth = auth;
|
||||
this.metrics = metrics;
|
||||
@@ -179,30 +207,64 @@ public final class FleetApp {
|
||||
ctx.status(200).contentType("text/plain; version=0.0.4; charset=utf-8").result(metrics.render());
|
||||
}
|
||||
|
||||
/** Liveness + herdr reachability. 200 when herdr answers ping, 503 otherwise. */
|
||||
/**
|
||||
* Liveness + herdr reachability. 200 only when BOTH daemons answer ping — 503 otherwise
|
||||
* (CB-185). With no {@code memberHerdrSocket} configured {@code memberHerdr == herdr}, so this
|
||||
* makes exactly the one {@code ping} call it always did and reports the same body; with a
|
||||
* 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}.
|
||||
*/
|
||||
private void healthz(Context ctx) {
|
||||
JsonNode pong;
|
||||
try {
|
||||
JsonNode pong = herdr.call("ping");
|
||||
ctx.status(200).json(Map.of(
|
||||
"status", "ok",
|
||||
"herdr", Map.of(
|
||||
"version", pong.path("version").asText(""),
|
||||
"protocol", pong.path("protocol").asInt())));
|
||||
pong = herdr.call("ping");
|
||||
} catch (HerdrException e) {
|
||||
ctx.status(503).json(Map.of(
|
||||
"status", "degraded",
|
||||
"herdr", "unreachable",
|
||||
"detail", e.getMessage()));
|
||||
return;
|
||||
}
|
||||
if (memberHerdr != herdr) {
|
||||
try {
|
||||
memberHerdr.call("ping");
|
||||
} catch (HerdrException e) {
|
||||
ctx.status(503).json(Map.of(
|
||||
"status", "degraded",
|
||||
"herdr", "member unreachable",
|
||||
"detail", e.getMessage()));
|
||||
return;
|
||||
}
|
||||
}
|
||||
ctx.status(200).json(Map.of(
|
||||
"status", "ok",
|
||||
"herdr", Map.of(
|
||||
"version", pong.path("version").asText(""),
|
||||
"protocol", pong.path("protocol").asInt())));
|
||||
}
|
||||
|
||||
/** Sessions view derived from herdr {@code workspace.list} (one workspace → one row). */
|
||||
/**
|
||||
* Sessions view derived from herdr {@code workspace.list} (one workspace → one row), merged
|
||||
* across both daemons (CB-185). With no {@code memberHerdrSocket} configured {@code
|
||||
* memberHerdr == herdr}, so this calls {@code workspace.list} exactly once, same as before the
|
||||
* router existed; with a second daemon configured, calling it twice would silently drop every
|
||||
* member workspace (they live on the member daemon only).
|
||||
*/
|
||||
private void sessions(Context ctx) {
|
||||
if (!allow(ctx, Authz.Action.READ, null)) {
|
||||
return;
|
||||
}
|
||||
JsonNode result = herdr.call("workspace.list");
|
||||
List<Map<String, Object>> out = new ArrayList<>();
|
||||
collectSessions(herdr, out);
|
||||
if (memberHerdr != herdr) {
|
||||
collectSessions(memberHerdr, out);
|
||||
}
|
||||
ctx.status(200).json(Map.of("sessions", out));
|
||||
}
|
||||
|
||||
private static void collectSessions(HerdrClient client, List<Map<String, Object>> out) {
|
||||
JsonNode result = client.call("workspace.list");
|
||||
for (JsonNode w : result.path("workspaces")) {
|
||||
out.add(Map.of(
|
||||
"id", w.path("workspace_id").asText(""),
|
||||
@@ -211,7 +273,6 @@ public final class FleetApp {
|
||||
"paneCount", w.path("pane_count").asInt(),
|
||||
"agentStatus", w.path("agent_status").asText("unknown")));
|
||||
}
|
||||
ctx.status(200).json(Map.of("sessions", out));
|
||||
}
|
||||
|
||||
/** Discovery: every agent herdr tracks, keyed by its Claude session UUID. */
|
||||
@@ -507,9 +568,9 @@ public final class FleetApp {
|
||||
|
||||
/**
|
||||
* Live lifecycle status of a worker (MCP `fleet_status` wraps this in CB-105), plus its
|
||||
* <em>readiness</em> (CB-113): {@code ready} is true once the worker's Claude has connected the
|
||||
* bridge MCP — the reliable "available to receive a task" signal, unlike bare {@code idle}, which
|
||||
* is also true during boot.
|
||||
* <em>readiness</em>: {@code ready} is true when the injector can deliver to the target. A
|
||||
* spawned member must connect the bridge MCP first, while a known lead is ready without member
|
||||
* presence. This differs from bare {@code idle}, which is also true during member boot.
|
||||
*/
|
||||
private void sessionStatus(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
@@ -520,7 +581,7 @@ public final class FleetApp {
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
body.put("sessionId", id);
|
||||
body.put("status", messages.status(id).name().toLowerCase());
|
||||
body.put("ready", presence.isPresent(id));
|
||||
body.put("ready", deliverable.test(id));
|
||||
// CB-582: a worker paused mid-turn in an async fleet_ask is otherwise invisible to a
|
||||
// status poll — surface the open question and how to answer it, same as fleet_poll's
|
||||
// Phase.ASKING view.
|
||||
|
||||
@@ -7,6 +7,8 @@ import java.io.BufferedReader;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStreamReader;
|
||||
import java.io.UncheckedIOException;
|
||||
import java.net.URI;
|
||||
import java.net.URISyntaxException;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
@@ -20,6 +22,7 @@ import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
@@ -64,6 +67,10 @@ public final class GitWorktrees implements Worktrees {
|
||||
/** What {@link #isolateToolSurface} writes for {@code .autoenv}: a valid, empty env file. */
|
||||
private static final String NEUTRAL_AUTOENV_CONFIG = "";
|
||||
|
||||
/** A credential helper command which reads only an environment variable at Git call time. */
|
||||
private static final String ENVIRONMENT_CREDENTIAL_HELPER = "!f() { if [ \"$1\" = get ]; then "
|
||||
+ "printf 'username=%s\\npassword=%s\\n\\n' git \"$WORKER_GITEA_TOKEN\"; fi; }; f";
|
||||
|
||||
/**
|
||||
* A tracked project config that is hostile in a provisioned worktree, and what to replace it
|
||||
* with. {@link #file} is the repo-relative path; {@link #stub} is a neutral but VALID payload for
|
||||
@@ -81,6 +88,7 @@ public final class GitWorktrees implements Worktrees {
|
||||
);
|
||||
|
||||
private final String configuredRoot;
|
||||
private final Consumer<String> afterWorktreeAdded;
|
||||
private final SecureRandom random = new SecureRandom();
|
||||
private final AtomicLong seq = new AtomicLong();
|
||||
|
||||
@@ -91,7 +99,13 @@ public final class GitWorktrees implements Worktrees {
|
||||
|
||||
/** @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of the repo root. */
|
||||
public GitWorktrees(String configuredRoot) {
|
||||
this(configuredRoot, _ -> {});
|
||||
}
|
||||
|
||||
/** Test seam for changing a real worktree between its creation and its security check. */
|
||||
GitWorktrees(String configuredRoot, Consumer<String> afterWorktreeAdded) {
|
||||
this.configuredRoot = configuredRoot;
|
||||
this.afterWorktreeAdded = afterWorktreeAdded == null ? _ -> {} : afterWorktreeAdded;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -107,11 +121,148 @@ public final class GitWorktrees implements Worktrees {
|
||||
}
|
||||
String wt = path.toAbsolutePath().toString();
|
||||
log.info("adding worktree branch={} path={} base={}", branch, wt, base);
|
||||
removeUserInfoFromHttpsOrigin(repoRoot);
|
||||
exec("git", "-C", repoRoot, "worktree", "add", wt, "-b", branch, base);
|
||||
afterWorktreeAdded.accept(wt);
|
||||
requireCredentialFreeHttpsOrigin(wt);
|
||||
configureEnvironmentCredentialHelper(repoRoot, wt);
|
||||
configureHttpsUrlRewriteForSshOrigin(repoRoot, wt);
|
||||
isolateToolSurface(wt);
|
||||
return wt;
|
||||
}
|
||||
|
||||
/**
|
||||
* A linked worktree shares its primary checkout's git config. Remove HTTPS user info before
|
||||
* adding one, so a credential accidentally embedded in that config cannot reach the member.
|
||||
*/
|
||||
private void removeUserInfoFromHttpsOrigin(String repoRoot) {
|
||||
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
|
||||
return;
|
||||
}
|
||||
String origin = exec("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
|
||||
URI uri;
|
||||
try {
|
||||
uri = new URI(origin);
|
||||
} catch (URISyntaxException e) {
|
||||
throw new WorktreeException("origin URL is invalid; cannot provision a safe worktree", e);
|
||||
}
|
||||
if (!"https".equalsIgnoreCase(uri.getScheme()) || uri.getUserInfo() == null) {
|
||||
return;
|
||||
}
|
||||
int schemeEnd = origin.indexOf("://") + 3;
|
||||
int userInfoEnd = origin.indexOf('@', schemeEnd);
|
||||
if (userInfoEnd < schemeEnd) {
|
||||
throw new WorktreeException("origin URL has invalid HTTPS user info; cannot provision safely");
|
||||
}
|
||||
String cleanOrigin = origin.substring(0, schemeEnd) + origin.substring(userInfoEnd + 1);
|
||||
exec("git", "-C", repoRoot, "remote", "set-url", "origin", cleanOrigin);
|
||||
log.info("removed HTTPS user info from forge origin before provisioning worktree");
|
||||
}
|
||||
|
||||
/** Refuse the worktree if Git still resolves any HTTPS origin URL with embedded credentials. */
|
||||
private void requireCredentialFreeHttpsOrigin(String worktreePath) {
|
||||
if (exitCode("git", "-C", worktreePath, "remote", "get-url", "--all", "origin") != 0) {
|
||||
return;
|
||||
}
|
||||
String origins = exec("git", "-C", worktreePath, "remote", "get-url", "--all", "origin");
|
||||
for (String origin : origins.split("\\R")) {
|
||||
try {
|
||||
URI uri = new URI(origin);
|
||||
if ("https".equalsIgnoreCase(uri.getScheme()) && uri.getUserInfo() != null) {
|
||||
throw new WorktreeException("worktree origin contains HTTPS user info; refusing provision");
|
||||
}
|
||||
} catch (URISyntaxException e) {
|
||||
throw new WorktreeException("worktree origin URL is invalid; refusing provision", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** 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");
|
||||
// An empty helper resets values inherited from the system or global config. Without it Git
|
||||
// asks the next helper after this one, which can expose an operator-level credential.
|
||||
exec("git", "-C", worktreePath, "config", "--worktree", "--replace-all", "credential.helper", "");
|
||||
exec("git", "-C", worktreePath, "config", "--worktree", "--add", "credential.helper",
|
||||
ENVIRONMENT_CREDENTIAL_HELPER);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link #configureEnvironmentCredentialHelper} only ever fires for an HTTPS origin — Git never
|
||||
* consults a {@code credential.helper} for an SSH transport. This repo's own origin is
|
||||
* {@code ssh://git@git.ltms.dev:2224/fleet/fleetd.git}, so a member sitting on that origin never
|
||||
* reaches the helper and the repo-scoped {@code WORKER_GITEA_TOKEN} is simply not used.
|
||||
*
|
||||
* <p>An earlier version of this javadoc justified the rewrite by claiming a member <em>cannot</em>
|
||||
* push once {@code memberCredentials.policy: allow-list} blocks {@code SSH_AUTH_SOCK}, because
|
||||
* "there is no private key file on this host, only an ssh-agent socket". That premise is false
|
||||
* (fleetd #184): {@code ssh -G} resolves a readable, passphrase-free {@code IdentityFile} outside
|
||||
* {@code ~/.ssh}, and a member — same OS user — pushes over SSH with the socket blanked. The
|
||||
* rewrite is still worth having, but for the reason below rather than that one: it routes the
|
||||
* member through its own scoped token instead of the operator's ssh identity, which is what makes
|
||||
* a member's pushes attributable and revocable.
|
||||
*
|
||||
* <p>The fix is a <em>worktree-scoped</em> URL rewrite: {@code url.<https-base>.insteadOf
|
||||
* <ssh-base>}, set with {@code --worktree} so it lands only in
|
||||
* {@code <worktree>/.git/worktrees/<name>/config.worktree} (enabled by
|
||||
* {@code extensions.worktreeConfig}, already turned on above) and never touches the shared
|
||||
* repo-level config the primary checkout also reads. {@code insteadOf} — not
|
||||
* {@code pushInsteadOf} — because a member may also need to fetch or rebase, and both should go
|
||||
* through the member's own token for the same reason.
|
||||
*
|
||||
* <p>The host (and, for the rewrite's SSH-side match, the port) come from parsing the origin
|
||||
* itself — never a hardcoded forge host, which is exactly what #177 removed. An origin that is
|
||||
* already {@code https://} is left alone; the credential helper already covers it. An origin
|
||||
* that is neither {@code ssh://} nor {@code https://} — including the scp-like shorthand
|
||||
* ({@code git@host:path}, no scheme) — is left untouched deliberately: that shorthand's
|
||||
* {@code host:path} split is defined by the user's ssh_config aliases, not by URI syntax, so
|
||||
* guessing at it risks rewriting to the wrong place. A repo provisioned from that form keeps
|
||||
* today's (broken, if the policy blocks the agent) SSH-only behaviour rather than a wrong rewrite.
|
||||
*/
|
||||
/**
|
||||
* Blank the user-info of a remote URL before it reaches a log. A remote URL is not obviously a
|
||||
* credential channel, which is exactly why one has leaked here three times ({@code git remote -v}
|
||||
* printing a token inline, and fleetd #157 / #182). An {@code ssh://} authority normally carries
|
||||
* only {@code git@}, so this usually changes nothing — it is here so that the one origin that
|
||||
* does carry a secret cannot print it. Matches every {@code ://…@} pair, not just the first.
|
||||
*/
|
||||
private static String redactUserInfo(String url) {
|
||||
return url == null ? null : url.replaceAll("://[^@/]*@", "://<redacted>@");
|
||||
}
|
||||
|
||||
private void configureHttpsUrlRewriteForSshOrigin(String repoRoot, String worktreePath) {
|
||||
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
|
||||
return;
|
||||
}
|
||||
String origin = exec("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
|
||||
URI uri;
|
||||
try {
|
||||
uri = new URI(origin);
|
||||
} catch (URISyntaxException e) {
|
||||
log.warn("origin URL {} is not a valid URI; skipping worktree HTTPS rewrite", redactUserInfo(origin));
|
||||
return;
|
||||
}
|
||||
String scheme = uri.getScheme();
|
||||
if (!"ssh".equalsIgnoreCase(scheme)) {
|
||||
// Already https:// (the credential helper covers it), or a scheme-less/scp-like origin
|
||||
// left alone on purpose — see the javadoc above.
|
||||
log.debug("origin scheme is not ssh ({}) — no worktree HTTPS rewrite needed", redactUserInfo(origin));
|
||||
return;
|
||||
}
|
||||
String host = uri.getHost();
|
||||
String authority = uri.getRawAuthority();
|
||||
if (host == null || host.isBlank() || authority == null || authority.isBlank()) {
|
||||
log.warn("ssh origin {} has no resolvable host; skipping worktree HTTPS rewrite", redactUserInfo(origin));
|
||||
return;
|
||||
}
|
||||
String sshBase = "ssh://" + authority + "/";
|
||||
String httpsBase = "https://" + host + "/";
|
||||
exec("git", "-C", worktreePath, "config", "--worktree", "--replace-all",
|
||||
"url." + httpsBase + ".insteadOf", sshBase);
|
||||
log.info("worktree {} rewrites {} to {} (worktree-scoped; parent checkout untouched)",
|
||||
worktreePath, redactUserInfo(sshBase), httpsBase);
|
||||
}
|
||||
|
||||
/**
|
||||
* Neutralize the worktree's worktree-hostile project configs so a worker inherits only the tools
|
||||
* and environment its launcher mounts (the bridge via {@code --mcp-config}, the opencode config
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-185: {@code ConnectionIdentity} must resolve a caller's pane on EITHER herdr daemon (a
|
||||
* lead's MCP connection resolves against the lead daemon; a member's against the member daemon).
|
||||
* Pinning {@code PaneLocator} to {@code memberHerdr} alone — the bug this guards against — leaves
|
||||
* every lead's own connection unresolvable ({@code callerTerminal == null}) the moment
|
||||
* {@code memberHerdrSocket} names a second daemon, which breaks {@code fleet_reply}/{@code
|
||||
* fleet_ask} and {@code fleet_whoami} for a lead. A unit test on {@link
|
||||
* dev.ltms.fleet.herdr.PaneLocator} alone (see {@code PaneLocatorTest}) proves the class CAN
|
||||
* search two clients, but not that {@code Fleetd.main} actually wires it that way — hence this
|
||||
* source-level assertion, the same technique {@code FleetdHerdrControlConstructionTest} uses.
|
||||
*/
|
||||
class FleetdConnectionIdentityConstructionTest {
|
||||
@Test
|
||||
void connectionIdentitySearchesBothDaemonsNotJustTheMemberOne() throws Exception {
|
||||
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
assertFalse(source.contains("new PaneLocator(memberHerdr)"),
|
||||
"PaneLocator must not be pinned to the member daemon alone — a lead's own "
|
||||
+ "connection resolves against the LEAD daemon and would never be found");
|
||||
assertTrue(source.contains("new PaneLocator(herdr, memberHerdr)"),
|
||||
"PaneLocator must search the lead daemon first, then the member daemon");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-185: {@code FleetApp} must be constructed with BOTH herdr clients (the lead's and the
|
||||
* member's), never the raw lead-only {@code herdr}. Passing only {@code herdr} — the bug this
|
||||
* guards against — makes {@code GET /healthz} green while the member daemon is down (so every
|
||||
* spawn fails invisibly) and silently drops every member workspace from {@code GET /sessions}.
|
||||
* A behavioural test on {@code FleetApp} alone (see {@code FleetAppTwoDaemonTest}) proves the
|
||||
* class merges/gates correctly when given two clients, but not that {@code Fleetd.main} actually
|
||||
* passes it two — hence this source-level assertion, mirroring
|
||||
* {@code FleetdHerdrControlConstructionTest}.
|
||||
*/
|
||||
class FleetdFleetAppConstructionTest {
|
||||
@Test
|
||||
void fleetAppIsConstructedWithBothHerdrDaemons() throws Exception {
|
||||
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
assertFalse(source.contains("new FleetApp(herdr, workers,"),
|
||||
"FleetApp must not be constructed with the lead-only herdr client");
|
||||
assertTrue(source.contains("new FleetApp(herdr, memberHerdr, workers,"),
|
||||
"FleetApp must be constructed with both the lead and the member herdr client");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
|
||||
class FleetdHerdrControlConstructionTest {
|
||||
@Test
|
||||
void fleetdDelegatesStatefulControlsToTheRouter() throws Exception {
|
||||
// AgentControl caches paneByTerminal, so the router must be its only production factory.
|
||||
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
assertFalse(source.contains("new AgentControl("));
|
||||
assertFalse(source.contains("new WorkspaceControl("));
|
||||
}
|
||||
}
|
||||
@@ -1,14 +1,18 @@
|
||||
package dev.ltms.fleet.health;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
@@ -17,10 +21,15 @@ import dev.ltms.fleet.config.FleetConfig;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledThreadPoolExecutor;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
@@ -40,7 +49,7 @@ class FleetHealthMonitorTest {
|
||||
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox());
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, sessions::roster, messages, scheduler, () -> 1, 60,
|
||||
(_, _) -> { });
|
||||
600, (_, _) -> { });
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
assertEquals(1, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
|
||||
@@ -52,7 +61,7 @@ class FleetHealthMonitorTest {
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, () -> 1, 60, (_, _) -> { });
|
||||
scheduler, () -> 1, 60, 600, (_, _) -> { });
|
||||
monitor.tick();
|
||||
herdr.healthy(true);
|
||||
monitor.tick();
|
||||
@@ -71,7 +80,7 @@ class FleetHealthMonitorTest {
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, () -> 1, 60, (_, _) -> { });
|
||||
scheduler, () -> 1, 60, 600, (_, _) -> { });
|
||||
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
|
||||
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
|
||||
monitor.stop();
|
||||
@@ -82,6 +91,148 @@ class FleetHealthMonitorTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test void failedListMarksEveryMemberControlLinkDownAndReschedules() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr().healthy(false);
|
||||
ScheduledThreadPoolExecutor scheduler = new ScheduledThreadPoolExecutor(1);
|
||||
FleetHealthMonitor monitor = monitor(herdr, List.of(
|
||||
member("term_one", MemberSession.State.READY, 0, 0),
|
||||
member("term_two", MemberSession.State.BUSY, 0, 0)), scheduler, () -> 1, 600,
|
||||
(_, _) -> { });
|
||||
|
||||
monitor.tick();
|
||||
|
||||
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
|
||||
.contains("member=term_one state=CONTROL_LINK_DOWN")).count());
|
||||
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
|
||||
.contains("member=term_two state=CONTROL_LINK_DOWN")).count());
|
||||
assertEquals(1, scheduler.getQueue().size());
|
||||
monitor.stop();
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test void missingRosterMemberGoesGoneAndFailsTargetOnceAcrossTicks() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = monitor(herdr,
|
||||
List.of(member("term_missing", MemberSession.State.READY, 0, 0)), scheduler,
|
||||
() -> 1, 600, failTarget);
|
||||
|
||||
monitor.tick();
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
|
||||
assertEquals(1, failTarget.calls.size());
|
||||
assertEquals("term_missing", failTarget.calls.get(0).target());
|
||||
assertTrue(failTarget.calls.get(0).reason().contains("GONE"));
|
||||
}
|
||||
|
||||
@Test void spawningMemberBecomesNeverReadyOnlyAfterReadinessGrace() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
AtomicLong clock = new AtomicLong(FleetHealthMonitor.READINESS_GRACE_NANOS - 1);
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = monitor(new FakeHerdr(),
|
||||
List.of(member("term_starting", MemberSession.State.SPAWNING, 0, 0)), scheduler,
|
||||
clock::get, 600, failTarget);
|
||||
|
||||
monitor.tick();
|
||||
assertEquals(0, failTarget.calls.size());
|
||||
clock.set(FleetHealthMonitor.READINESS_GRACE_NANOS);
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
|
||||
assertEquals(1, failTarget.calls.size());
|
||||
assertTrue(failTarget.calls.get(0).reason().contains("NEVER_READY"));
|
||||
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
|
||||
.contains("state=NEVER_READY previous=STARTING")));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test void workingSuspectConfigChangesTheStallBoundary() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
long nowNanos = TimeUnit.SECONDS.toNanos(600);
|
||||
FakeHerdr herdr = new FakeHerdr()
|
||||
.withAgent("short", "term_short", "pane_short", "tab_short")
|
||||
.withAgent("long", "term_long", "pane_long", "tab_long");
|
||||
FleetConfig.Health shortConfig = new FleetConfig.Health(true, 30, 300, null, null);
|
||||
FleetConfig.Health longConfig = new FleetConfig.Health(true, 30, 601, null, null);
|
||||
var shortScheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
var longScheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor shortMonitor = monitor(herdr,
|
||||
List.of(member("term_short", MemberSession.State.BUSY, 0, 0)), shortScheduler,
|
||||
() -> nowNanos, shortConfig.workingSuspectAfterOrDefault(), (_, _) -> { });
|
||||
FleetHealthMonitor longMonitor = monitor(herdr,
|
||||
List.of(member("term_long", MemberSession.State.BUSY, 0, 0)), longScheduler,
|
||||
() -> nowNanos, longConfig.workingSuspectAfterOrDefault(), (_, _) -> { });
|
||||
|
||||
shortMonitor.tick();
|
||||
longMonitor.tick();
|
||||
shortMonitor.stop();
|
||||
longMonitor.stop();
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
|
||||
.contains("member=term_short state=STALL_SUSPECTED")));
|
||||
assertEquals(0, appender.list.stream().filter(event -> event.getFormattedMessage()
|
||||
.contains("member=term_long state=STALL_SUSPECTED")).count());
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test void workingSuspectConfigKeepsItsDefaultAndFloor() {
|
||||
assertEquals(600, new FleetConfig.Health(true, null, null, null, null)
|
||||
.workingSuspectAfterOrDefault());
|
||||
assertEquals(300, new FleetConfig.Health(true, null, 1, null, null)
|
||||
.workingSuspectAfterOrDefault());
|
||||
}
|
||||
|
||||
@Test void goneMemberRecoveryLogsOnceWithoutRefiringTargetFailure() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
Level previousLevel = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = monitor(herdr,
|
||||
List.of(member("term_recovered", MemberSession.State.READY, 0, 0)), scheduler,
|
||||
() -> 1, 600, failTarget);
|
||||
|
||||
monitor.tick();
|
||||
herdr.withAgent("recovered", "term_recovered", "pane_recovered", "tab_recovered");
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
|
||||
assertEquals(1, failTarget.calls.size());
|
||||
assertTrue(appender.list.stream().anyMatch(event -> event.getFormattedMessage()
|
||||
.contains("member=term_recovered recovered state=IDLE previous=GONE")));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(previousLevel);
|
||||
}
|
||||
}
|
||||
|
||||
// --- CB-580: a member that reaches GONE/NEVER_READY must fail its waiting tickets
|
||||
|
||||
private static FleetHealthMonitor monitorWith(BiConsumer<String, String> failTarget) {
|
||||
@@ -90,7 +241,23 @@ class FleetHealthMonitorTest {
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
return new FleetHealthMonitor(agents, java.util.List::of,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, () -> 1, 60, failTarget);
|
||||
scheduler, () -> 1, 60, 600, failTarget);
|
||||
}
|
||||
|
||||
private static FleetHealthMonitor monitor(FakeHerdr herdr, List<MemberSession> roster,
|
||||
java.util.concurrent.ScheduledExecutorService scheduler,
|
||||
LongSupplier clock, long workingSuspectAfterSeconds,
|
||||
BiConsumer<String, String> failTarget) {
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
return new FleetHealthMonitor(agents, () -> roster,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, clock, 60, workingSuspectAfterSeconds, failTarget);
|
||||
}
|
||||
|
||||
private static MemberSession member(String terminalId, MemberSession.State state,
|
||||
long spawnedAtNanos, long lastActivityAtNanos) {
|
||||
return new MemberSession("pane-" + terminalId, terminalId, "test", MemberRole.DEV, "/tmp", null,
|
||||
spawnedAtNanos, lastActivityAtNanos, 0, state, null, null);
|
||||
}
|
||||
|
||||
@Test void terminalTransitionFailsTheTargetOnce() {
|
||||
@@ -167,4 +334,104 @@ class FleetHealthMonitorTest {
|
||||
throw new RuntimeException("boom");
|
||||
}
|
||||
}
|
||||
|
||||
// --- CB-643: the two-tick gate on an orphaned delegation ------------------------------------
|
||||
|
||||
private static final String ORPHAN_WARNING = "member=term_a state=DELEGATION_ORPHANED";
|
||||
|
||||
/**
|
||||
* Everything one orphan test needs: a live target whose async ticket really is orphaned. The
|
||||
* recipe is the one MessageServiceTest proves for {@code hasOrphanedDelegation} — an unanswered
|
||||
* {@code fleet_ask} lapses, so the ticket returns to PENDING while its forward waiter is already
|
||||
* closed. {@code rendezvous} is exposed so a test can turn the fact off and on again.
|
||||
*/
|
||||
private record OrphanFleet(FakeHerdr herdr, Rendezvous rendezvous, MessageService messages,
|
||||
FleetHealthMonitor monitor) {
|
||||
}
|
||||
|
||||
private static OrphanFleet orphanedTarget(BiConsumer<String, String> failTarget) throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().withAgent("worker", "term_a", "pane-term_a", "tab_a")
|
||||
.readText("$ prompt");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
Injector injector = new Injector(agents);
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
inbox.own("term_a");
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
||||
|
||||
messages.sendAsync("term_a", "task that asks");
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting("term_a"), "the async send should have opened its waiter");
|
||||
injector.onStatus("term_a", AgentStatus.IDLE); // deliver the task
|
||||
injector.onStatus("term_a", AgentStatus.WORKING); // the worker picks it up
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT,
|
||||
messages.ask("term_a", "which config?", 200).outcome());
|
||||
assertTrue(messages.hasOrphanedDelegation("term_a"),
|
||||
"the lapsed ask should leave a PENDING ticket with nothing in flight");
|
||||
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, () -> List.of(
|
||||
member("term_a", MemberSession.State.READY, 0, 0)),
|
||||
messages, new ScheduledThreadPoolExecutor(1), () -> 1, 60, 600, failTarget);
|
||||
return new OrphanFleet(herdr, rendezvous, messages, monitor);
|
||||
}
|
||||
|
||||
private static long orphanWarnings(ListAppender<ILoggingEvent> appender) {
|
||||
return appender.list.stream()
|
||||
.filter(event -> event.getFormattedMessage().contains(ORPHAN_WARNING)).count();
|
||||
}
|
||||
|
||||
@Test void oneOrphanObservationIsNotYetReportedButTwoAre() throws Exception {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
OrphanFleet fleet = orphanedTarget((_, _) -> { });
|
||||
try {
|
||||
fleet.monitor().tick();
|
||||
assertEquals(0, orphanWarnings(appender),
|
||||
"one observation can be an ordinary race, so health must not report it yet");
|
||||
|
||||
fleet.monitor().tick();
|
||||
assertEquals(1, orphanWarnings(appender),
|
||||
"a second consecutive observation confirms the orphan and is reported once");
|
||||
} finally {
|
||||
fleet.messages().abandon("term_a", "test over");
|
||||
fleet.monitor().stop();
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test void aSingleCleanTickResetsTheOrphanStreak() throws Exception {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
OrphanFleet fleet = orphanedTarget((_, _) -> { });
|
||||
try {
|
||||
fleet.monitor().tick(); // first observation: streak 1
|
||||
|
||||
// Something is in flight for the target again, so the orphan fact reads false.
|
||||
var waiter = fleet.rendezvous().open("term_a");
|
||||
fleet.monitor().tick();
|
||||
assertEquals(0, orphanWarnings(appender), "a clean tick must clear the streak");
|
||||
|
||||
// Deregister that waiter — resolving it is not enough, the sender's close() is what
|
||||
// removes it — so the ticket is orphaned again, from a streak of zero.
|
||||
fleet.rendezvous().close("term_a", waiter);
|
||||
assertTrue(fleet.messages().hasOrphanedDelegation("term_a"));
|
||||
fleet.monitor().tick();
|
||||
assertEquals(0, orphanWarnings(appender),
|
||||
"the streak restarted, so this first observation is not reported either");
|
||||
|
||||
fleet.monitor().tick();
|
||||
assertEquals(1, orphanWarnings(appender), "two consecutive observations report once");
|
||||
} finally {
|
||||
fleet.messages().abandon("term_a", "test over");
|
||||
fleet.monitor().stop();
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -40,6 +40,7 @@ public final class FakeHerdr implements HerdrClient {
|
||||
private int workerTabPaneCount = 1;
|
||||
private String paneCloseErrorCode = null;
|
||||
private String agentSendErrorCode = null;
|
||||
private boolean noPanes = false;
|
||||
private volatile String agentStatus = "idle"; // steady-state agent.get status
|
||||
private volatile String readText = "worker transcript tail"; // canned agent.read output
|
||||
private int pinnedStarts = 0; // how many upcoming agent.start calls report a fixed pane
|
||||
@@ -75,6 +76,15 @@ public final class FakeHerdr implements HerdrClient {
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code pane.list} report no panes at all — models a second herdr daemon (CB-185) that
|
||||
* simply does not host the pane a {@link PaneLocator} is searching for.
|
||||
*/
|
||||
public FakeHerdr withNoPanes() {
|
||||
this.noPanes = true;
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Set the {@code agent_status} that {@code agent.get} reports (drives the injector). */
|
||||
public FakeHerdr agentStatus(String status) {
|
||||
this.agentStatus = status;
|
||||
@@ -268,7 +278,9 @@ public final class FakeHerdr implements HerdrClient {
|
||||
case "pane.get" -> mapper.readTree("""
|
||||
{"type":"pane_info","pane":{"pane_id":"w9:pW","workspace_id":"w9",
|
||||
"tab_id":"w9:t2","agent_status":"idle"}}""");
|
||||
case "pane.list" -> mapper.readTree("""
|
||||
case "pane.list" -> noPanes
|
||||
? mapper.readTree("{\"type\":\"pane_list\",\"panes\":[]}")
|
||||
: mapper.readTree("""
|
||||
{"type":"pane_list","panes":[
|
||||
{"pane_id":"w2:p7","terminal_id":"term_a","workspace_id":"w2","tab_id":"w2:t7","agent":"claude"},
|
||||
{"pane_id":"w2:p9","terminal_id":"term_shell","workspace_id":"w2","tab_id":"w2:t8"}]}""");
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
package dev.ltms.fleet.herdr;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertSame;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotSame;
|
||||
|
||||
class HerdrRouterTest {
|
||||
@Test
|
||||
void absentMemberClientSharesControlsForEveryTarget() {
|
||||
FakeHerdr client = new FakeHerdr();
|
||||
HerdrRouter router = new HerdrRouter(client, null, id -> id.equals("lead"));
|
||||
|
||||
assertSame(router.leadAgents(), router.memberAgents());
|
||||
assertSame(router.leadSpaces(), router.memberSpaces());
|
||||
assertSame(router.leadAgents(), router.agentsFor("lead"));
|
||||
assertSame(router.leadAgents(), router.agentsFor("member"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void separateClientsRouteLeadAndMemberTargets() {
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
HerdrRouter router = new HerdrRouter(lead, member, id -> id.equals("lead"));
|
||||
|
||||
assertNotSame(router.leadAgents(), router.memberAgents());
|
||||
assertNotSame(router.leadSpaces(), router.memberSpaces());
|
||||
assertSame(router.leadAgents(), router.agentsFor("lead"));
|
||||
assertSame(router.memberAgents(), router.agentsFor("member"));
|
||||
}
|
||||
}
|
||||
@@ -24,4 +24,43 @@ class PaneLocatorTest {
|
||||
assertNull(loc.terminalForPid(0));
|
||||
assertNull(loc.terminalForPid(-1));
|
||||
}
|
||||
|
||||
// --- two-daemon fallback (CB-185) -----------------------------------------
|
||||
|
||||
@Test
|
||||
void fallsBackToTheMemberClientWhenTheLeadHasNoMatch() {
|
||||
// The caller's pane lives on the member daemon only (e.g. the caller is a spawned
|
||||
// member) — the lead client reports no panes at all, so the locator must fall back.
|
||||
HerdrClient lead = new FakeHerdr().withNoPanes();
|
||||
HerdrClient member = new FakeHerdr();
|
||||
PaneLocator two = new PaneLocator(lead, member);
|
||||
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID));
|
||||
}
|
||||
|
||||
@Test
|
||||
void searchesTheLeadClientBeforeTheMemberClient() {
|
||||
// The caller's pane lives on the LEAD daemon (e.g. the caller is a peer lead) — with two
|
||||
// daemons, resolving it must not depend on the member client having a matching pane.
|
||||
HerdrClient lead = new FakeHerdr();
|
||||
HerdrClient member = new FakeHerdr().withNoPanes();
|
||||
PaneLocator two = new PaneLocator(lead, member);
|
||||
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID));
|
||||
}
|
||||
|
||||
@Test
|
||||
void nullWhenNeitherClientHasTheMatch() {
|
||||
PaneLocator two = new PaneLocator(new FakeHerdr().withNoPanes(), new FakeHerdr().withNoPanes());
|
||||
assertNull(two.terminalForPid(FakeHerdr.WORKER_PID));
|
||||
}
|
||||
|
||||
@Test
|
||||
void collapsesToOneScanWhenLeadAndMemberAreTheSameClient() {
|
||||
// The single-daemon deployment (no memberHerdrSocket configured): the two-arg constructor
|
||||
// must behave exactly like the one-arg constructor, including making only one herdr call.
|
||||
FakeHerdr shared = new FakeHerdr();
|
||||
PaneLocator two = new PaneLocator(shared, shared);
|
||||
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID));
|
||||
long paneListCalls = shared.calls.stream().filter(c -> c.method().equals("pane.list")).count();
|
||||
assertEquals(1, paneListCalls, "same-object lead/member must scan exactly once, not twice");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -197,10 +197,15 @@ class CompletionResolverTest {
|
||||
void resolvesSynchronouslyBeforePostTurnContextClearing() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
// fleetd#164: an injectable clock, so this real (non-crash) turn lands outside MIN_TURN_NANOS
|
||||
// — captureBaseline and resolveBeforePostAction below run back-to-back with no real delay.
|
||||
long[] clock = {1_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter));
|
||||
herdr.readText("⏺ answer that /clear would erase\n❯ ");
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // this turn took longer than the floor
|
||||
|
||||
resolver.resolveBeforePostAction("term_a");
|
||||
|
||||
@@ -219,7 +224,11 @@ class CompletionResolverTest {
|
||||
String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(longBlock);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
// fleetd#164: an injectable clock so captureBaseline and resolve (back-to-back, no real
|
||||
// delay) don't trip the too-fast-turn floor — this test is about the suppression guard, not timing.
|
||||
long[] clock = {1_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
|
||||
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block
|
||||
@@ -227,6 +236,7 @@ class CompletionResolverTest {
|
||||
assertEquals(CompletionResolver.MAX_SCRAPE_CHARS, turn.baseline().length(),
|
||||
"the delivery baseline is clipped to the same cap resolve() applies to the tail");
|
||||
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // outside the floor
|
||||
resolver.resolve("term_a", turn); // scrape unchanged → clipped tail == baseline → suppress
|
||||
|
||||
assertFalse(waiter.isDone(),
|
||||
@@ -249,12 +259,13 @@ class CompletionResolverTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolvesWhenTheScrapeItselfFailsEvenWithABaselinePresent() {
|
||||
// The most important branch of the CB-115 guard: a failed read means the resolver could not
|
||||
// SEE the screen — "couldn't see", not "no change". It must still resolve the send (an empty
|
||||
// tail beats hanging until the caller's timeout), even though a baseline was captured. The
|
||||
// baseline here is "" (an empty pane at delivery), so without the !scrapeFailed clause the
|
||||
// byte-identical guard would wrongly match the empty tail and suppress.
|
||||
void aFailedScrapeResolvesAsAFailureEvenWithABaselinePresent() {
|
||||
// fleetd#164: this test used to assert that a failed read resolved the send as a SUCCESS
|
||||
// carrying an empty string ("an empty tail beats hanging until the caller's timeout") — that
|
||||
// was the bug this ticket fixes: a lost turn and a genuine empty answer looked identical to
|
||||
// every caller. This test encoded the bug and is changed here: a failed read must fail the
|
||||
// send instead, naming the member, whether or not a baseline was captured (the baseline here
|
||||
// is "", an empty pane at delivery — proof this isn't the CB-115 misattribution path either).
|
||||
FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
@@ -265,8 +276,82 @@ class CompletionResolverTest {
|
||||
|
||||
assertTrue(waiter.isDone(),
|
||||
"a failed scrape must still resolve the send, not hang until the caller's timeout");
|
||||
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind());
|
||||
assertEquals("", waiter.getNow(null).text(), "the tail is empty because the screen was unreadable");
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"a failed read is a lost turn, not a successful empty reply");
|
||||
assertTrue(waiter.getNow(null).text().contains("term_a"),
|
||||
"the failure names the member: " + waiter.getNow(null).text());
|
||||
assertTrue(waiter.getNow(null).text().contains("could not be read"),
|
||||
"the failure explains the scrape could not be read: " + waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
// --- fleetd#164: an empty (but readable) scrape must never resolve as a success --------------
|
||||
|
||||
@Test
|
||||
void anEmptyScrapeResolvesAsAFailureNamingTheMember() {
|
||||
// The core defect: a scrape that read CLEANLY but produced zero characters used to resolve
|
||||
// the send as a SUCCESS carrying "" — indistinguishable, to every caller, from a worker that
|
||||
// genuinely finished with nothing to say. A lost turn must never look like a real empty reply.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("");
|
||||
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(), "an empty scrape must still resolve the send, not hang");
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"an empty scrape is a lost turn, not a successful empty reply");
|
||||
assertTrue(waiter.getNow(null).text().contains("term_a"),
|
||||
"the failure names the member: " + waiter.getNow(null).text());
|
||||
assertTrue(waiter.getNow(null).text().toLowerCase().contains("empty"),
|
||||
"the failure says the scrape was empty: " + waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
// --- fleetd#164: a BUSY -> DONE transition inside the floor is a crash, not a fast answer ------
|
||||
|
||||
@Test
|
||||
void aBusyToDoneTransitionInsideTheFloorResolvesAsAFailure() {
|
||||
// The exact fleetd#164 scenario: the backend returned an HTTP 400 before the worker did
|
||||
// anything, and the member went BUSY -> DONE in ~1s. That transition alone is indistinguishable
|
||||
// from a genuine (if unusually fast) completion, so the resolver leans on the floor to catch it.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ HTTP 400: invalid request\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
long[] clock = {10_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]); // delivered "now"
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // 1ns inside the floor
|
||||
|
||||
resolver.resolve("term_a", turn);
|
||||
|
||||
assertTrue(waiter.isDone(), "a suspiciously fast turn must still resolve (as a failure)");
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
|
||||
assertTrue(waiter.getNow(null).text().contains("term_a"),
|
||||
"the failure names the member: " + waiter.getNow(null).text());
|
||||
assertTrue(waiter.getNow(null).text().contains("HTTP 400"),
|
||||
"the failure carries whatever was on screen: " + waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBusyToDoneTransitionJustOutsideTheFloorResolvesNormally() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ a real, if quick, answer\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
long[] clock = {10_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]); // delivered "now"
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // 1ns outside the floor
|
||||
|
||||
resolver.resolve("term_a", turn);
|
||||
|
||||
assertTrue(waiter.isDone());
|
||||
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(),
|
||||
"a turn that took longer than the floor resolves normally");
|
||||
assertEquals("a real, if quick, answer", waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
// --- CB-115/CB-116 fail guard: an already-done or absent waiter is left alone ---------
|
||||
|
||||
@@ -0,0 +1,75 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
import dev.ltms.fleet.msg.TestTurnTokens;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-185: with a router split across two herdr daemons, {@link StatusPoller} must refine a raw
|
||||
* {@code UNKNOWN} status by reading the pane content from the SAME daemon the status was sampled
|
||||
* from — the lead daemon for a lead target, the member daemon for a member target. Reading the
|
||||
* wrong daemon never finds the pane, classification stays {@code UNKNOWN} forever, and the
|
||||
* status-gated {@link Injector} wedges: a queued message is never delivered.
|
||||
*
|
||||
* <p>This exercises the real production classes ({@code StatusPoller(HerdrRouter, ...)},
|
||||
* {@code Injector(HerdrRouter, ...)}) wired together, not a hand-built object graph — the earlier
|
||||
* three CB-185 bugs all passed exactly that kind of test while the real wiring stayed broken.
|
||||
*/
|
||||
class StatusPollerRoutingTest {
|
||||
|
||||
private static final String LEAD_TARGET = "term_a";
|
||||
|
||||
@Test
|
||||
void refinesALeadTargetFromTheLeadDaemonAndDelivers() throws Exception {
|
||||
// The lead daemon's pane is at a settled idle prompt; the member daemon's pane content is
|
||||
// unclassifiable garbage. A correct refiner reads the LEAD daemon and delivers.
|
||||
FakeHerdr leadHerdr = new FakeHerdr().agentStatus("unknown").readText("⏺ answer\n❯ ");
|
||||
FakeHerdr memberHerdr = new FakeHerdr().agentStatus("unknown")
|
||||
.readText("garbled ansi noise with no prompt");
|
||||
HerdrRouter router = new HerdrRouter(leadHerdr, memberHerdr, LEAD_TARGET::equals);
|
||||
Injector injector = new Injector(router, TurnListener.NOOP, _ -> true, _ -> {
|
||||
});
|
||||
StatusPoller poller = new StatusPoller(router, injector, 10);
|
||||
poller.start();
|
||||
try {
|
||||
CompletableFuture<Void> delivered =
|
||||
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET));
|
||||
// Must resolve quickly: refining against the WRONG daemon (member) never classifies
|
||||
// out of UNKNOWN, so this would time out under the bug.
|
||||
delivered.get(2, TimeUnit.SECONDS);
|
||||
} finally {
|
||||
poller.stop();
|
||||
}
|
||||
assertTrue(leadHerdr.called("agent.read"), "refine must probe the LEAD daemon's pane content");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aLeadTargetNeverDeliversWhenOnlyTheMemberDaemonIsClassifiable() throws Exception {
|
||||
// Inverted control: the member daemon's content WOULD classify to idle, but this is a lead
|
||||
// target — a correct implementation must not use it, so delivery must NOT happen.
|
||||
FakeHerdr leadHerdr = new FakeHerdr().agentStatus("unknown")
|
||||
.readText("garbled ansi noise with no prompt");
|
||||
FakeHerdr memberHerdr = new FakeHerdr().agentStatus("unknown").readText("⏺ answer\n❯ ");
|
||||
HerdrRouter router = new HerdrRouter(leadHerdr, memberHerdr, LEAD_TARGET::equals);
|
||||
Injector injector = new Injector(router, TurnListener.NOOP, _ -> true, _ -> {
|
||||
});
|
||||
StatusPoller poller = new StatusPoller(router, injector, 10);
|
||||
poller.start();
|
||||
try {
|
||||
CompletableFuture<Void> delivered =
|
||||
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET));
|
||||
assertThrows(TimeoutException.class, () -> delivered.get(500, TimeUnit.MILLISECONDS),
|
||||
"a lead target must never be refined from the member daemon's pane content");
|
||||
} finally {
|
||||
poller.stop();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -85,4 +85,34 @@ class StatusRefinerTest {
|
||||
|
||||
assertEquals(AgentStatus.UNKNOWN, refiner.refine("term_a", AgentStatus.UNKNOWN));
|
||||
}
|
||||
|
||||
// --- refine(target, raw, control) — CB-185 per-call routing ---------------
|
||||
|
||||
@Test
|
||||
void threeArgRefineReadsThroughTheGivenControlNotTheConstructedOne() {
|
||||
// The refiner is CONSTRUCTED with one control (standing in for "the member daemon"), but
|
||||
// a call names a DIFFERENT control (standing in for "the lead daemon") — the read must go
|
||||
// to the one passed to the call, since that is the daemon the raw status came from.
|
||||
FakeHerdr constructedWith = new FakeHerdr().readText("nothing recognizable here");
|
||||
FakeHerdr passedToCall = new FakeHerdr().readText("⏺ answer\n❯ ");
|
||||
StatusRefiner refiner = new StatusRefiner(new AgentControl(constructedWith));
|
||||
|
||||
AgentStatus result = refiner.refine("term_a", AgentStatus.UNKNOWN, new AgentControl(passedToCall));
|
||||
|
||||
assertEquals(AgentStatus.IDLE, result, "must classify from the PASSED control's pane content");
|
||||
assertTrue(passedToCall.called("agent.read"));
|
||||
assertFalse(constructedWith.called("agent.read"),
|
||||
"the control fixed at construction must not be read when a call-site control is given");
|
||||
}
|
||||
|
||||
@Test
|
||||
void twoArgRefineStillReadsTheConstructedControl() {
|
||||
// The legacy 2-arg overload (single-daemon callers) must keep using the constructed
|
||||
// control — this is refine(target, raw, control) called with the field as `control`.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ answer\n❯ ");
|
||||
StatusRefiner refiner = new StatusRefiner(new AgentControl(herdr));
|
||||
|
||||
assertEquals(AgentStatus.IDLE, refiner.refine("term_a", AgentStatus.UNKNOWN));
|
||||
assertTrue(herdr.called("agent.read"));
|
||||
}
|
||||
}
|
||||
|
||||
+128
-1
@@ -1,5 +1,9 @@
|
||||
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;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
@@ -8,6 +12,7 @@ import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
@@ -108,6 +113,122 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null);
|
||||
}
|
||||
|
||||
/** Same as {@link #allowList()} but with an operator-configured {@code allow:} list. */
|
||||
private static Supplier<FleetConfig.MemberCredentials> allowListWithAllow(List<String> allow) {
|
||||
return () -> new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, allow, List.of(), null);
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up: a name that lives ONLY in {@code memberCredentials.allow:} — no profile
|
||||
* mentions it — must survive the scrub the real spawn path generates. Calling {@code
|
||||
* MemberEnvAllowList.derive} directly (as {@link MemberEnvAllowListTest} does) would pass even
|
||||
* if {@code HerdrPeerLauncher} never threaded {@code allow:} into the derivation at all; this
|
||||
* test goes through {@link HerdrPeerLauncher#spawn}, the method the daemon actually calls at
|
||||
* spawn time, so it proves the union is wired in, not just correct in isolation.
|
||||
*/
|
||||
@Test
|
||||
void spawningUnderAllowListPolicyIncludesAnOperatorConfiguredAllowName() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr,
|
||||
allowListWithAllow(List.of("OPERATOR_ONLY_NAME")));
|
||||
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
Path dir = Path.of(launcher.env.get("ZDOTDIR"));
|
||||
assertTrue(readAll(dir.resolve(EnvAllowListScrub.SCRUB_FILE)).contains("OPERATOR_ONLY_NAME"),
|
||||
"a name only in memberCredentials.allow: must reach the generated scrub through the "
|
||||
+ "real launcher spawn path");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code SSH_AUTH_SOCK} is a live ssh-agent handle, not a value — it must stay blocked under
|
||||
* {@code allow-list} even when the operator lists it under {@code allow:}, because {@code
|
||||
* sshAuthSock} defaults to blocked. Governed ONLY by {@code memberCredentials.sshAuthSock}.
|
||||
*/
|
||||
@Test
|
||||
void sshAuthSockStaysBlockedEvenWhenListedInMemberCredentialsAllow() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr,
|
||||
allowListWithAllow(List.of("SSH_AUTH_SOCK")));
|
||||
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
Path dir = Path.of(launcher.env.get("ZDOTDIR"));
|
||||
String scrub = readAll(dir.resolve(EnvAllowListScrub.SCRUB_FILE));
|
||||
assertFalse(scrub.contains("'SSH_AUTH_SOCK'"),
|
||||
"SSH_AUTH_SOCK must not be on the derived allow-list just because the operator put "
|
||||
+ "it under allow: — sshAuthSock is unset here, so it defaults to block");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up criterion 3: on every allow-list spawn the daemon logs one INFO line, shaped
|
||||
* "member credentials: allowed N of M", with real counts — not constants. Real path: the count
|
||||
* is asserted after a real {@link HerdrPeerLauncher#spawn} call, reading the log the production
|
||||
* code actually emits.
|
||||
*/
|
||||
@Test
|
||||
void logsAnAllowedCountLineAgainstTheHostEnvironmentOnEverySpawn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
// INJECTED is a key of this launch's own env map, so it always survives; the other two are
|
||||
// neither derived from the profile nor configured anywhere, so they are blanked. Real
|
||||
// N=1 (INJECTED), real M=3 (all three names) — neither number is hardcoded in the assertion
|
||||
// by coincidence, they follow directly from this fixture.
|
||||
Set<String> hostEnvNames = Set.of(INJECTED, "SOME_UNRELATED_NAME", "ANOTHER_UNRELATED_NAME");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames);
|
||||
|
||||
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 {
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
assertTrue(appender.list.stream()
|
||||
.anyMatch(e -> "member credentials: allowed 1 of 3".equals(e.getFormattedMessage())),
|
||||
"expected 'member credentials: allowed 1 of 3', got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* Lead-review fix: on a NON-zsh shell no scrub ever runs (bash ignores {@code ZDOTDIR}), so the
|
||||
* "allowed N of M" line — which describes what the scrub does — must not be printed there either.
|
||||
* Before this fix the line was logged BEFORE the zsh gate, so a non-zsh host printed e.g.
|
||||
* "allowed 1 of 3" while blocking nothing at all, telling an operator a control ran when it did
|
||||
* not. Real path: goes through {@link HerdrPeerLauncher#spawn}, same as the sibling test above,
|
||||
* with the shell fixed to bash so the fallback branch is the one exercised.
|
||||
*/
|
||||
@Test
|
||||
void noAllowedCountLineIsEmittedOnTheNonZshFallbackPath() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Set<String> hostEnvNames = Set.of(INJECTED, "SOME_UNRELATED_NAME", "ANOTHER_UNRELATED_NAME");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/bash", () -> hostEnvNames);
|
||||
|
||||
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 {
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
assertFalse(appender.list.stream()
|
||||
.anyMatch(e -> e.getFormattedMessage().startsWith("member credentials: allowed ")),
|
||||
"no scrub runs on a non-zsh shell, so no 'allowed N of M' count may be printed — got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
private static String readAll(Path p) {
|
||||
try {
|
||||
return Files.readString(p);
|
||||
@@ -134,10 +255,16 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
}
|
||||
|
||||
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell) {
|
||||
this(herdr, creds, shell, null);
|
||||
}
|
||||
|
||||
/** Plus an injectable {@code hostEnvNames} source, for the "allowed N of M" log line test. */
|
||||
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
|
||||
Supplier<Set<String>> hostEnvNames) {
|
||||
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of("test", profile()), "test",
|
||||
name -> "SHELL".equals(name) ? shell : null,
|
||||
0, () -> 0L, () -> { }, null, creds);
|
||||
0, () -> 0L, () -> { }, null, creds, hostEnvNames);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -88,6 +88,39 @@ class MemberEnvAllowListTest {
|
||||
assertTrue(after.containsAll(Set.of("TOKEN_SECOND", "GIT_TOK", "SECOND_KEY")));
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up: a name that appears ONLY in {@code memberCredentials.allow:} — no profile
|
||||
* mentions it at all — must still survive the derivation. Before this fix {@code derive} never
|
||||
* saw {@code allow:}, so setting {@code policy: allow-list} silently blanked exactly this name.
|
||||
*/
|
||||
@Test
|
||||
void aNameOnlyInMemberCredentialsAllowSurvivesDerivation() {
|
||||
FleetConfig.Profile p = profile("p", "TOKEN_A", null, null, Map.of("KEY_A", "v"));
|
||||
|
||||
Set<String> derived = MemberEnvAllowList.derive(List.of(p), Set.of("OPERATOR_ONLY_NAME"));
|
||||
|
||||
assertTrue(derived.contains("OPERATOR_ONLY_NAME"),
|
||||
"memberCredentials.allow: must be unioned in, not ignored");
|
||||
// and the profile-derived half must still be present — this is a union, not a replacement.
|
||||
assertTrue(derived.contains("KEY_A"));
|
||||
assertTrue(derived.contains("TOKEN_A"));
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code SSH_AUTH_SOCK} is a live handle to the operator's ssh-agent, never a value — so it must
|
||||
* stay excluded from the derived set even when the operator lists it under {@code allow:} for an
|
||||
* unrelated reason. It is governed ONLY by {@code memberCredentials.sshAuthSock}, applied
|
||||
* separately by the caller ({@code HerdrPeerLauncher}).
|
||||
*/
|
||||
@Test
|
||||
void sshAuthSockInMemberCredentialsAllowIsStillExcluded() {
|
||||
Set<String> derived = MemberEnvAllowList.derive(List.of(), Set.of("SSH_AUTH_SOCK", "OTHER_NAME"));
|
||||
|
||||
assertFalse(derived.contains("SSH_AUTH_SOCK"),
|
||||
"SSH_AUTH_SOCK must never ride in on the generic allow: list");
|
||||
assertTrue(derived.contains("OTHER_NAME"), "other allow: names are unaffected");
|
||||
}
|
||||
|
||||
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
|
||||
@Test
|
||||
void keepsMatchesExactlyPlusTheLocalePrefixRule() {
|
||||
|
||||
@@ -0,0 +1,149 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import com.rabbitmq.client.AMQP;
|
||||
import com.rabbitmq.client.Channel;
|
||||
import com.rabbitmq.client.Connection;
|
||||
import com.rabbitmq.client.DeliverCallback;
|
||||
import com.rabbitmq.client.Delivery;
|
||||
import com.rabbitmq.client.Envelope;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.lang.reflect.InvocationHandler;
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.ArrayDeque;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
|
||||
/**
|
||||
* Pins the manual-ack prefetch behaviour without a broker. The fake channel models a broker that
|
||||
* sends no more than its QoS window of unacked deliveries. If {@link AmqpReplyInbox} starts acking
|
||||
* messages while it adds them to {@code held}, this test drains the whole fake queue instead.
|
||||
*/
|
||||
class AmqpReplyInboxPrefetchTest {
|
||||
|
||||
@Test
|
||||
void unackedDeliveriesKeepTheHeldBacklogAtThePrefetchWindow() {
|
||||
int prefetch = 3;
|
||||
int published = 8;
|
||||
PrefetchBroker broker = new PrefetchBroker();
|
||||
for (int i = 0; i < published; i++) {
|
||||
broker.publish("m" + i, "payload " + i);
|
||||
}
|
||||
|
||||
try (AmqpReplyInbox inbox = new AmqpReplyInbox(connectionFor(broker.channel()), prefetch)) {
|
||||
inbox.own("worker");
|
||||
|
||||
assertEquals(prefetch, inbox.peek("worker").size(),
|
||||
"held messages must stop at the unacked prefetch window");
|
||||
assertEquals(published - prefetch, broker.queuedCount(),
|
||||
"messages beyond the window must remain on the broker");
|
||||
assertEquals(0, broker.ackCount(), "receipt must not ack a held message");
|
||||
|
||||
inbox.ack("worker", "m0");
|
||||
|
||||
assertEquals(prefetch, inbox.peek("worker").size(),
|
||||
"one caller ack frees exactly one slot for the broker");
|
||||
assertEquals(published - prefetch - 1, broker.queuedCount(),
|
||||
"only one queued message may enter after one caller ack");
|
||||
assertEquals(1, broker.ackCount(), "only the caller ack may reach the broker");
|
||||
}
|
||||
}
|
||||
|
||||
private static Connection connectionFor(Channel consumeChannel) {
|
||||
Channel publishChannel = (Channel) Proxy.newProxyInstance(
|
||||
AmqpReplyInboxPrefetchTest.class.getClassLoader(), new Class<?>[] {Channel.class},
|
||||
(proxy, method, args) -> defaultValue(method.getReturnType()));
|
||||
AtomicInteger channelCalls = new AtomicInteger();
|
||||
InvocationHandler handler = (proxy, method, args) -> {
|
||||
if (method.getName().equals("createChannel") && (args == null || args.length == 0)) {
|
||||
return channelCalls.getAndIncrement() == 0 ? consumeChannel : publishChannel;
|
||||
}
|
||||
return defaultValue(method.getReturnType());
|
||||
};
|
||||
return (Connection) Proxy.newProxyInstance(AmqpReplyInboxPrefetchTest.class.getClassLoader(),
|
||||
new Class<?>[] {Connection.class}, handler);
|
||||
}
|
||||
|
||||
private static final class PrefetchBroker implements InvocationHandler {
|
||||
private final ArrayDeque<Delivery> queued = new ArrayDeque<>();
|
||||
private final Map<Long, Delivery> unacked = new LinkedHashMap<>();
|
||||
private DeliverCallback consumer;
|
||||
private int prefetch;
|
||||
private int acks;
|
||||
private long nextTag = 1;
|
||||
|
||||
Channel channel() {
|
||||
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxPrefetchTest.class.getClassLoader(),
|
||||
new Class<?>[] {Channel.class}, this);
|
||||
}
|
||||
|
||||
void publish(String msgId, String content) {
|
||||
queued.add(new Delivery(new Envelope(nextTag++, false, "", ""),
|
||||
new AMQP.BasicProperties.Builder().messageId(msgId).build(),
|
||||
content.getBytes(StandardCharsets.UTF_8)));
|
||||
}
|
||||
|
||||
int queuedCount() {
|
||||
return queued.size();
|
||||
}
|
||||
|
||||
int ackCount() {
|
||||
return acks;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object invoke(Object proxy, java.lang.reflect.Method method, Object[] args) throws IOException {
|
||||
switch (method.getName()) {
|
||||
case "basicQos" -> {
|
||||
prefetch = (int) args[0];
|
||||
return null;
|
||||
}
|
||||
case "basicConsume" -> {
|
||||
assertFalse((boolean) args[1], "the inbox consumer must use manual acknowledgements");
|
||||
consumer = (DeliverCallback) args[2];
|
||||
deliverAvailable();
|
||||
return "consumer";
|
||||
}
|
||||
case "basicAck" -> {
|
||||
unacked.remove((long) args[0]);
|
||||
acks++;
|
||||
deliverAvailable();
|
||||
return null;
|
||||
}
|
||||
default -> {
|
||||
return defaultValue(method.getReturnType());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void deliverAvailable() throws IOException {
|
||||
while (consumer != null && unacked.size() < prefetch && !queued.isEmpty()) {
|
||||
Delivery delivery = queued.removeFirst();
|
||||
unacked.put(delivery.getEnvelope().getDeliveryTag(), delivery);
|
||||
consumer.handle("consumer", delivery);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static Object defaultValue(Class<?> type) {
|
||||
if (!type.isPrimitive() || type == void.class) {
|
||||
return null;
|
||||
}
|
||||
if (type == boolean.class) {
|
||||
return false;
|
||||
}
|
||||
if (type == long.class) {
|
||||
return 0L;
|
||||
}
|
||||
if (type == int.class) {
|
||||
return 0;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
@@ -33,8 +33,17 @@ class MessageServiceTest {
|
||||
private final FakeHerdr herdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
private final Rendezvous rendezvous = new Rendezvous();
|
||||
/**
|
||||
* fleetd#164: this fixture drives a delivery and its completion back-to-back with no real time
|
||||
* between them, so the real clock would trip {@link CompletionResolver#MIN_TURN_NANOS} on every
|
||||
* completion-fallback test here. An ever-advancing fake clock stands in for the model-latency and
|
||||
* herdr round-trips a real turn would spend, so each delivery-then-resolve pair still lands
|
||||
* outside the floor.
|
||||
*/
|
||||
private final java.util.concurrent.atomic.AtomicLong resolverClock = new java.util.concurrent.atomic.AtomicLong();
|
||||
private final CompletionResolver completion =
|
||||
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none(),
|
||||
() -> resolverClock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1));
|
||||
private final Injector injector = new Injector(agents, completion);
|
||||
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
||||
@@ -1081,4 +1090,147 @@ class MessageServiceTest {
|
||||
assertEquals(phase, view.phase());
|
||||
return view;
|
||||
}
|
||||
|
||||
// --- CB-640: fleet health evidence accessors --------------------------------------------
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryIsFalseForAnUnknownTarget() {
|
||||
assertFalse(messages.hasQueuedDelivery("nobody-ever-sent-here"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryIsFalseBeforeAnyTimeout() {
|
||||
assertFalse(messages.hasQueuedDelivery(T));
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryIsTrueAfterAnUndeliveredSendTimesOut() {
|
||||
// Nothing ever delivers the message (never goes IDLE/BLOCKED), so the send times out with
|
||||
// TIMED_OUT_QUEUED — same setup as sendTimesOutBeforeDeliveryIsQueuedNotWorking above.
|
||||
MessageService.Reply r = messages.send(T, "never delivered", 50);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome());
|
||||
|
||||
assertTrue(messages.hasQueuedDelivery(T),
|
||||
"a TIMED_OUT_QUEUED send leaves the message still queued in the injector");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
|
||||
assertTrue(messages.hasQueuedDelivery(T));
|
||||
|
||||
// A fresh send accepts delivery (opens its own waiter) — the stale queued fact is cleared.
|
||||
CompletableFuture<MessageService.Reply> second = sendAsync();
|
||||
awaitWaiting();
|
||||
assertFalse(messages.hasQueuedDelivery(T),
|
||||
"a fresh accepted delivery supersedes the earlier queued fact");
|
||||
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
assertTrue(rendezvous.resolve(T, "second done"));
|
||||
second.get(5, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryClearsOnAbandon() {
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
|
||||
assertTrue(messages.hasQueuedDelivery(T));
|
||||
|
||||
messages.abandon(T, "session released");
|
||||
assertFalse(messages.hasQueuedDelivery(T), "a torn-down target has nothing left queued for it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyIsFalseForAnUnknownTarget() {
|
||||
assertFalse(messages.hasStrandedReply("nobody-ever-sent-here"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyIsFalseWhenTheReplyResolvedALiveSend() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
awaitWaiting();
|
||||
|
||||
assertTrue(messages.reply(T, "resolved-live"));
|
||||
assertFalse(messages.hasStrandedReply(T), "a reply that resolved an open send is not stranded");
|
||||
|
||||
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.REPLIED, r.outcome());
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyIsTrueWhenNoSendWasWaiting() {
|
||||
// No send is open for T — the reply queues into the inbox and is recorded as stranded.
|
||||
assertTrue(messages.reply(T, "nobody was waiting"));
|
||||
assertTrue(messages.hasStrandedReply(T),
|
||||
"a reply with no open send strands, even though it is safely queued in the inbox");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
|
||||
assertTrue(messages.reply(T, "stray"));
|
||||
assertTrue(messages.hasStrandedReply(T));
|
||||
|
||||
// The next accepted delivery for T clears the stale stranding fact — the one case the
|
||||
// ticket calls out as the one that matters.
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
awaitWaiting();
|
||||
assertFalse(messages.hasStrandedReply(T),
|
||||
"a stranded reply must clear once the target's delivery is accepted again");
|
||||
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
send.get(5, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyClearsOnAbandon() {
|
||||
assertTrue(messages.reply(T, "stray"));
|
||||
assertTrue(messages.hasStrandedReply(T));
|
||||
|
||||
messages.abandon(T, "session released");
|
||||
assertFalse(messages.hasStrandedReply(T), "a torn-down target has nothing left to strand");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasOrphanedDelegationIsFalseForAnUnknownTarget() {
|
||||
assertFalse(messages.hasOrphanedDelegation("nobody-ever-sent-here"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasOrphanedDelegationIsFalseWhilePendingTicketsHaveAnAcceptedDelivery() throws Exception {
|
||||
// Mirrors abandonFailsEveryPendingAsyncTicketForTheReleasedTarget above: "first" holds the
|
||||
// session lock and its waiter is open, so the target genuinely has something in flight even
|
||||
// though "second" and "third" are themselves parked (PENDING) behind the lock.
|
||||
messages.sendAsync(T, "first task");
|
||||
awaitWaiting(); // first task owns the target lock and rendezvous waiter
|
||||
String second = messages.sendAsync(T, "second task");
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(second).phase());
|
||||
|
||||
assertFalse(messages.hasOrphanedDelegation(T),
|
||||
"the target has an accepted delivery in flight (first), so nothing here is orphaned");
|
||||
|
||||
assertTrue(messages.abandon(T, "session released")); // release the lock and the parked tickets
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasOrphanedDelegationIsTrueOnceAnUnansweredAskLapsesBackToPending() throws Exception {
|
||||
// Same setup as unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget above:
|
||||
// once the ask lapses, the ticket goes back to PENDING but send() already closed the
|
||||
// forward waiter the instant the question surfaced — nothing is left in flight for T.
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT,
|
||||
messages.ask(T, "which config?", 200).outcome());
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
|
||||
|
||||
assertFalse(messages.hasAcceptedDelivery(T), "the forward waiter closed when the question surfaced");
|
||||
assertFalse(messages.hasQueuedDelivery(T), "this ticket never timed out as queued");
|
||||
assertTrue(messages.hasOrphanedDelegation(T),
|
||||
"a PENDING ticket with no accepted or queued delivery for its target is orphaned");
|
||||
|
||||
assertTrue(messages.abandon(T, "session released")); // clean up the still-open ticket
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,6 +32,7 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
import java.util.function.Predicate;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -64,6 +65,11 @@ class FleetAppTest {
|
||||
}
|
||||
|
||||
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement, Worktrees worktrees) {
|
||||
return start(herdr, workerBaseUrl, allow, placement, worktrees, ignored -> false);
|
||||
}
|
||||
|
||||
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
|
||||
Worktrees worktrees, Predicate<String> deliverable) {
|
||||
FleetConfig.Profile wcfg = new FleetConfig.Profile(
|
||||
"ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
placement, "fleet", "worker: {profile} #{n}", null, null, null);
|
||||
@@ -84,7 +90,8 @@ class FleetAppTest {
|
||||
// it directly so the inbox contract holds for those endpoints.
|
||||
inbox.own("term_a");
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
||||
app = new FleetApp(herdr, workers, sessions, messages, this.presence, null)
|
||||
app = new FleetApp(herdr, workers, sessions, messages, this.presence, null,
|
||||
null, null, id -> this.presence.isPresent(id) || deliverable.test(id))
|
||||
.build().start("127.0.0.1", 0);
|
||||
return app.port();
|
||||
}
|
||||
@@ -484,6 +491,16 @@ class FleetAppTest {
|
||||
assertTrue(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean());
|
||||
}
|
||||
|
||||
@Test
|
||||
void sessionStatusReportsRegisteredLeadAsReady() throws Exception {
|
||||
Map<String, String> leads = Map.of("term_lead", "terra");
|
||||
int port = start(new FakeHerdr(), "http://gx00.gw:8000", Set.of("gx00.gw"), "tab",
|
||||
new GitWorktrees(), leads::containsKey);
|
||||
|
||||
JsonNode body = mapper.readTree(req(port, "GET", "/sessions/term_lead/status").body());
|
||||
assertTrue(body.get("ready").asBoolean());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-582: a lead polling {@code GET /sessions/{id}/status} on its normal cadence — not the
|
||||
* ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code fleet_ask}
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
package dev.ltms.fleet.rest;
|
||||
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.net.URI;
|
||||
import java.net.http.HttpClient;
|
||||
import java.net.http.HttpRequest;
|
||||
import java.net.http.HttpResponse;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* CB-185: with a router split across two herdr daemons (lead + {@code memberHerdrSocket}),
|
||||
* {@link FleetApp#healthz} must require BOTH daemons to answer and {@link FleetApp#sessions}
|
||||
* (which the {@code GET /sessions} route calls) must merge workspaces from both — the bug this
|
||||
* guards against had {@code FleetApp} constructed with the raw lead-only client, so a down member
|
||||
* daemon was invisible behind a green {@code /healthz} (every spawn then fails) and every member
|
||||
* workspace was silently dropped from {@code GET /sessions}.
|
||||
*
|
||||
* <p>Builds the real {@link FleetApp} directly (not a hand-rolled stand-in) against only the two
|
||||
* herdr clients — the other collaborators are unused by the two routes under test here.
|
||||
*/
|
||||
class FleetAppTwoDaemonTest {
|
||||
|
||||
private final HttpClient http = HttpClient.newHttpClient();
|
||||
private Javalin app;
|
||||
|
||||
@AfterEach
|
||||
void stop() {
|
||||
if (app != null) app.stop();
|
||||
}
|
||||
|
||||
private int start(HerdrClient lead, HerdrClient member) {
|
||||
app = new FleetApp(lead, member, null, null, null, null, null, null, null, ignored -> false)
|
||||
.build().start("127.0.0.1", 0);
|
||||
return app.port();
|
||||
}
|
||||
|
||||
private HttpResponse<String> get(int port, String path) throws Exception {
|
||||
HttpRequest req = HttpRequest.newBuilder(URI.create("http://127.0.0.1:" + port + path)).GET().build();
|
||||
return http.send(req, HttpResponse.BodyHandlers.ofString());
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzIsGreenWhenBothDaemonsAnswer() throws Exception {
|
||||
int port = start(new FakeHerdr(), new FakeHerdr());
|
||||
assertEquals(200, get(port, "/healthz").statusCode());
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzIsDegradedWhenOnlyTheMemberDaemonIsDown() throws Exception {
|
||||
int port = start(new FakeHerdr(), new FakeHerdr().healthy(false));
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
assertEquals(503, res.statusCode(),
|
||||
"a down MEMBER daemon must not be masked by a healthy lead — every spawn goes "
|
||||
+ "through the member daemon");
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzIsDegradedWhenOnlyTheLeadDaemonIsDown() throws Exception {
|
||||
int port = start(new FakeHerdr().healthy(false), new FakeHerdr());
|
||||
assertEquals(503, get(port, "/healthz").statusCode());
|
||||
}
|
||||
|
||||
@Test
|
||||
void healthzMakesExactlyOneCallWhenLeadAndMemberAreTheSameClient() throws Exception {
|
||||
// Single-daemon deployment (no memberHerdrSocket) — must be byte-for-byte the old
|
||||
// behaviour: one ping call, 200 on success.
|
||||
FakeHerdr shared = new FakeHerdr();
|
||||
int port = start(shared, shared);
|
||||
assertEquals(200, get(port, "/healthz").statusCode());
|
||||
long pings = shared.calls.stream().filter(c -> c.method().equals("ping")).count();
|
||||
assertEquals(1, pings, "single-daemon deployment must make exactly one ping call");
|
||||
}
|
||||
|
||||
@Test
|
||||
void sessionsMergesWorkspacesFromBothDaemons() throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
FakeHerdr member = new FakeHerdr().withWorkspace("w9", "member-only-workspace");
|
||||
int port = start(lead, member);
|
||||
HttpResponse<String> res = get(port, "/sessions");
|
||||
assertEquals(200, res.statusCode(), res.body());
|
||||
assertTrue(res.body().contains("member-only-workspace"),
|
||||
"GET /sessions must not silently drop the member daemon's workspaces");
|
||||
}
|
||||
|
||||
@Test
|
||||
void sessionsMakesExactlyOneWorkspaceListCallWhenLeadAndMemberAreTheSameClient() throws Exception {
|
||||
FakeHerdr shared = new FakeHerdr();
|
||||
int port = start(shared, shared);
|
||||
assertEquals(200, get(port, "/sessions").statusCode());
|
||||
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");
|
||||
}
|
||||
}
|
||||
@@ -55,12 +55,21 @@ class GitWorktreesTest {
|
||||
}
|
||||
|
||||
private static void git(Path cwd, String... args) throws Exception {
|
||||
gitOutput(cwd, args);
|
||||
}
|
||||
|
||||
private static String gitOutput(Path cwd, String... args) throws Exception {
|
||||
List<String> cmd = new java.util.ArrayList<>(List.of("git"));
|
||||
cmd.addAll(List.of(args));
|
||||
Process p = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true).start();
|
||||
ProcessBuilder pb = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true);
|
||||
pb.environment().put("GIT_CONFIG_GLOBAL", "/dev/null");
|
||||
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
|
||||
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
|
||||
Process p = pb.start();
|
||||
String out = new String(p.getInputStream().readAllBytes());
|
||||
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: " + String.join(" ", cmd));
|
||||
assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out);
|
||||
return out;
|
||||
}
|
||||
|
||||
/** Pending changes to {@code file} in {@code cwd}, empty when git considers it unmodified. */
|
||||
@@ -205,6 +214,144 @@ class GitWorktreesTest {
|
||||
"expected an explicitly empty server map, got:\n" + body);
|
||||
}
|
||||
|
||||
/**
|
||||
* The worktree command shares the primary checkout's config, so this checks the URL git actually
|
||||
* reads after {@link GitWorktrees#add}, rather than checking only a URL formatting helper.
|
||||
*/
|
||||
@Test
|
||||
void aProvisionedWorktreeUsesACleanHttpsOrigin(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://synthetic-test-token@git.ltms.dev/akb/kb.git");
|
||||
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "fleetd-157-safe-origin", "HEAD");
|
||||
|
||||
Path worktree = Path.of(wt);
|
||||
String origin = gitOutput(worktree, "config", "--get", "remote.origin.url").trim();
|
||||
assertEquals("https://git.ltms.dev/akb/kb.git", origin);
|
||||
assertFalse(origin.contains("synthetic-test-token"), "provisioned worktree kept user info");
|
||||
assertFalse(gitOutput(worktree, "remote", "-v").contains("synthetic-test-token"),
|
||||
"git remote -v exposed user info");
|
||||
assertFalse(gitOutput(worktree, "config", "--list").contains("synthetic-test-token"),
|
||||
"git config --list exposed user info");
|
||||
String helper = gitOutput(worktree, "config", "--worktree", "--get", "credential.helper");
|
||||
assertTrue(helper.contains("WORKER_GITEA_TOKEN"), "credential helper does not read the member environment");
|
||||
assertFalse(helper.contains("synthetic-test-token"), "credential helper stored user info");
|
||||
}
|
||||
|
||||
@Test
|
||||
void worktreeCredentialHelperCompletesWithoutUsingAnInheritedHelper(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "fleetd-157-helper", "HEAD");
|
||||
Path globalConfig = tmp.resolve("global.gitconfig");
|
||||
Files.writeString(globalConfig, """
|
||||
[credential]
|
||||
helper = !f() { printf 'username=%s\\npassword=%s\\n\\n' operator operator-secret; }; f
|
||||
""");
|
||||
|
||||
ProcessBuilder pb = new ProcessBuilder("git", "credential", "fill")
|
||||
.directory(Path.of(wt).toFile()).redirectErrorStream(true);
|
||||
pb.environment().put("GIT_CONFIG_GLOBAL", globalConfig.toString());
|
||||
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
|
||||
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
|
||||
pb.environment().put("WORKER_GITEA_TOKEN", "synthetic-worker-value");
|
||||
Process p = pb.start();
|
||||
p.getOutputStream().write("protocol=https\nhost=git.ltms.dev\n\n".getBytes(StandardCharsets.UTF_8));
|
||||
p.getOutputStream().close();
|
||||
String credential = new String(p.getInputStream().readAllBytes(), StandardCharsets.UTF_8);
|
||||
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git credential fill timed out");
|
||||
assertEquals(0, p.exitValue(), "git credential fill failed");
|
||||
assertTrue(credential.contains("username=git"), "helper did not return its fixed username");
|
||||
assertTrue(credential.contains("password=synthetic-worker-value"),
|
||||
"helper did not return the worker token as the password");
|
||||
assertFalse(credential.contains("operator-secret"), "Git used the inherited global helper");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #157 follow-up. An SSH origin never consults {@code credential.helper} — the environment
|
||||
* credential helper set by {@link GitWorktrees#add} is therefore useless when a member's origin is
|
||||
* SSH, which is exactly this repo's shape. The worktree must instead get a worktree-scoped
|
||||
* {@code url.<https>.insteadOf <ssh>} rewrite so both fetch and push resolve to HTTPS, while the
|
||||
* parent checkout — sharing the same repo-level origin config — must resolve the original SSH URL
|
||||
* completely unchanged. The host/port here (a synthetic {@code forge.example.test:2222}, not
|
||||
* {@code git.ltms.dev}) proves the rewrite is derived from the origin, not a hardcoded constant.
|
||||
*/
|
||||
@Test
|
||||
void aProvisionedWorktreeRewritesAnSshOriginToHttpsWorktreeScoped(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "ssh://git@forge.example.test:2222/acme/proj.git");
|
||||
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "fleetd-157-ssh-rewrite", "HEAD");
|
||||
|
||||
Path worktree = Path.of(wt);
|
||||
assertEquals("https://forge.example.test/acme/proj.git",
|
||||
gitOutput(worktree, "remote", "get-url", "origin").trim(),
|
||||
"worktree fetch URL was not rewritten to HTTPS");
|
||||
assertEquals("https://forge.example.test/acme/proj.git",
|
||||
gitOutput(worktree, "remote", "get-url", "--push", "origin").trim(),
|
||||
"worktree push URL was not rewritten to HTTPS");
|
||||
// The raw config value is unchanged — only the resolved URL is rewritten, via insteadOf.
|
||||
assertEquals("ssh://git@forge.example.test:2222/acme/proj.git",
|
||||
gitOutput(worktree, "config", "--get", "remote.origin.url").trim());
|
||||
|
||||
assertEquals("ssh://git@forge.example.test:2222/acme/proj.git",
|
||||
gitOutput(repo, "remote", "get-url", "origin").trim(),
|
||||
"the parent checkout's fetch URL must be untouched");
|
||||
assertEquals("ssh://git@forge.example.test:2222/acme/proj.git",
|
||||
gitOutput(repo, "remote", "get-url", "--push", "origin").trim(),
|
||||
"the parent checkout's push URL must be untouched");
|
||||
}
|
||||
|
||||
/** An origin already on HTTPS is left alone — the environment credential helper already covers it. */
|
||||
@Test
|
||||
void aProvisionedWorktreeLeavesAnHttpsOriginAlone(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
|
||||
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "fleetd-157-https-noop", "HEAD");
|
||||
|
||||
Path worktree = Path.of(wt);
|
||||
assertEquals("https://git.ltms.dev/akb/kb.git",
|
||||
gitOutput(worktree, "remote", "get-url", "origin").trim());
|
||||
assertEquals(1, exitCode("git", "-C", wt, "config", "--worktree", "--get-regexp", "^url\\."),
|
||||
"no url.*.insteadOf rewrite should be added for an already-HTTPS origin");
|
||||
}
|
||||
|
||||
/** Test-local exit-code probe, mirroring {@link GitWorktrees#exitCode} for an assertion the
|
||||
* production class does not expose. */
|
||||
private static int exitCode(String... command) throws Exception {
|
||||
ProcessBuilder pb = new ProcessBuilder(command).redirectErrorStream(true);
|
||||
pb.environment().put("GIT_CONFIG_GLOBAL", "/dev/null");
|
||||
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
|
||||
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
|
||||
Process p = pb.start();
|
||||
p.getInputStream().readAllBytes();
|
||||
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "command timed out: " + String.join(" ", command));
|
||||
return p.exitValue();
|
||||
}
|
||||
|
||||
@Test
|
||||
void provisioningRefusesAWorktreeWhoseOriginStillHasHttpsUserInfo(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
|
||||
GitWorktrees worktrees = new GitWorktrees(tmp.resolve("wts").toString(), worktreePath -> {
|
||||
try {
|
||||
git(Path.of(worktreePath), "remote", "set-url", "origin",
|
||||
"https://synthetic-test-token@git.ltms.dev/akb/kb.git");
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
});
|
||||
|
||||
WorktreeException error = assertThrows(WorktreeException.class,
|
||||
() -> worktrees.add(repo.toString(), "fleetd-157-refuse-origin", "HEAD"));
|
||||
assertEquals("worktree origin contains HTTPS user info; refusing provision", error.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