Compare commits

...

9 Commits

Author SHA1 Message Date
Dai Ha 5da912ae8e CB-633: sshAuthSock: block must win over an SSH_AUTH_SOCK allow: entry, and report the credential gap on the non-zsh fallback path
CI / contract (pull_request) Successful in 1m14s
CI / build (pull_request) Successful in 1m39s
Fixes two defects in PR #174:

1. applyEnvironmentAllowListPolicy unioned memberCredentials.allow directly into the
   derived allow-list, so an operator who wrote both `sshAuthSock: block` and
   SSH_AUTH_SOCK on `allow:` (the live fleetd.yaml shape) got the block silently
   defeated. SSH_AUTH_SOCK is now excluded from that union and governed only by
   sshAuthSock:, with a one-time WARN when the two controls conflict.

2. The non-zsh login-shell fallback branch dropped its logCredentialGap call, so an
   allow-list spawn on a non-zsh shell produced no gap report at all — exactly the
   weakest, overlay-only path that most needs one. The call is restored, and
   logCredentialGap now picks its WARN/INFO wording from whether a scrub-derived
   allow-list set was actually computed (effectiveAllowed != null) rather than from
   the policy alone, so this path correctly gets the WARN ("inherited UNBLOCKED")
   wording instead of the allow-list INFO wording.

3. credentialGapLogged was one shared AtomicBoolean guarding both report kinds;
   split into allowListGapLogged/unprotectedGapLogged so an early INFO on one spawn
   can no longer suppress a later WARN on a live-reloaded policy.

Each fix is proven by reverting it and observing the corresponding test fail, then
restoring it.
2026-08-31 08:44:30 +07:00
Dai Ha f6c150e99a CB-633: honor explicit member credential keeps
CI / build (pull_request) Successful in 1m0s
CI / contract (pull_request) Successful in 1m19s
2026-08-28 05:25:46 +07:00
Dai Ha 7c4170ff6d CB-643: join the message-layer evidence to the health monitor
CI / build (push) Successful in 1m8s
CI / contract (push) Successful in 1m9s
CB-640 published the three message-layer facts and CB-641 wired the herdr
and time ones. This joins them, so every HealthSnapshot field now carries
real evidence and the NOT_YET_OBSERVED placeholder is gone. That constant
is what made 8 of the 9 fault states unreachable, GONE and NEVER_READY
included, which is why CB-580's failTarget never fired.

hasOrphanedDelegation is a true snapshot, but it can read true for one
tick during an ordinary race: an async ticket exists before its virtual
thread reaches rendezvous.open, so for that instant nothing is accepted or
queued behind it. decide maps the field straight to DELEGATION_ORPHANED
with no smoothing, so one racy read would log a fault that clears on the
next tick. The monitor now requires two consecutive observations. That
costs one interval on a real orphan and removes the false positive.

Two tests drive real ticks against a genuinely orphaned ticket (an
unanswered fleet_ask that lapsed back to PENDING), not the seam: one tick
reports nothing, two report once, and a single clean tick in between
resets the streak.

Also correct two config comments. paneProbeIntervalSeconds is parsed and
read by nothing, so its "minimum 60" note promised a floor that does not
exist.

970 tests green.
2026-08-27 22:15:39 +07:00
Dai Ha a6095743f0 Merge CB-641: wire herdr and time evidence into the fleet health monitor 2026-08-27 22:06:50 +07:00
Dai Ha 66e776d178 CB-641: wire herdr health evidence
CI / contract (pull_request) Successful in 1m9s
CI / build (pull_request) Successful in 1m11s
2026-08-27 22:03:20 +07:00
Dai Ha 26f64cba45 Merge CB-640: MessageService evidence accessors for the health monitor
CI / build (push) Successful in 1m4s
CI / contract (push) Successful in 1m28s
2026-08-27 22:02:36 +07:00
Dai Ha 51047848f1 CB-642: make the fleets-status redaction global — a non-global sed leaks a second URI on the same line
CI / contract (push) Successful in 41s
CI / build (push) Successful in 1m29s
2026-08-27 22:01:15 +07:00
Dai Ha d8c0b657e8 Merge CB-642: /fleets-status skill — multi-fleet status over the shared LavinMQ 2026-08-27 22:00:49 +07:00
Dai Ha 312c0584ce CB-640: add MessageService message-layer health evidence accessors
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m11s
hasQueuedDelivery/hasStrandedReply/hasOrphanedDelegation surface three of the
message-layer facts FleetHealthMonitor needs but currently hardcodes to
NOT_YET_OBSERVED. Additive only — no existing public method's signature or
behavior changes.
2026-08-27 21:58:21 +07:00
12 changed files with 937 additions and 64 deletions
+7 -2
View File
@@ -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
```
+8 -9
View File
@@ -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
@@ -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
}
}