Compare commits
9 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 5da912ae8e | |||
| f6c150e99a | |||
| 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
|
||||
```
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -449,7 +449,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)) {
|
||||
|
||||
@@ -642,6 +642,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); }
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -1000,8 +1000,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
*
|
||||
* <p>CB-633: under {@code policy: allow-list} this overlay is NOT the control anymore — it is
|
||||
* applied before the login shell runs and a sourced file can (and did) undo it. The control is
|
||||
* the ZDOTDIR scrub ({@link #applyEnvironmentAllowListPolicy}); {@code known}/{@code allow}
|
||||
* remain as reporting only via {@link #logCredentialGap}.
|
||||
* the ZDOTDIR scrub ({@link #applyEnvironmentAllowListPolicy}); {@code allow} adds explicit
|
||||
* keeps to its derived base, while {@code known} remains reporting metadata.
|
||||
*/
|
||||
private void applyMemberCredentialPolicy(Map<String, String> workerEnv) {
|
||||
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
|
||||
@@ -1010,8 +1010,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
}
|
||||
if (!creds.isAllowList()) {
|
||||
overlayBlockedCredentials(workerEnv, creds);
|
||||
logCredentialGap(creds, null);
|
||||
}
|
||||
logCredentialGap(creds);
|
||||
}
|
||||
|
||||
/** Put {@link #BLOCKED_CREDENTIAL_SENTINEL} over every blocked name in the pane-creation env map. */
|
||||
@@ -1034,10 +1034,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* 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.
|
||||
* MemberEnvAllowList#derive}) UNIONed with {@code memberCredentials.allow} and the exact keys of
|
||||
* THIS launch's env map. This lets the operator keep inherited variables by name without putting
|
||||
* their values in config, and names the daemon itself injects survive its own control. {@code
|
||||
* SSH_AUTH_SOCK} is EXCLUDED from that union and added ONLY when {@code
|
||||
* memberCredentials.sshAuthSock: allow} — {@code sshAuthSock} is the sole control for that one
|
||||
* name, so it cannot be widened back in by naming it on {@code allow:} too; a conflicting entry
|
||||
* there gets a WARN ({@link #warnSshAuthSockConflict}) and loses. By default {@code
|
||||
* SSH_AUTH_SOCK} is absent, so the scrub blanks it like any other non-derived name.
|
||||
*/
|
||||
private Path applyEnvironmentAllowListPolicy(FleetConfig.Profile cfg, Launch launch) {
|
||||
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
|
||||
@@ -1053,17 +1057,31 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
// it), but strictly better than nothing.
|
||||
warnNonZsh(loginShell);
|
||||
overlayBlockedCredentials(launch.env(), creds);
|
||||
logCredentialGap(creds);
|
||||
// Nothing is scrubbed on this path — it is the CB-596 overlay only, which a sourced
|
||||
// startup file can undo. Report the gap with the WARN wording (an inherited-unblocked
|
||||
// exposure), not the allow-list INFO wording, which would falsely claim a scrub blanks it.
|
||||
logCredentialGap(creds, null);
|
||||
return null;
|
||||
}
|
||||
Set<String> allowed = new java.util.TreeSet<>(MemberEnvAllowList.derive(profiles.values()));
|
||||
for (String name : creds.allow()) {
|
||||
// SSH_AUTH_SOCK is governed ONLY by sshAuthSock: below, never by allow: — an operator
|
||||
// who writes both `sshAuthSock: block` and `SSH_AUTH_SOCK` on `allow:` means the block
|
||||
// to win, not to be silently widened away by the second control.
|
||||
if (!SSH_AUTH_SOCK.equals(name)) {
|
||||
allowed.add(name);
|
||||
}
|
||||
}
|
||||
if (creds.sshAuthSockAllowed()) {
|
||||
allowed.add(SSH_AUTH_SOCK);
|
||||
} else if (creds.allow().contains(SSH_AUTH_SOCK)) {
|
||||
warnSshAuthSockConflict();
|
||||
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
|
||||
allowed.addAll(launch.env().keySet());
|
||||
logCredentialGap(creds, allowed);
|
||||
Path dir = EnvAllowListScrub.generate(Path.of(System.getProperty("java.io.tmpdir")), allowed);
|
||||
launch.env().put("ZDOTDIR", dir.toAbsolutePath().toString());
|
||||
log.info("memberCredentials policy=allow-list: profile={} generated ZDOTDIR {} — derived "
|
||||
log.info("memberCredentials policy=allow-list: profile={} generated ZDOTDIR {} — effective "
|
||||
+ "allow-list holds {} name(s); the pane reports allowed N of M at release",
|
||||
cfg.profile(), dir.getFileName(), allowed.size());
|
||||
return dir;
|
||||
@@ -1072,6 +1090,23 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
/** The operator ssh-agent handle — kept ONLY by explicit config decision, never by default. */
|
||||
private static final String SSH_AUTH_SOCK = "SSH_AUTH_SOCK";
|
||||
|
||||
/** Guards {@link #warnSshAuthSockConflict} to one WARN per launcher instance, not one per spawn. */
|
||||
private final AtomicBoolean sshAuthSockConflictWarned = new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* Defect fix: {@code SSH_AUTH_SOCK} named on {@code memberCredentials.allow} while
|
||||
* {@code sshAuthSock} is not {@code allow} is a conflict between the two controls — say which
|
||||
* one wins, once per launcher instance, without printing any value.
|
||||
*/
|
||||
private void warnSshAuthSockConflict() {
|
||||
if (sshAuthSockConflictWarned.compareAndSet(false, true)) {
|
||||
log.warn("memberCredentials policy=allow-list: SSH_AUTH_SOCK is named on allow: but "
|
||||
+ "sshAuthSock is not 'allow' — sshAuthSock: block wins and SSH_AUTH_SOCK "
|
||||
+ "stays blocked. Remove it from allow: or set sshAuthSock: allow if the "
|
||||
+ "member should keep it.");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633: a non-zsh login shell means the allow-list control CANNOT run — say so once per
|
||||
* launcher instance, naming the shell, instead of failing silently.
|
||||
@@ -1131,29 +1166,59 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
private static final Pattern CREDENTIAL_SHAPED_NAME =
|
||||
Pattern.compile("(?i).*(TOKEN|SECRET|_KEY|APIKEY|PASSWORD|CREDENTIAL|AUTH).*");
|
||||
|
||||
/** Guards {@link #logCredentialGap} to one WARN per launcher instance, not one per spawn. */
|
||||
private final AtomicBoolean credentialGapLogged = new AtomicBoolean();
|
||||
/**
|
||||
* Guards {@link #logCredentialGap}'s two report kinds separately — one flag per wording, not
|
||||
* one shared flag — so an early INFO ("will be blanked by the scrub") on one spawn can never
|
||||
* suppress a later, more serious WARN ("inherited UNBLOCKED") on another. {@code
|
||||
* memberCredentials} is live-reloadable, so the policy really can change between spawns on the
|
||||
* same launcher instance.
|
||||
*/
|
||||
private final AtomicBoolean allowListGapLogged = new AtomicBoolean();
|
||||
|
||||
/** See {@link #allowListGapLogged} — the WARN-wording counterpart. */
|
||||
private final AtomicBoolean unprotectedGapLogged = new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* CB-596 criterion 4: a credential-shaped host env var name on neither {@code known} nor
|
||||
* {@code allow} is not silently allowed — it is reported. {@link #hostEnvNames} enumerates the
|
||||
* CB-596 criterion 4: report credential-shaped host env var names that the active policy does
|
||||
* not classify. Callers that pass a non-null {@code effectiveAllowed} are reporting against a
|
||||
* kept-name set a scrub will actually enforce — an absent name gets an INFO stating that the
|
||||
* scrub will blank it, because that is the control working rather than an exposure. Every
|
||||
* other caller — deny-by-default, and an allow-list spawn whose scrub could not run (non-zsh
|
||||
* login shell, falls back to the overlay) — passes {@code null} and gets the original WARN,
|
||||
* because in both cases a gap name is inherited unblocked. {@link #hostEnvNames} enumerates the
|
||||
* daemon's own environment (see that field's javadoc for why the daemon's env is read rather
|
||||
* than the spawned pane's, which the daemon has no channel to inspect at spawn time); this logs
|
||||
* every such NAME, at WARN, at most once per launcher instance — never a value, a prefix of a
|
||||
* value, or a hash of a value, so the log itself cannot leak anything.
|
||||
* than the spawned pane's). This logs names only, at most once per launcher instance per report
|
||||
* kind — never a value, a prefix of a value, or a hash of a value, so the log itself cannot
|
||||
* leak anything.
|
||||
*/
|
||||
private void logCredentialGap(FleetConfig.MemberCredentials creds) {
|
||||
Set<String> covered = new HashSet<>(creds.known());
|
||||
covered.addAll(creds.allow());
|
||||
private void logCredentialGap(FleetConfig.MemberCredentials creds, Set<String> effectiveAllowed) {
|
||||
Set<String> covered = effectiveAllowed == null
|
||||
? new HashSet<>(creds.known()) : effectiveAllowed;
|
||||
if (effectiveAllowed == null) {
|
||||
covered.addAll(creds.allow());
|
||||
}
|
||||
List<String> gap = hostEnvNames.get().stream()
|
||||
.filter(name -> CREDENTIAL_SHAPED_NAME.matcher(name).matches())
|
||||
.filter(name -> !covered.contains(name))
|
||||
.filter(name -> effectiveAllowed == null
|
||||
? !covered.contains(name) : !MemberEnvAllowList.keeps(covered, name))
|
||||
.sorted()
|
||||
.toList();
|
||||
if (gap.isEmpty()) {
|
||||
return;
|
||||
}
|
||||
if (credentialGapLogged.compareAndSet(false, true)) {
|
||||
// effectiveAllowed != null means a derived allow-list set was actually computed (the
|
||||
// scrub will run against it) — that is the INFO case. A null effectiveAllowed means no
|
||||
// scrub protects this spawn (deny-by-default, or an allow-list spawn that fell back to
|
||||
// the overlay-only path because the login shell is not zsh) — that is the WARN case.
|
||||
// See #allowListGapLogged for why the two kinds use separate guards.
|
||||
if (effectiveAllowed != null) {
|
||||
if (allowListGapLogged.compareAndSet(false, true)) {
|
||||
log.info("memberCredentials allow-list: {} credential-shaped host env var name(s) "
|
||||
+ "are not kept and will be blanked by the scrub — {}. Add any name "
|
||||
+ "a member legitimately needs to memberCredentials.allow.",
|
||||
gap.size(), gap);
|
||||
}
|
||||
} else if (unprotectedGapLogged.compareAndSet(false, true)) {
|
||||
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||
+ "known: nor allow: — every member pane inherits them UNBLOCKED — {}. "
|
||||
+ "Add each to memberCredentials.known (blocked by default) or .allow "
|
||||
|
||||
@@ -7,14 +7,16 @@ import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
|
||||
/**
|
||||
* CB-633: the set of environment variable NAMES a spawned member is allowed to keep under
|
||||
* {@code memberCredentials.policy: allow-list} — DERIVED from what the launcher itself injects,
|
||||
* never hand-typed.
|
||||
* CB-633: the base set of environment variable NAMES a spawned member is allowed to keep under
|
||||
* {@code memberCredentials.policy: allow-list}. This class derives the base from what the launcher
|
||||
* itself injects. The launcher then adds the operator's explicitly named {@code
|
||||
* memberCredentials.allow} keeps and the exact keys of this launch's env map.
|
||||
*
|
||||
* <p>A hand-typed allow-list is the defect this class exists to prevent: a name an operator forgets
|
||||
* to type is a credential that passes through to every member, and a profile added to config later
|
||||
* would silently break spawns whose scrub did not know its names. Derivation closes both ends. The
|
||||
* kept-name set is the union of:
|
||||
* <p>A hand-typed list that <em>replaces</em> derivation is the defect this class exists to prevent:
|
||||
* a profile added to config later would silently break spawns whose scrub did not know its names.
|
||||
* An explicit list that only adds keeps is safe because adding a name can only widen the set; it
|
||||
* cannot make another spawn lose a name that derivation already kept. The derived base is the union
|
||||
* of:
|
||||
*
|
||||
* <ul>
|
||||
* <li>every configured {@link FleetConfig.Profile profile}'s {@code gitTokenEnv},
|
||||
@@ -26,11 +28,11 @@ import java.util.TreeSet;
|
||||
* or the agent binary genuinely needs to function.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>Because the union spans EVERY profile (not just the one spawning), adding a new profile can
|
||||
* only ever widen the list — it cannot break another spawn's scrub. And because the launcher also
|
||||
* unions in the exact keys of each spawn's own env map at generation time (see {@code
|
||||
* HerdrPeerLauncher}), anything the daemon deliberately injects for THIS spawn survives its own
|
||||
* control.
|
||||
* <p>Because the base spans EVERY profile (not just the one spawning), adding a new profile can only
|
||||
* ever widen the list — it cannot break another spawn's scrub. The launcher then adds the operator's
|
||||
* explicit keeps and the exact keys of each spawn's own env map at generation time (see {@code
|
||||
* HerdrPeerLauncher}), so member binaries can keep named inherited variables without putting their
|
||||
* secret values in config, and anything the daemon injects for THIS spawn survives its own control.
|
||||
*
|
||||
* <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
|
||||
|
||||
@@ -204,6 +204,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());
|
||||
@@ -271,6 +290,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 +356,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 +412,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 +520,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 +543,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) {
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1025,7 +1025,7 @@ class ClaudeCodeLauncherTest {
|
||||
* {@code System.getenv()} so the test is deterministic.
|
||||
*/
|
||||
@Test
|
||||
void aCredentialShapedNameOnNeitherListIsLoggedAsAGap() {
|
||||
void denyByDefaultKeepsTheExactCredentialGapWarn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
@@ -1045,16 +1045,104 @@ class ClaudeCodeLauncherTest {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
|
||||
String expected = "memberCredentials gap: 1 credential-shaped env var name(s) are on neither "
|
||||
+ "known: nor allow: — every member pane inherits them UNBLOCKED — "
|
||||
+ "[A_BRAND_NEW_SECRET_TOKEN]. Add each to memberCredentials.known (blocked by default) "
|
||||
+ "or .allow (if a member legitimately needs it).";
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getFormattedMessage().contains("memberCredentials gap")
|
||||
&& e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN")),
|
||||
"the gap must name the unrecognized credential-shaped var, never a value");
|
||||
e.getLevel() == ch.qos.logback.classic.Level.WARN
|
||||
&& e.getFormattedMessage().equals(expected)),
|
||||
"deny-by-default must keep the exact existing WARN text");
|
||||
assertNotNull(startEnv(herdr).get("GITEA_ACCESS_TOKEN"),
|
||||
"deny-by-default must still overlay a known blocked credential");
|
||||
assertFalse(appender.list.stream().anyMatch(e -> e.getFormattedMessage().contains("PATH")),
|
||||
"PATH/HOME are not credential-shaped and must not be reported as a gap");
|
||||
assertFalse(appender.list.stream().anyMatch(e -> e.getFormattedMessage().contains("AI_GATEWAY_TOKEN")),
|
||||
"a name already on allow: is covered, not a gap");
|
||||
}
|
||||
|
||||
@Test
|
||||
void allowListReportsAnUnkeptCredentialAsBlankedWithoutSayingUnblocked() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
FleetConfig.MemberCredentials creds = new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null);
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
name -> "SHELL".equals(name) ? "/bin/zsh" : null,
|
||||
0, System::currentTimeMillis, () -> {}, null, () -> creds,
|
||||
() -> Set.of("PATH", "HOME", "A_BRAND_NEW_SECRET_TOKEN"));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
ch.qos.logback.classic.Level original = logger.getLevel();
|
||||
logger.setLevel(ch.qos.logback.classic.Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
svc.spawn();
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
ILoggingEvent gap = appender.list.stream()
|
||||
.filter(e -> e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN"))
|
||||
.findFirst().orElseThrow(() -> new AssertionError("allow-list must report the blanked name"));
|
||||
assertEquals(ch.qos.logback.classic.Level.INFO, gap.getLevel(),
|
||||
"a scrubbed credential name is the allow-list control working, not a WARN");
|
||||
assertTrue(gap.getFormattedMessage().contains("will be blanked by the scrub"));
|
||||
assertFalse(gap.getFormattedMessage().contains("UNBLOCKED"),
|
||||
"the allow-list report must not claim that a blanked name is exposed");
|
||||
}
|
||||
|
||||
/**
|
||||
* Defect fix: {@code memberCredentials} is live-reloadable, so the policy really can change
|
||||
* between spawns on the SAME launcher instance. An earlier INFO gap report (allow-list) must
|
||||
* not permanently suppress a later, more serious WARN gap report (deny-by-default) — the two
|
||||
* report kinds need separate guards, not one shared flag.
|
||||
*/
|
||||
@Test
|
||||
void anEarlierInfoGapDoesNotSuppressALaterWarnGapOnTheSameLauncher() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("claude"), "tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
AtomicReference<FleetConfig.MemberCredentials> creds = new AtomicReference<>(
|
||||
new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null));
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
name -> "SHELL".equals(name) ? "/bin/zsh" : null,
|
||||
0, System::currentTimeMillis, () -> {}, null, creds::get,
|
||||
() -> Set.of("PATH", "HOME", "A_BRAND_NEW_SECRET_TOKEN"));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
ch.qos.logback.classic.Level original = logger.getLevel();
|
||||
logger.setLevel(ch.qos.logback.classic.Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
svc.spawn(); // allow-list: logs the INFO gap report first
|
||||
creds.set(new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT, List.of(), List.of(), null));
|
||||
svc.spawn(); // deny-by-default: must still log its own WARN gap report
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getLevel() == ch.qos.logback.classic.Level.WARN
|
||||
&& e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN")
|
||||
&& e.getFormattedMessage().contains("UNBLOCKED")),
|
||||
"a prior INFO gap report must not suppress a later WARN gap report on the same "
|
||||
+ "launcher instance");
|
||||
}
|
||||
|
||||
/**
|
||||
* The half of CB-592 that can actually survive the pane's login shell. BRIDGED_MEMBER is a name
|
||||
* secrets.sh never exports, so nothing overwrites it — measured: GITEA_TOKEN is injected the
|
||||
|
||||
+160
-2
@@ -1,5 +1,8 @@
|
||||
package dev.ltms.fleet.member;
|
||||
|
||||
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,7 +11,10 @@ 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.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.HashMap;
|
||||
@@ -21,6 +27,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
import static org.junit.jupiter.api.Assumptions.assumeTrue;
|
||||
|
||||
/**
|
||||
* CB-633: proves the allow-list scrub is actually WIRED INTO the spawn path — not merely that its
|
||||
@@ -43,6 +50,8 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
|
||||
/** A name the daemon itself injects — it must survive its own scrub, so it must be allowed. */
|
||||
private static final String INJECTED = "ANTHROPIC_BASE_URL";
|
||||
private static final String OPERATOR_KEEP = "CONTEXT7_TOKEN";
|
||||
private static final String UNLISTED_CREDENTIAL = "A_BRAND_NEW_SECRET_TOKEN";
|
||||
|
||||
@Test
|
||||
void spawningUnderAllowListPolicyGivesThePaneAGeneratedZdotdir() {
|
||||
@@ -74,6 +83,43 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
+ "), or the daemon's own configuration is blanked by its own control");
|
||||
}
|
||||
|
||||
/** Explicit operator keeps widen the derived base, but all other exported names stay blocked. */
|
||||
@Test
|
||||
void operatorKeepSurvivesWhileAnUnlistedCredentialIsBlanked(@TempDir Path home) throws Exception {
|
||||
Path zsh = Path.of("/bin/zsh");
|
||||
assumeTrue(Files.isExecutable(zsh), "/bin/zsh not present — nothing to prove here");
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(List.of(OPERATOR_KEEP)));
|
||||
|
||||
var spawned = launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
try {
|
||||
Path zdotdir = Path.of(launcher.env.get("ZDOTDIR"));
|
||||
ProcessBuilder probe = new ProcessBuilder(zsh.toString(), "-i", "-c",
|
||||
"[[ -n \"$" + OPERATOR_KEEP + "\" ]] && print -r -- kept || print -r -- missing; "
|
||||
+ "[[ -z \"$" + UNLISTED_CREDENTIAL
|
||||
+ "\" ]] && print -r -- blanked || print -r -- leaked");
|
||||
probe.environment().clear();
|
||||
probe.environment().putAll(Map.of(
|
||||
"HOME", home.toString(),
|
||||
"PATH", "/usr/bin:/bin",
|
||||
"SHELL", zsh.toString(),
|
||||
"ZDOTDIR", zdotdir.toString(),
|
||||
OPERATOR_KEEP, "needed-by-member-binary",
|
||||
UNLISTED_CREDENTIAL, "must-not-survive"));
|
||||
probe.redirectError(ProcessBuilder.Redirect.DISCARD);
|
||||
|
||||
Process process = probe.start();
|
||||
String output = new String(process.getInputStream().readAllBytes(), StandardCharsets.UTF_8);
|
||||
assertTrue(process.waitFor(60, java.util.concurrent.TimeUnit.SECONDS),
|
||||
"the zsh scrub probe did not exit within 60 seconds");
|
||||
assertEquals(0, process.exitValue(), "the zsh scrub probe must exit cleanly");
|
||||
assertEquals("kept\nblanked\n", output,
|
||||
"memberCredentials.allow must add a keep, while an unlisted credential stays blanked");
|
||||
} finally {
|
||||
launcher.stop(spawned.id());
|
||||
}
|
||||
}
|
||||
|
||||
/** The default policy must not generate anything — an upgrade changes nothing until asked. */
|
||||
@Test
|
||||
void spawningUnderTheDefaultPolicyGeneratesNoZdotdir() {
|
||||
@@ -103,9 +149,116 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
"bash ignores ZDOTDIR; setting it would be protection theatre");
|
||||
}
|
||||
|
||||
/**
|
||||
* Defect fix: under {@code policy: allow-list} with a non-zsh login shell, nothing is
|
||||
* scrubbed — the overlay-only fallback is the whole protection. The gap must still be
|
||||
* reported, and with the deny-by-default WARN wording ("inherited UNBLOCKED"), never the
|
||||
* allow-list INFO wording ("will be blanked by the scrub"), because nothing is scrubbed here.
|
||||
*/
|
||||
@Test
|
||||
void nonZshAllowListFallbackStillReportsTheGapAsAWarn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/bash",
|
||||
() -> Set.of("PATH", "HOME", UNLISTED_CREDENTIAL));
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
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);
|
||||
}
|
||||
|
||||
ILoggingEvent gap = appender.list.stream()
|
||||
.filter(e -> e.getFormattedMessage().contains(UNLISTED_CREDENTIAL))
|
||||
.findFirst().orElseThrow(() -> new AssertionError(
|
||||
"the non-zsh allow-list fallback must still report the credential gap"));
|
||||
assertEquals(ch.qos.logback.classic.Level.WARN, gap.getLevel(),
|
||||
"nothing is scrubbed on this path, so the report must use the WARN wording, not "
|
||||
+ "the allow-list INFO wording");
|
||||
assertTrue(gap.getFormattedMessage().contains("UNBLOCKED"),
|
||||
"the fallback path scrubs nothing, so the wording must say so");
|
||||
assertFalse(gap.getFormattedMessage().contains("will be blanked by the scrub"),
|
||||
"nothing is scrubbed on this path — that wording would be false here");
|
||||
}
|
||||
|
||||
/**
|
||||
* The live operator config has {@code SSH_AUTH_SOCK} on BOTH {@code allow:} AND
|
||||
* {@code sshAuthSock: block}. Defect fix: {@code sshAuthSock: block} must win — the allow:
|
||||
* entry must not silently widen the derived allow-list to include it.
|
||||
*/
|
||||
@Test
|
||||
void sshAuthSockBlockWinsOverAnAllowListEntry(@TempDir Path home) throws Exception {
|
||||
Path zsh = Path.of("/bin/zsh");
|
||||
assumeTrue(Files.isExecutable(zsh), "/bin/zsh not present");
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Supplier<FleetConfig.MemberCredentials> live = () -> new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST,
|
||||
List.of("AI_GATEWAY_TOKEN", "WORKER_GITEA_TOKEN", "CONTEXT7_TOKEN", "GITEA_HOST",
|
||||
"OPENCODE_AUTOMODE_MODEL", "SSH_AUTH_SOCK", "CLAUDE_CODE_MESSAGING_TOKEN"),
|
||||
List.of(), "block");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, live);
|
||||
var spawned = launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
try {
|
||||
Path zdotdir = Path.of(launcher.env.get("ZDOTDIR"));
|
||||
ProcessBuilder probe = new ProcessBuilder(zsh.toString(), "-i", "-c",
|
||||
"[[ -n \"$SSH_AUTH_SOCK\" ]] && print -r -- LEAKED || print -r -- blanked");
|
||||
probe.environment().clear();
|
||||
probe.environment().putAll(Map.of(
|
||||
"HOME", home.toString(), "PATH", "/usr/bin:/bin", "SHELL", zsh.toString(),
|
||||
"ZDOTDIR", zdotdir.toString(), "SSH_AUTH_SOCK", "/tmp/agent.sock"));
|
||||
probe.redirectError(ProcessBuilder.Redirect.DISCARD);
|
||||
Process pr = probe.start();
|
||||
String out = new String(pr.getInputStream().readAllBytes(), StandardCharsets.UTF_8);
|
||||
pr.waitFor(60, java.util.concurrent.TimeUnit.SECONDS);
|
||||
assertEquals("blanked\n", out, "sshAuthSock: block must win over allow:");
|
||||
} finally {
|
||||
launcher.stop(spawned.id());
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Mirror of {@link #sshAuthSockBlockWinsOverAnAllowListEntry}: with {@code sshAuthSock: allow}
|
||||
* (no conflict with {@code allow:}), {@code SSH_AUTH_SOCK} IS kept. Without this case, a fix
|
||||
* that simply deletes {@code SSH_AUTH_SOCK} everywhere would also pass the block test.
|
||||
*/
|
||||
@Test
|
||||
void sshAuthSockAllowKeepsTheVariable(@TempDir Path home) throws Exception {
|
||||
Path zsh = Path.of("/bin/zsh");
|
||||
assumeTrue(Files.isExecutable(zsh), "/bin/zsh not present");
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Supplier<FleetConfig.MemberCredentials> creds = () -> new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST,
|
||||
List.of("SSH_AUTH_SOCK"), List.of(), "allow");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, creds);
|
||||
var spawned = launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
try {
|
||||
Path zdotdir = Path.of(launcher.env.get("ZDOTDIR"));
|
||||
ProcessBuilder probe = new ProcessBuilder(zsh.toString(), "-i", "-c",
|
||||
"[[ -n \"$SSH_AUTH_SOCK\" ]] && print -r -- kept || print -r -- blanked");
|
||||
probe.environment().clear();
|
||||
probe.environment().putAll(Map.of(
|
||||
"HOME", home.toString(), "PATH", "/usr/bin:/bin", "SHELL", zsh.toString(),
|
||||
"ZDOTDIR", zdotdir.toString(), "SSH_AUTH_SOCK", "/tmp/agent.sock"));
|
||||
probe.redirectError(ProcessBuilder.Redirect.DISCARD);
|
||||
Process pr = probe.start();
|
||||
String out = new String(pr.getInputStream().readAllBytes(), StandardCharsets.UTF_8);
|
||||
pr.waitFor(60, java.util.concurrent.TimeUnit.SECONDS);
|
||||
assertEquals("kept\n", out, "sshAuthSock: allow must keep SSH_AUTH_SOCK");
|
||||
} finally {
|
||||
launcher.stop(spawned.id());
|
||||
}
|
||||
}
|
||||
|
||||
private static Supplier<FleetConfig.MemberCredentials> allowList() {
|
||||
return allowList(List.of());
|
||||
}
|
||||
|
||||
private static Supplier<FleetConfig.MemberCredentials> allowList(List<String> allow) {
|
||||
return () -> new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null);
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, allow, List.of(), null);
|
||||
}
|
||||
|
||||
private static String readAll(Path p) {
|
||||
@@ -134,10 +287,15 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
}
|
||||
|
||||
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell) {
|
||||
this(herdr, creds, shell, null);
|
||||
}
|
||||
|
||||
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
|
||||
|
||||
@@ -1081,4 +1081,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
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user