Compare commits
28 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 31d5516991 | |||
| 5ba05d0bdb | |||
| 25726a5ae7 | |||
| 11c3ff67b6 | |||
| d867c87100 | |||
| b0c4cedfab | |||
| 85417d5215 | |||
| 4accc746bd | |||
| 21c4c8cbef | |||
| 430f5b0dae | |||
| 42731833d0 | |||
| 65acf066ad | |||
| ee8f570fd7 | |||
| fa3f910d44 | |||
| 46ac6e4e38 | |||
| 3bfa82839b | |||
| 65ccf2e4ad | |||
| 7a3b27f76f | |||
| 82e7be564c | |||
| c5e24197bf | |||
| bcb402b688 | |||
| 7c4170ff6d | |||
| a6095743f0 | |||
| 66e776d178 | |||
| 26f64cba45 | |||
| 51047848f1 | |||
| d8c0b657e8 | |||
| 312c0584ce |
@@ -27,9 +27,14 @@ can make every broker probe look empty.
|
||||
reaches the report:
|
||||
|
||||
```bash
|
||||
sed -E 's#://[^@]*@#://<redacted>@#'
|
||||
sed -E 's#://[^@]*@#://<redacted>@#g'
|
||||
```
|
||||
|
||||
**The `g` flag is not optional.** Without it `sed` replaces only the first match on each line, so a
|
||||
line carrying two URIs leaks the second one. `scripts/redeploy-fleetd.sh --check` prints lines like
|
||||
that. Checked on 2026-08-27: without `g`, `amqp://u1:p1@h1/mac and http://u2:p2@h2:15672/api`
|
||||
redacts the first pair and prints `u2:p2` in the clear.
|
||||
|
||||
Keep `pipefail` on when applying that filter. Otherwise the filter can hide a failed probe. Apply
|
||||
the same no-print rule to the management password below, even though it is not in an AMQP URI.
|
||||
|
||||
@@ -44,7 +49,7 @@ Do not copy those checks into new shell code. The script reads `LAVINMQ_URI`, so
|
||||
```bash
|
||||
set -o pipefail
|
||||
scripts/redeploy-fleetd.sh --check 2>&1 \
|
||||
| sed -E 's#://[^@]*@#://<redacted>@#'
|
||||
| sed -E 's#://[^@]*@#://<redacted>@#g'
|
||||
git rev-parse HEAD
|
||||
```
|
||||
|
||||
@@ -108,8 +113,19 @@ PY
|
||||
|
||||
**What this tier cannot see:** it proves facts only about the Mac daemon at `127.0.0.1:8765`.
|
||||
It cannot show the fleet01 daemon, broker queue depth, or broker consumers. The fleet01 REST service
|
||||
at `10.10.20.13:8765` is not reachable from the Mac, and SSH as `dai.ha@10.10.20.13` is denied.
|
||||
Say this in the report rather than omitting fleet01.
|
||||
at `10.10.20.13:8765` is not reachable from the Mac. Say this in the report rather than omitting
|
||||
fleet01.
|
||||
|
||||
**But fleet01 IS reachable over SSH — checked 2026-08-28.** An older version of this line said SSH
|
||||
was denied. That is true only for the user `dai.ha`. The host alias `fleet01` maps to user `ltms`,
|
||||
and `ssh fleet01` works with key auth:
|
||||
|
||||
```bash
|
||||
ssh -o BatchMode=yes -o ConnectTimeout=6 fleet01 'echo $(id -un)@$(hostname)'
|
||||
```
|
||||
|
||||
So fleet01's daemon PID, uptime, jar and `/healthz` **can** be reported — over SSH, not over REST.
|
||||
Do that rather than writing `not reachable`. `ltms` also has passwordless sudo there.
|
||||
|
||||
## 3. Tier 2 — the shared broker (run when management access exists)
|
||||
|
||||
|
||||
@@ -93,21 +93,20 @@ bind:
|
||||
# intervalSeconds → how often a tick runs (default 30). ENFORCED floor of 15: the code computes
|
||||
# Math.max(15, intervalSeconds), so a lower value is silently raised, not
|
||||
# rejected.
|
||||
# workingSuspectAfterSeconds, paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by
|
||||
# anything — the dormant monitor only consumes intervalSeconds today (CB-573
|
||||
# shipped ahead of the evidence publishers these two knobs are for). Setting
|
||||
# them changes nothing right now, and no minimum is enforced on either, because
|
||||
# nothing reads them to enforce one. They exist so a later build can start
|
||||
# honouring them without another config-shape change.
|
||||
# workingSuspectAfterSeconds → age before a BUSY member is suspected of a stall (default 600).
|
||||
# ENFORCED floor of 300: a lower value is silently raised.
|
||||
# paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by anything. Setting it changes
|
||||
# nothing right now. It exists so a later build can start honouring it without
|
||||
# another config-shape change.
|
||||
# notifications.mode → "webhook" flips what fleet_list REPORTS (healthCoverage: "full" instead
|
||||
# of "detection-only") — it does NOT make fleetd send any webhook call; no
|
||||
# delivery mechanism is implemented yet. Any other value, or omitting the
|
||||
# block, reports "detection-only".
|
||||
# health:
|
||||
# enabled: true
|
||||
# intervalSeconds: 30
|
||||
# workingSuspectAfterSeconds: 600
|
||||
# paneProbeIntervalSeconds: 60
|
||||
# intervalSeconds: 30 # floor 15
|
||||
# workingSuspectAfterSeconds: 600 # floor 300 — how long BUSY with no activity means STALL_SUSPECTED
|
||||
# paneProbeIntervalSeconds: 60 # parsed, but nothing reads it yet — changing it changes nothing
|
||||
# notifications:
|
||||
# mode: disabled
|
||||
|
||||
|
||||
@@ -380,7 +380,8 @@ public final class Fleetd {
|
||||
sessions.onTurnFailed(target);
|
||||
}
|
||||
};
|
||||
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
|
||||
Predicate<String> deliverable = deliverableTo(presence, leads);
|
||||
Injector injector = new Injector(agents, turnListener, deliverable,
|
||||
presence::forget);
|
||||
StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS);
|
||||
poller.start();
|
||||
@@ -449,7 +450,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)) {
|
||||
@@ -591,7 +593,7 @@ public final class Fleetd {
|
||||
}));
|
||||
|
||||
Javalin app = new FleetApp(herdr, workers, sessions, messages, presence, mcp.servlet(),
|
||||
callers, metrics).build();
|
||||
callers, metrics, deliverable).build();
|
||||
app.start(cfg.bind().host(), cfg.bind().port());
|
||||
log.info("fleetd listening on {}:{}, herdr socket {}",
|
||||
cfg.bind().host(), cfg.bind().port(), socket);
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -38,6 +38,11 @@ public final class AgentControl {
|
||||
this.herdr = herdr;
|
||||
}
|
||||
|
||||
/** The herdr daemon this control object sends its agent calls to. */
|
||||
public HerdrClient herdr() {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
/** One agent-targeted call, translating a terminal id to its pane id (retrying once fresh). */
|
||||
private JsonNode agentCall(String method, String target, Map<String, Object> extra) {
|
||||
String resolved = resolveTarget(target);
|
||||
|
||||
@@ -6,12 +6,14 @@ import dev.ltms.fleet.msg.TurnToken;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.regex.Pattern;
|
||||
|
||||
/**
|
||||
@@ -61,6 +63,17 @@ public final class CompletionResolver implements TurnListener {
|
||||
/** Cap the scraped tail so a long transcript can't return an unbounded blob. */
|
||||
static final int MAX_SCRAPE_CHARS = 4000;
|
||||
|
||||
/**
|
||||
* fleetd#164: the floor below which a {@code BUSY -> DONE} transition cannot be real work. A
|
||||
* backend that rejects a turn outright (e.g. an HTTP 400 from the model, before the worker read
|
||||
* a single file or produced a token) drives the exact same confirmed {@code working -> idle}
|
||||
* transition a genuine completion does — just in about a second instead of the many seconds a
|
||||
* real turn costs. {@link #onTurnComplete} cannot tell those two cases apart from the transition
|
||||
* alone, so a turn that settles inside this floor is treated as a crash signature and resolved
|
||||
* as a failure, never as a (possibly empty) success.
|
||||
*/
|
||||
public static final long MIN_TURN_NANOS = Duration.ofSeconds(2).toNanos();
|
||||
|
||||
private static final String CLIPPED_PANE_TAIL_MARKER =
|
||||
"[Pane tail clipped: member did not call fleet_reply.]";
|
||||
|
||||
@@ -68,6 +81,7 @@ public final class CompletionResolver implements TurnListener {
|
||||
private final Rendezvous rendezvous;
|
||||
private final ExhaustedPatternLookup exhaustedPatterns;
|
||||
private final ExhaustionSink exhaustionSink;
|
||||
private final LongSupplier nowNanos;
|
||||
|
||||
/**
|
||||
* Per-target record of the turn currently in flight: the exact {@link Rendezvous} waiter its
|
||||
@@ -79,8 +93,21 @@ public final class CompletionResolver implements TurnListener {
|
||||
* reference: a completion scrape equal to it means the worker produced no new output (the previous
|
||||
* turn's wind-down sampled as this boundary), so it is suppressed. Overwritten on each delivery;
|
||||
* cleared when the turn resolves. Package-private so tests can capture and replay a specific turn.
|
||||
*
|
||||
* <p>{@code deliveredAtNanos} (fleetd#164) is the {@link #nowNanos} reading taken at delivery —
|
||||
* the other half of the {@link #MIN_TURN_NANOS} floor check, compared against a fresh reading at
|
||||
* resolution time.
|
||||
*/
|
||||
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline) {
|
||||
record InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline, long deliveredAtNanos) {
|
||||
|
||||
/**
|
||||
* Convenience for tests exercising scrape/suppression logic that don't care about turn
|
||||
* timing: back-dates the delivery far enough that {@link #MIN_TURN_NANOS} can never fire.
|
||||
* Not used by production code — {@link #captureBaseline} always records a real reading.
|
||||
*/
|
||||
InFlight(CompletableFuture<Rendezvous.Resolution> waiter, String baseline) {
|
||||
this(waiter, baseline, Long.MIN_VALUE / 2);
|
||||
}
|
||||
}
|
||||
|
||||
private final ConcurrentHashMap<String, InFlight> inFlight = new ConcurrentHashMap<>();
|
||||
@@ -97,10 +124,25 @@ public final class CompletionResolver implements TurnListener {
|
||||
*/
|
||||
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
|
||||
ExhaustionSink exhaustionSink) {
|
||||
this(agents, rendezvous, exhaustedPatterns, exhaustionSink, System::nanoTime);
|
||||
}
|
||||
|
||||
/**
|
||||
* Test constructor with an injectable clock (fleetd#164), matching the {@code LongSupplier}
|
||||
* pattern {@link dev.ltms.fleet.session.SessionManager} and {@link dev.ltms.fleet.msg.MessageService}
|
||||
* already use: lets a test place a turn's delivery and its resolution at an exact, controllable
|
||||
* distance apart around the {@link #MIN_TURN_NANOS} floor, without a real sleep. Public (rather
|
||||
* than package-private like those two) because callers that wire a full {@code MessageService}
|
||||
* fixture — e.g. {@code MessageServiceTest} — construct this resolver directly from another
|
||||
* package.
|
||||
*/
|
||||
public CompletionResolver(AgentControl agents, Rendezvous rendezvous, ExhaustedPatternLookup exhaustedPatterns,
|
||||
ExhaustionSink exhaustionSink, LongSupplier nowNanos) {
|
||||
this.agents = agents;
|
||||
this.rendezvous = rendezvous;
|
||||
this.exhaustedPatterns = Objects.requireNonNull(exhaustedPatterns, "exhaustedPatterns");
|
||||
this.exhaustionSink = Objects.requireNonNull(exhaustionSink, "exhaustionSink");
|
||||
this.nowNanos = Objects.requireNonNull(nowNanos, "nowNanos");
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -130,7 +172,7 @@ public final class CompletionResolver implements TurnListener {
|
||||
baseline = null; // fail open: no baseline ⇒ no suppression
|
||||
log.debug("delivery baseline for {} failed: {}", target, e.getMessage());
|
||||
}
|
||||
inFlight.put(target, new InFlight(waiter, baseline));
|
||||
inFlight.put(target, new InFlight(waiter, baseline, nowNanos.getAsLong()));
|
||||
}
|
||||
|
||||
/** The turn currently baselined for {@code target}, or {@code null} — a test hook for the captureBaseline path. */
|
||||
@@ -176,6 +218,15 @@ public final class CompletionResolver implements TurnListener {
|
||||
inFlight.remove(target, turn);
|
||||
return;
|
||||
}
|
||||
// fleetd#164: a BUSY -> DONE transition inside the floor cannot be real work — it's a crash
|
||||
// signature (e.g. a backend HTTP 400 before the worker did anything), not a fast answer. Fail
|
||||
// it before spending a scrape on the ordinary path; the reason still carries whatever is on
|
||||
// screen, since that is usually the backend's own error.
|
||||
long elapsedNanos = nowNanos.getAsLong() - turn.deliveredAtNanos();
|
||||
if (elapsedNanos < MIN_TURN_NANOS) {
|
||||
fail(target, turn, tooFastReason(target, elapsedNanos));
|
||||
return;
|
||||
}
|
||||
String tail;
|
||||
String assistantBlock = null;
|
||||
int originalLength = 0;
|
||||
@@ -187,20 +238,25 @@ public final class CompletionResolver implements TurnListener {
|
||||
clipped = originalLength > MAX_SCRAPE_CHARS;
|
||||
tail = clip(assistantBlock);
|
||||
} catch (RuntimeException e) {
|
||||
// The worker finished but we couldn't read its screen — still resolve the send so the
|
||||
// caller unblocks; an empty tail beats hanging until the caller's timeout.
|
||||
log.warn("completion scrape for {} failed; resolving with an empty tail: {}",
|
||||
target, e.getMessage());
|
||||
log.warn("completion scrape for {} failed: {}", target, e.getMessage());
|
||||
tail = "";
|
||||
scrapeFailed = true;
|
||||
}
|
||||
// fleetd#164: a scrape nobody could read, and a scrape that read cleanly but produced nothing,
|
||||
// both used to resolve the send as a SUCCESS carrying "" — indistinguishable from a worker that
|
||||
// genuinely finished with nothing to say. That is the defect: fail loudly instead, naming the
|
||||
// member, so a caller (including a lead deciding whether to delegate again) can tell a lost
|
||||
// turn from a real empty answer.
|
||||
if (scrapeFailed || tail.isEmpty()) {
|
||||
fail(target, turn, emptyScrapeReason(target, scrapeFailed));
|
||||
return;
|
||||
}
|
||||
// Misattribution guard (CB-115): if the scrape is byte-identical to the pane content at
|
||||
// delivery, this turn produced no new output — the boundary belongs to the previous turn's
|
||||
// wind-down (common on rapid back-to-back sends). Suppress rather than resolve the send with
|
||||
// a stale answer; the real fleet_reply (or a later genuine completion) resolves it instead.
|
||||
// A scrape that failed to read is exempt — an empty tail there is "couldn't see", not "no change".
|
||||
String baseline = turn.baseline();
|
||||
if (!scrapeFailed && baseline != null && baseline.equals(tail)) {
|
||||
if (baseline != null && baseline.equals(tail)) {
|
||||
log.debug("suppressing misattributed completion for {} (no output change since delivery)",
|
||||
target);
|
||||
return; // keep the in-flight record: a later genuine completion still needs it
|
||||
@@ -208,21 +264,19 @@ public final class CompletionResolver implements TurnListener {
|
||||
// CB-578 stage A: a turn that ended with no fleet_reply AND whose scrape matches the
|
||||
// backend's configured usage-limit pattern is a refusal, not an answer. Classify it as
|
||||
// BACKEND_EXHAUSTED rather than handing the caller a scrape that reads like a real reply.
|
||||
if (!scrapeFailed) {
|
||||
Pattern exhausted = exhaustedPatterns.patternFor(target);
|
||||
String matchedLine = exhausted == null ? null : firstMatchingLine(assistantBlock, exhausted);
|
||||
if (matchedLine != null) {
|
||||
String reason = "backend exhausted (usage limit): " + matchedLine;
|
||||
if (rendezvous.resolveExhausted(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("completion for {} classified BACKEND_EXHAUSTED (no fleet_reply; scrape "
|
||||
+ "matched the profile's exhausted pattern): {}", target, reason);
|
||||
// CB-578 stage B: only on the resolution that actually won the race — a late
|
||||
// duplicate must never quarantine a credential twice for one refusal.
|
||||
exhaustionSink.onExhausted(target, reason);
|
||||
}
|
||||
return;
|
||||
Pattern exhausted = exhaustedPatterns.patternFor(target);
|
||||
String matchedLine = exhausted == null ? null : firstMatchingLine(assistantBlock, exhausted);
|
||||
if (matchedLine != null) {
|
||||
String reason = "backend exhausted (usage limit): " + matchedLine;
|
||||
if (rendezvous.resolveExhausted(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("completion for {} classified BACKEND_EXHAUSTED (no fleet_reply; scrape "
|
||||
+ "matched the profile's exhausted pattern): {}", target, reason);
|
||||
// CB-578 stage B: only on the resolution that actually won the race — a late
|
||||
// duplicate must never quarantine a credential twice for one refusal.
|
||||
exhaustionSink.onExhausted(target, reason);
|
||||
}
|
||||
return;
|
||||
}
|
||||
String completion = clipped ? tail + "\n" + CLIPPED_PANE_TAIL_MARKER : tail;
|
||||
if (rendezvous.resolveCompletion(waiter, completion)) {
|
||||
@@ -272,6 +326,36 @@ public final class CompletionResolver implements TurnListener {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd#164: the failure reason for a turn that settled inside {@link #MIN_TURN_NANOS} — names
|
||||
* the member and both timings, and appends whatever the pane shows (usually the backend's own
|
||||
* error) so the caller sees the cause, not just "it failed".
|
||||
*/
|
||||
private String tooFastReason(String target, long elapsedNanos) {
|
||||
String scrape;
|
||||
try {
|
||||
scrape = clip(agents.read(target, SCRAPE_SOURCE));
|
||||
} catch (RuntimeException e) {
|
||||
scrape = "";
|
||||
}
|
||||
String reason = String.format(
|
||||
"member %s went BUSY -> DONE in %dms (floor %dms) — too fast to be real work, most "
|
||||
+ "likely a backend error before any work started",
|
||||
target, elapsedNanos / 1_000_000, MIN_TURN_NANOS / 1_000_000);
|
||||
return scrape.isBlank() ? reason : reason + ": " + scrape;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd#164: the failure reason for a scrape that produced zero characters — names the member
|
||||
* and says plainly that the turn produced nothing, so a caller (a lead deciding whether to
|
||||
* delegate again included) never mistakes a lost turn for a genuinely empty reply.
|
||||
*/
|
||||
private static String emptyScrapeReason(String target, boolean scrapeFailed) {
|
||||
return "member " + target + " turn completed with an empty scrape (0 chars) — "
|
||||
+ (scrapeFailed ? "its pane could not be read; " : "")
|
||||
+ "treating as a lost turn, not a real answer";
|
||||
}
|
||||
|
||||
/**
|
||||
* The first line of {@code text} matching {@code pattern}, stripped — the CB-578 stage A
|
||||
* evidence carried in a {@code BACKEND_EXHAUSTED} reason so the operator sees the real refusal
|
||||
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.fleet.member;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.peer.PeerHandle;
|
||||
@@ -21,6 +22,7 @@ import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
import java.util.EnumSet;
|
||||
import java.util.HashSet;
|
||||
import java.util.IdentityHashMap;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -44,12 +46,13 @@ import java.util.stream.Collectors;
|
||||
* the single adapter that declares it. Profiles partition cleanly across adapters: the
|
||||
* constructor rejects a name claimed by two.</li>
|
||||
* <li><strong>By pane id</strong> — {@link #stop} routes to the adapter that spawned that pane
|
||||
* (recorded at spawn time). A pane the composite never spawned (only real for a caller that
|
||||
* hand-rolls an id) falls back to the first delegate; teardown is pane-id addressed and
|
||||
* tab cleanup is single-occupant guarded, so it is safe either way.</li>
|
||||
* (recorded at spawn time). A pane the composite never spawned can use the fallback route
|
||||
* in a one-daemon fleet. With more than one herdr daemon, its owner is unknown, so stop refuses
|
||||
* the ambiguous id rather than closing a pane on an arbitrary herdr daemon.</li>
|
||||
* <li><strong>Fleet-wide</strong> — {@link #reapOrphanWorkers} and {@link #capabilities} fan out
|
||||
* and combine. {@link #list} is deduplicated by pane id because every herdr-backed delegate
|
||||
* shares one herdr connection and so reports the same global agent set.</li>
|
||||
* and combine. {@link #list} is deduplicated by (owning daemon, pane id): delegates that share
|
||||
* one herdr connection report the same global agent set, but two daemons can each hold a pane
|
||||
* called {@code w1:p1}, so the daemon has to be part of the key.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>CB-518: an unqualified spawn is routed through a {@link PlacementPolicy}. The default
|
||||
@@ -435,12 +438,32 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
HerdrPeerLauncher d = spawnedBy.remove(id);
|
||||
HerdrPeerLauncher d = spawnedBy.get(id);
|
||||
if (d == null) {
|
||||
log.debug("stop({}) — no recorded owner, routing to the first adapter (pane-addressed)", id);
|
||||
if (herdrDaemonCount() != 1) {
|
||||
throw new IllegalArgumentException("ambiguous paneId '" + id
|
||||
+ "': no owning herdr daemon was recorded");
|
||||
}
|
||||
log.debug("stop({}) — no recorded owner in a single-daemon fleet", id);
|
||||
d = delegates.getFirst();
|
||||
}
|
||||
// Drop the owner record only after the delegate accepted the stop. Removing it first meant a
|
||||
// delegate that threw left the pane alive with its owner forgotten, so the retry fell into
|
||||
// the ambiguous branch above and refused the id for good.
|
||||
d.stop(id);
|
||||
spawnedBy.remove(id);
|
||||
}
|
||||
|
||||
/**
|
||||
* Count actual herdr daemons, not peer adapter kinds. Identity is intentional: separate client
|
||||
* objects may represent different daemons even if a client later implements value equality.
|
||||
*/
|
||||
private int herdrDaemonCount() {
|
||||
Set<HerdrClient> daemons = Collections.newSetFromMap(new IdentityHashMap<>());
|
||||
for (HerdrPeerLauncher delegate : delegates) {
|
||||
daemons.add(delegate.herdr());
|
||||
}
|
||||
return daemons.size();
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -475,14 +498,24 @@ public final class CompositePeerLauncher implements PeerLauncher {
|
||||
return route(profileName).capabilities();
|
||||
}
|
||||
|
||||
/** Every herdr agent, deduplicated by pane id (all delegates share one herdr and list globally). */
|
||||
/**
|
||||
* Every herdr agent, deduplicated by (owning daemon, pane id).
|
||||
*
|
||||
* <p>Delegates that share one {@link HerdrClient} see the same global agent set, so listing them
|
||||
* both would report every agent twice — that is what the dedupe is for. But pane ids are
|
||||
* per-daemon counters, so two daemons really can both hold {@code w1:p1} on different panes.
|
||||
* Keying on the pane id alone would silently drop one of them from {@code fleet_list} and from
|
||||
* every status view built on it. The daemon is part of the key for exactly that reason.
|
||||
*/
|
||||
@Override
|
||||
public List<Agent> list() {
|
||||
Map<HerdrClient, Integer> daemonIndex = new IdentityHashMap<>();
|
||||
Map<String, Agent> byPane = new LinkedHashMap<>();
|
||||
for (HerdrPeerLauncher d : delegates) {
|
||||
int daemon = daemonIndex.computeIfAbsent(d.herdr(), _ -> daemonIndex.size());
|
||||
for (Agent a : d.list()) {
|
||||
if (a.paneId() != null) {
|
||||
byPane.putIfAbsent(a.paneId(), a);
|
||||
byPane.putIfAbsent(daemon + "\u0000" + a.paneId(), a);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package dev.ltms.fleet.member;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.Tab;
|
||||
import dev.ltms.fleet.herdr.Workspace;
|
||||
@@ -520,6 +521,11 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
req.sessionName(), spawned.agentSessionId(), spawned.receipt());
|
||||
}
|
||||
|
||||
/** The herdr daemon that owns this launcher's pane coordinates. */
|
||||
public HerdrClient herdr() {
|
||||
return agents.herdr();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
return effectiveCwd(req.profileName(), req.requestedCwd(), req.callerCwd());
|
||||
@@ -1033,34 +1039,40 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* pass it through {@code tab.create}/{@code pane.split}. Returns the directory for teardown
|
||||
* registration, or {@code null} when the policy does not apply.
|
||||
*
|
||||
* <p>The allow-list handed to the generator is the derived profile set ({@link
|
||||
* MemberEnvAllowList#derive}) UNIONed with the exact keys of THIS launch's env map — names the
|
||||
* daemon itself injects must survive its own control. {@code SSH_AUTH_SOCK} is added ONLY when
|
||||
* the config explicitly allows it; by default it is absent, so the scrub blanks it like any
|
||||
* other non-derived name.
|
||||
* <p>The allow-list handed to the generator is the derived profile set UNIONed with the
|
||||
* operator's own {@code memberCredentials.allow:} names ({@link MemberEnvAllowList#derive(
|
||||
* Collection, Set)} — CB-633 follow-up) and with the exact keys of THIS launch's env map —
|
||||
* names the daemon itself injects must survive its own control. {@code SSH_AUTH_SOCK} is added
|
||||
* ONLY when the config explicitly allows it, EVEN IF the operator also listed it under
|
||||
* {@code allow:}; by default it is absent, so the scrub blanks it like any other non-derived
|
||||
* name. It stays a one-off decision because it is a live handle to the operator's ssh-agent, not
|
||||
* a value — a member holding it can sign with every key the agent holds, so letting it ride in
|
||||
* on the generic {@code allow:} list would hand that out for an unrelated reason.
|
||||
*/
|
||||
private Path applyEnvironmentAllowListPolicy(FleetConfig.Profile cfg, Launch launch) {
|
||||
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
|
||||
if (creds == null || !creds.isAllowList()) {
|
||||
return null;
|
||||
}
|
||||
Set<String> allowed = derivedAllowedNames(creds, launch);
|
||||
String loginShell = resolveEnv("SHELL");
|
||||
boolean zsh = loginShell != null && (loginShell.endsWith("/zsh") || loginShell.equals("zsh"));
|
||||
if (!zsh) {
|
||||
// A non-zsh login shell ignores ZDOTDIR entirely: NO scrub would run, so pretending
|
||||
// otherwise would be worse than saying so. Warn loudly and fall back to the CB-596
|
||||
// sentinel overlay over the enumerated known: names — weaker (a sourced file can undo
|
||||
// it), but strictly better than nothing.
|
||||
// it), but strictly better than nothing. Deliberately no "allowed N of M" line here: the
|
||||
// scrub this count describes does not run on this path, so printing it would tell an
|
||||
// operator that a fraction of names were blocked when the real number blocked is zero.
|
||||
// logCredentialGap's WARN (below) is the only signal for this path.
|
||||
warnNonZsh(loginShell);
|
||||
overlayBlockedCredentials(launch.env(), creds);
|
||||
logCredentialGap(creds);
|
||||
return null;
|
||||
}
|
||||
Set<String> allowed = new java.util.TreeSet<>(MemberEnvAllowList.derive(profiles.values()));
|
||||
if (creds.sshAuthSockAllowed()) {
|
||||
allowed.add(SSH_AUTH_SOCK);
|
||||
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
|
||||
allowed.addAll(launch.env().keySet());
|
||||
// Only reached when the scrub is actually about to run — the count below describes that
|
||||
// scrub, so it must not be logged before this gate (see the non-zsh branch above).
|
||||
logAllowListCoverage(allowed);
|
||||
Path dir = EnvAllowListScrub.generate(Path.of(System.getProperty("java.io.tmpdir")), allowed);
|
||||
launch.env().put("ZDOTDIR", dir.toAbsolutePath().toString());
|
||||
log.info("memberCredentials policy=allow-list: profile={} generated ZDOTDIR {} — derived "
|
||||
@@ -1069,8 +1081,42 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
return dir;
|
||||
}
|
||||
|
||||
/**
|
||||
* The full kept-name set for this spawn: the profile-derived names, unioned with {@code
|
||||
* memberCredentials.allow:} (CB-633 follow-up — previously ignored by this whole policy), the
|
||||
* ssh-agent handle when explicitly allowed, and the exact keys of THIS launch's own env map.
|
||||
*/
|
||||
private Set<String> derivedAllowedNames(FleetConfig.MemberCredentials creds, Launch launch) {
|
||||
Set<String> allowed = new java.util.TreeSet<>(
|
||||
MemberEnvAllowList.derive(profiles.values(), creds.allowSet()));
|
||||
if (creds.sshAuthSockAllowed()) {
|
||||
allowed.add(SSH_AUTH_SOCK);
|
||||
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
|
||||
allowed.addAll(launch.env().keySet());
|
||||
return allowed;
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up: one INFO line per allow-list spawn WHOSE SCRUB ACTUALLY RUNS, so an operator
|
||||
* can read a single log line and know the scrub ran and how much of the visible environment it
|
||||
* will keep. Callable ONLY from the zsh branch of {@link #applyEnvironmentAllowListPolicy}, after
|
||||
* the shell gate — logging it before that gate (or on the non-zsh fallback, where nothing is
|
||||
* scrubbed) would tell an operator a fraction of names were blocked when the real number blocked
|
||||
* is zero, which is worse than not logging at all. {@code M} is {@link #hostEnvNames}' size (the
|
||||
* daemon's own environment — see that field's javadoc for why it stands in for the pane's, which
|
||||
* the daemon has no channel to inspect at spawn time) and {@code N} is how many of those names
|
||||
* survive {@code allowed} (including the {@code LC_*} prefix rule). Neither number is a constant:
|
||||
* both come from the actual derived set and the actual environment this spawn sees. Never logs a
|
||||
* variable NAME or VALUE — only the counts.
|
||||
*/
|
||||
private void logAllowListCoverage(Set<String> allowed) {
|
||||
Set<String> hostNames = hostEnvNames.get();
|
||||
long kept = hostNames.stream().filter(name -> MemberEnvAllowList.keeps(allowed, name)).count();
|
||||
log.info("member credentials: allowed {} of {}", kept, hostNames.size());
|
||||
}
|
||||
|
||||
/** The operator ssh-agent handle — kept ONLY by explicit config decision, never by default. */
|
||||
private static final String SSH_AUTH_SOCK = "SSH_AUTH_SOCK";
|
||||
private static final String SSH_AUTH_SOCK = MemberEnvAllowList.SSH_AUTH_SOCK;
|
||||
|
||||
/**
|
||||
* CB-633: a non-zsh login shell means the allow-list control CANNOT run — say so once per
|
||||
|
||||
@@ -35,9 +35,27 @@ import java.util.TreeSet;
|
||||
* <p>{@code SSH_AUTH_SOCK} is deliberately NOT here. It is a handle to the operator's ssh-agent — a
|
||||
* member holding it can sign with the operator's keys — so keeping it is a config decision
|
||||
* ({@code memberCredentials.sshAuthSock: allow}), not a derivation default.
|
||||
*
|
||||
* <p><b>CB-633 follow-up:</b> the union also includes {@code memberCredentials.allow:} — the
|
||||
* operator's own explicit list. Before this, {@code policy: allow-list} silently ignored every name
|
||||
* an operator wrote under {@code allow:} unless a profile happened to carry it too, which meant
|
||||
* turning the policy on could blank credentials working members already depended on. {@code
|
||||
* SSH_AUTH_SOCK} is the one exception: even when the operator lists it under {@code allow:}, it is
|
||||
* excluded here and added back ONLY by the caller when {@code sshAuthSock: allow} is explicitly set
|
||||
* (see {@link #SSH_AUTH_SOCK}'s javadoc) — it is a live handle to the operator's own ssh-agent, not
|
||||
* a value, so treating it like any other allow-listed name would hand a member every key the
|
||||
* operator's agent holds the moment they typed the name under {@code allow:} for an unrelated
|
||||
* reason.
|
||||
*/
|
||||
public final class MemberEnvAllowList {
|
||||
|
||||
/**
|
||||
* The operator's ssh-agent socket path. Deliberately excluded from {@link #derive}'s union of
|
||||
* {@code memberCredentials.allow:} — see the class javadoc's CB-633 follow-up note. Governed
|
||||
* ONLY by {@code memberCredentials.sshAuthSock}, never by appearing in {@code allow:}.
|
||||
*/
|
||||
public static final String SSH_AUTH_SOCK = "SSH_AUTH_SOCK";
|
||||
|
||||
/**
|
||||
* Names that are not credentials and that a login shell or agent binary genuinely needs.
|
||||
*
|
||||
@@ -73,9 +91,21 @@ public final class MemberEnvAllowList {
|
||||
|
||||
/**
|
||||
* Derive the allowed NAME set from the given profiles plus {@link #INFRASTRUCTURE_PASSTHROUGH}.
|
||||
* Deterministic (sorted) so generated scrub files are diffable run-to-run.
|
||||
* Equivalent to {@link #derive(Collection, Set)} with no operator-configured names — kept for
|
||||
* callers (and existing tests) that only care about the profile-derived half.
|
||||
*/
|
||||
public static Set<String> derive(Collection<FleetConfig.Profile> profiles) {
|
||||
return derive(profiles, Set.of());
|
||||
}
|
||||
|
||||
/**
|
||||
* Derive the allowed NAME set: the profile-derived union above, PLUS {@code configuredAllow} —
|
||||
* the operator's own {@code memberCredentials.allow:} list (CB-633 follow-up). {@code
|
||||
* SSH_AUTH_SOCK} is dropped from {@code configuredAllow} even if the operator listed it there;
|
||||
* see the class javadoc for why. Deterministic (sorted) so generated scrub files are diffable
|
||||
* run-to-run.
|
||||
*/
|
||||
public static Set<String> derive(Collection<FleetConfig.Profile> profiles, Set<String> configuredAllow) {
|
||||
Set<String> derived = new TreeSet<>(INFRASTRUCTURE_PASSTHROUGH);
|
||||
if (profiles != null) {
|
||||
for (FleetConfig.Profile p : profiles) {
|
||||
@@ -87,6 +117,13 @@ public final class MemberEnvAllowList {
|
||||
}
|
||||
}
|
||||
}
|
||||
if (configuredAllow != null) {
|
||||
for (String name : configuredAllow) {
|
||||
if (name != null && !name.isBlank() && !SSH_AUTH_SOCK.equals(name)) {
|
||||
derived.add(name);
|
||||
}
|
||||
}
|
||||
}
|
||||
return Set.copyOf(derived);
|
||||
}
|
||||
|
||||
|
||||
@@ -29,8 +29,9 @@ import java.util.concurrent.TimeoutException;
|
||||
*
|
||||
* <p><strong>Mapping — consume-and-hold with deferred manual ack.</strong> Each target has a durable
|
||||
* queue {@code agent.<target>.inbox}. The gateway that owns the target starts a manual-ack consumer
|
||||
* ({@link #own}) that pulls persistent messages off that queue into an in-memory <em>held</em> map
|
||||
* (keyed by {@code msgId}) but does <em>not</em> ack them. {@link #peek} returns that snapshot;
|
||||
* ({@link #own}) that pulls persistent messages, up to its prefetch window, off that queue into an
|
||||
* in-memory <em>held</em> map (keyed by {@code msgId}) but does <em>not</em> ack them.
|
||||
* {@link #peek} returns that snapshot;
|
||||
* {@link #ack} acks the broker delivery-tag and drops the entry. Because messages stay unacked until
|
||||
* the owning gateway actually drains them, a crash (or a {@code java -jar} bounce) before caller-ack
|
||||
* leaves them on the broker — it redelivers on reconnect. That is the durability the in-memory
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -31,6 +31,7 @@ import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Predicate;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
@@ -58,7 +59,7 @@ public final class FleetApp {
|
||||
private final PeerLauncher workers;
|
||||
private final SessionManager sessions; // CB-301: authoritative session registry
|
||||
private final MessageService messages;
|
||||
private final MemberPresence presence; // CB-113: which workers are MCP-connected (available)
|
||||
private final Predicate<String> deliverable;
|
||||
private final HttpServlet mcpServlet; // MCP Streamable-HTTP endpoint, mounted at /mcp (nullable)
|
||||
private final CallerResolver auth; // CB-501: null → authz not enforced (legacy behaviour)
|
||||
private final Metrics metrics; // CB-502: null → /metrics not exposed
|
||||
@@ -70,9 +71,9 @@ public final class FleetApp {
|
||||
* behaviour without each needing an auth fixture.
|
||||
*/
|
||||
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet) {
|
||||
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null);
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet) {
|
||||
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null, presence::isPresent);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -82,13 +83,23 @@ public final class FleetApp {
|
||||
* the endpoint
|
||||
*/
|
||||
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) {
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) {
|
||||
this(herdr, workers, sessions, messages, presence, mcpServlet, auth, metrics, presence::isPresent);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param deliverable the injector's readiness gate, shared so status reports its real result
|
||||
*/
|
||||
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
|
||||
MessageService messages, MemberPresence presence,
|
||||
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics,
|
||||
Predicate<String> deliverable) {
|
||||
this.herdr = herdr;
|
||||
this.workers = workers;
|
||||
this.sessions = sessions;
|
||||
this.messages = messages;
|
||||
this.presence = presence;
|
||||
this.deliverable = deliverable;
|
||||
this.mcpServlet = mcpServlet;
|
||||
this.auth = auth;
|
||||
this.metrics = metrics;
|
||||
@@ -507,9 +518,9 @@ public final class FleetApp {
|
||||
|
||||
/**
|
||||
* Live lifecycle status of a worker (MCP `fleet_status` wraps this in CB-105), plus its
|
||||
* <em>readiness</em> (CB-113): {@code ready} is true once the worker's Claude has connected the
|
||||
* bridge MCP — the reliable "available to receive a task" signal, unlike bare {@code idle}, which
|
||||
* is also true during boot.
|
||||
* <em>readiness</em>: {@code ready} is true when the injector can deliver to the target. A
|
||||
* spawned member must connect the bridge MCP first, while a known lead is ready without member
|
||||
* presence. This differs from bare {@code idle}, which is also true during member boot.
|
||||
*/
|
||||
private void sessionStatus(Context ctx) {
|
||||
String id = ctx.pathParam("id");
|
||||
@@ -520,7 +531,7 @@ public final class FleetApp {
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
body.put("sessionId", id);
|
||||
body.put("status", messages.status(id).name().toLowerCase());
|
||||
body.put("ready", presence.isPresent(id));
|
||||
body.put("ready", deliverable.test(id));
|
||||
// CB-582: a worker paused mid-turn in an async fleet_ask is otherwise invisible to a
|
||||
// status poll — surface the open question and how to answer it, same as fleet_poll's
|
||||
// Phase.ASKING view.
|
||||
|
||||
@@ -7,6 +7,8 @@ import java.io.BufferedReader;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStreamReader;
|
||||
import java.io.UncheckedIOException;
|
||||
import java.net.URI;
|
||||
import java.net.URISyntaxException;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
@@ -20,6 +22,7 @@ import java.util.Optional;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
@@ -64,6 +67,10 @@ public final class GitWorktrees implements Worktrees {
|
||||
/** What {@link #isolateToolSurface} writes for {@code .autoenv}: a valid, empty env file. */
|
||||
private static final String NEUTRAL_AUTOENV_CONFIG = "";
|
||||
|
||||
/** A credential helper command which reads only an environment variable at Git call time. */
|
||||
private static final String ENVIRONMENT_CREDENTIAL_HELPER = "!f() { if [ \"$1\" = get ]; then "
|
||||
+ "printf 'username=%s\\npassword=%s\\n\\n' git \"$WORKER_GITEA_TOKEN\"; fi; }; f";
|
||||
|
||||
/**
|
||||
* A tracked project config that is hostile in a provisioned worktree, and what to replace it
|
||||
* with. {@link #file} is the repo-relative path; {@link #stub} is a neutral but VALID payload for
|
||||
@@ -81,6 +88,7 @@ public final class GitWorktrees implements Worktrees {
|
||||
);
|
||||
|
||||
private final String configuredRoot;
|
||||
private final Consumer<String> afterWorktreeAdded;
|
||||
private final SecureRandom random = new SecureRandom();
|
||||
private final AtomicLong seq = new AtomicLong();
|
||||
|
||||
@@ -91,7 +99,13 @@ public final class GitWorktrees implements Worktrees {
|
||||
|
||||
/** @param configuredRoot nullable absolute or relative path; null/blank derives a sibling of the repo root. */
|
||||
public GitWorktrees(String configuredRoot) {
|
||||
this(configuredRoot, _ -> {});
|
||||
}
|
||||
|
||||
/** Test seam for changing a real worktree between its creation and its security check. */
|
||||
GitWorktrees(String configuredRoot, Consumer<String> afterWorktreeAdded) {
|
||||
this.configuredRoot = configuredRoot;
|
||||
this.afterWorktreeAdded = afterWorktreeAdded == null ? _ -> {} : afterWorktreeAdded;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -107,11 +121,148 @@ public final class GitWorktrees implements Worktrees {
|
||||
}
|
||||
String wt = path.toAbsolutePath().toString();
|
||||
log.info("adding worktree branch={} path={} base={}", branch, wt, base);
|
||||
removeUserInfoFromHttpsOrigin(repoRoot);
|
||||
exec("git", "-C", repoRoot, "worktree", "add", wt, "-b", branch, base);
|
||||
afterWorktreeAdded.accept(wt);
|
||||
requireCredentialFreeHttpsOrigin(wt);
|
||||
configureEnvironmentCredentialHelper(repoRoot, wt);
|
||||
configureHttpsUrlRewriteForSshOrigin(repoRoot, wt);
|
||||
isolateToolSurface(wt);
|
||||
return wt;
|
||||
}
|
||||
|
||||
/**
|
||||
* A linked worktree shares its primary checkout's git config. Remove HTTPS user info before
|
||||
* adding one, so a credential accidentally embedded in that config cannot reach the member.
|
||||
*/
|
||||
private void removeUserInfoFromHttpsOrigin(String repoRoot) {
|
||||
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
|
||||
return;
|
||||
}
|
||||
String origin = exec("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
|
||||
URI uri;
|
||||
try {
|
||||
uri = new URI(origin);
|
||||
} catch (URISyntaxException e) {
|
||||
throw new WorktreeException("origin URL is invalid; cannot provision a safe worktree", e);
|
||||
}
|
||||
if (!"https".equalsIgnoreCase(uri.getScheme()) || uri.getUserInfo() == null) {
|
||||
return;
|
||||
}
|
||||
int schemeEnd = origin.indexOf("://") + 3;
|
||||
int userInfoEnd = origin.indexOf('@', schemeEnd);
|
||||
if (userInfoEnd < schemeEnd) {
|
||||
throw new WorktreeException("origin URL has invalid HTTPS user info; cannot provision safely");
|
||||
}
|
||||
String cleanOrigin = origin.substring(0, schemeEnd) + origin.substring(userInfoEnd + 1);
|
||||
exec("git", "-C", repoRoot, "remote", "set-url", "origin", cleanOrigin);
|
||||
log.info("removed HTTPS user info from forge origin before provisioning worktree");
|
||||
}
|
||||
|
||||
/** Refuse the worktree if Git still resolves any HTTPS origin URL with embedded credentials. */
|
||||
private void requireCredentialFreeHttpsOrigin(String worktreePath) {
|
||||
if (exitCode("git", "-C", worktreePath, "remote", "get-url", "--all", "origin") != 0) {
|
||||
return;
|
||||
}
|
||||
String origins = exec("git", "-C", worktreePath, "remote", "get-url", "--all", "origin");
|
||||
for (String origin : origins.split("\\R")) {
|
||||
try {
|
||||
URI uri = new URI(origin);
|
||||
if ("https".equalsIgnoreCase(uri.getScheme()) && uri.getUserInfo() != null) {
|
||||
throw new WorktreeException("worktree origin contains HTTPS user info; refusing provision");
|
||||
}
|
||||
} catch (URISyntaxException e) {
|
||||
throw new WorktreeException("worktree origin URL is invalid; refusing provision", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Configure a per-worktree helper that supplies a token from the member environment at call time. */
|
||||
private void configureEnvironmentCredentialHelper(String repoRoot, String worktreePath) {
|
||||
exec("git", "-C", repoRoot, "config", "extensions.worktreeConfig", "true");
|
||||
// An empty helper resets values inherited from the system or global config. Without it Git
|
||||
// asks the next helper after this one, which can expose an operator-level credential.
|
||||
exec("git", "-C", worktreePath, "config", "--worktree", "--replace-all", "credential.helper", "");
|
||||
exec("git", "-C", worktreePath, "config", "--worktree", "--add", "credential.helper",
|
||||
ENVIRONMENT_CREDENTIAL_HELPER);
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link #configureEnvironmentCredentialHelper} only ever fires for an HTTPS origin — Git never
|
||||
* consults a {@code credential.helper} for an SSH transport. This repo's own origin is
|
||||
* {@code ssh://git@git.ltms.dev:2224/fleet/fleetd.git}, so a member sitting on that origin never
|
||||
* reaches the helper and the repo-scoped {@code WORKER_GITEA_TOKEN} is simply not used.
|
||||
*
|
||||
* <p>An earlier version of this javadoc justified the rewrite by claiming a member <em>cannot</em>
|
||||
* push once {@code memberCredentials.policy: allow-list} blocks {@code SSH_AUTH_SOCK}, because
|
||||
* "there is no private key file on this host, only an ssh-agent socket". That premise is false
|
||||
* (fleetd #184): {@code ssh -G} resolves a readable, passphrase-free {@code IdentityFile} outside
|
||||
* {@code ~/.ssh}, and a member — same OS user — pushes over SSH with the socket blanked. The
|
||||
* rewrite is still worth having, but for the reason below rather than that one: it routes the
|
||||
* member through its own scoped token instead of the operator's ssh identity, which is what makes
|
||||
* a member's pushes attributable and revocable.
|
||||
*
|
||||
* <p>The fix is a <em>worktree-scoped</em> URL rewrite: {@code url.<https-base>.insteadOf
|
||||
* <ssh-base>}, set with {@code --worktree} so it lands only in
|
||||
* {@code <worktree>/.git/worktrees/<name>/config.worktree} (enabled by
|
||||
* {@code extensions.worktreeConfig}, already turned on above) and never touches the shared
|
||||
* repo-level config the primary checkout also reads. {@code insteadOf} — not
|
||||
* {@code pushInsteadOf} — because a member may also need to fetch or rebase, and both should go
|
||||
* through the member's own token for the same reason.
|
||||
*
|
||||
* <p>The host (and, for the rewrite's SSH-side match, the port) come from parsing the origin
|
||||
* itself — never a hardcoded forge host, which is exactly what #177 removed. An origin that is
|
||||
* already {@code https://} is left alone; the credential helper already covers it. An origin
|
||||
* that is neither {@code ssh://} nor {@code https://} — including the scp-like shorthand
|
||||
* ({@code git@host:path}, no scheme) — is left untouched deliberately: that shorthand's
|
||||
* {@code host:path} split is defined by the user's ssh_config aliases, not by URI syntax, so
|
||||
* guessing at it risks rewriting to the wrong place. A repo provisioned from that form keeps
|
||||
* today's (broken, if the policy blocks the agent) SSH-only behaviour rather than a wrong rewrite.
|
||||
*/
|
||||
/**
|
||||
* Blank the user-info of a remote URL before it reaches a log. A remote URL is not obviously a
|
||||
* credential channel, which is exactly why one has leaked here three times ({@code git remote -v}
|
||||
* printing a token inline, and fleetd #157 / #182). An {@code ssh://} authority normally carries
|
||||
* only {@code git@}, so this usually changes nothing — it is here so that the one origin that
|
||||
* does carry a secret cannot print it. Matches every {@code ://…@} pair, not just the first.
|
||||
*/
|
||||
private static String redactUserInfo(String url) {
|
||||
return url == null ? null : url.replaceAll("://[^@/]*@", "://<redacted>@");
|
||||
}
|
||||
|
||||
private void configureHttpsUrlRewriteForSshOrigin(String repoRoot, String worktreePath) {
|
||||
if (exitCode("git", "-C", repoRoot, "config", "--get", "remote.origin.url") != 0) {
|
||||
return;
|
||||
}
|
||||
String origin = exec("git", "-C", repoRoot, "config", "--get", "remote.origin.url").trim();
|
||||
URI uri;
|
||||
try {
|
||||
uri = new URI(origin);
|
||||
} catch (URISyntaxException e) {
|
||||
log.warn("origin URL {} is not a valid URI; skipping worktree HTTPS rewrite", redactUserInfo(origin));
|
||||
return;
|
||||
}
|
||||
String scheme = uri.getScheme();
|
||||
if (!"ssh".equalsIgnoreCase(scheme)) {
|
||||
// Already https:// (the credential helper covers it), or a scheme-less/scp-like origin
|
||||
// left alone on purpose — see the javadoc above.
|
||||
log.debug("origin scheme is not ssh ({}) — no worktree HTTPS rewrite needed", redactUserInfo(origin));
|
||||
return;
|
||||
}
|
||||
String host = uri.getHost();
|
||||
String authority = uri.getRawAuthority();
|
||||
if (host == null || host.isBlank() || authority == null || authority.isBlank()) {
|
||||
log.warn("ssh origin {} has no resolvable host; skipping worktree HTTPS rewrite", redactUserInfo(origin));
|
||||
return;
|
||||
}
|
||||
String sshBase = "ssh://" + authority + "/";
|
||||
String httpsBase = "https://" + host + "/";
|
||||
exec("git", "-C", worktreePath, "config", "--worktree", "--replace-all",
|
||||
"url." + httpsBase + ".insteadOf", sshBase);
|
||||
log.info("worktree {} rewrites {} to {} (worktree-scoped; parent checkout untouched)",
|
||||
worktreePath, redactUserInfo(sshBase), httpsBase);
|
||||
}
|
||||
|
||||
/**
|
||||
* Neutralize the worktree's worktree-hostile project configs so a worker inherits only the tools
|
||||
* and environment its launcher mounts (the bridge via {@code --mcp-config}, the opencode config
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -197,10 +197,15 @@ class CompletionResolverTest {
|
||||
void resolvesSynchronouslyBeforePostTurnContextClearing() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ previous answer\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
// fleetd#164: an injectable clock, so this real (non-crash) turn lands outside MIN_TURN_NANOS
|
||||
// — captureBaseline and resolveBeforePostAction below run back-to-back with no real delay.
|
||||
long[] clock = {1_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter));
|
||||
herdr.readText("⏺ answer that /clear would erase\n❯ ");
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // this turn took longer than the floor
|
||||
|
||||
resolver.resolveBeforePostAction("term_a");
|
||||
|
||||
@@ -219,7 +224,11 @@ class CompletionResolverTest {
|
||||
String longBlock = "⏺ " + "x".repeat(CompletionResolver.MAX_SCRAPE_CHARS + 500) + "\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(longBlock);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
// fleetd#164: an injectable clock so captureBaseline and resolve (back-to-back, no real
|
||||
// delay) don't trip the too-fast-turn floor — this test is about the suppression guard, not timing.
|
||||
long[] clock = {1_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
|
||||
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block
|
||||
@@ -227,6 +236,7 @@ class CompletionResolverTest {
|
||||
assertEquals(CompletionResolver.MAX_SCRAPE_CHARS, turn.baseline().length(),
|
||||
"the delivery baseline is clipped to the same cap resolve() applies to the tail");
|
||||
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // outside the floor
|
||||
resolver.resolve("term_a", turn); // scrape unchanged → clipped tail == baseline → suppress
|
||||
|
||||
assertFalse(waiter.isDone(),
|
||||
@@ -249,12 +259,13 @@ class CompletionResolverTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolvesWhenTheScrapeItselfFailsEvenWithABaselinePresent() {
|
||||
// The most important branch of the CB-115 guard: a failed read means the resolver could not
|
||||
// SEE the screen — "couldn't see", not "no change". It must still resolve the send (an empty
|
||||
// tail beats hanging until the caller's timeout), even though a baseline was captured. The
|
||||
// baseline here is "" (an empty pane at delivery), so without the !scrapeFailed clause the
|
||||
// byte-identical guard would wrongly match the empty tail and suppress.
|
||||
void aFailedScrapeResolvesAsAFailureEvenWithABaselinePresent() {
|
||||
// fleetd#164: this test used to assert that a failed read resolved the send as a SUCCESS
|
||||
// carrying an empty string ("an empty tail beats hanging until the caller's timeout") — that
|
||||
// was the bug this ticket fixes: a lost turn and a genuine empty answer looked identical to
|
||||
// every caller. This test encoded the bug and is changed here: a failed read must fail the
|
||||
// send instead, naming the member, whether or not a baseline was captured (the baseline here
|
||||
// is "", an empty pane at delivery — proof this isn't the CB-115 misattribution path either).
|
||||
FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
@@ -265,8 +276,82 @@ class CompletionResolverTest {
|
||||
|
||||
assertTrue(waiter.isDone(),
|
||||
"a failed scrape must still resolve the send, not hang until the caller's timeout");
|
||||
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind());
|
||||
assertEquals("", waiter.getNow(null).text(), "the tail is empty because the screen was unreadable");
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"a failed read is a lost turn, not a successful empty reply");
|
||||
assertTrue(waiter.getNow(null).text().contains("term_a"),
|
||||
"the failure names the member: " + waiter.getNow(null).text());
|
||||
assertTrue(waiter.getNow(null).text().contains("could not be read"),
|
||||
"the failure explains the scrape could not be read: " + waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
// --- fleetd#164: an empty (but readable) scrape must never resolve as a success --------------
|
||||
|
||||
@Test
|
||||
void anEmptyScrapeResolvesAsAFailureNamingTheMember() {
|
||||
// The core defect: a scrape that read CLEANLY but produced zero characters used to resolve
|
||||
// the send as a SUCCESS carrying "" — indistinguishable, to every caller, from a worker that
|
||||
// genuinely finished with nothing to say. A lost turn must never look like a real empty reply.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertTrue(waiter.isDone(), "an empty scrape must still resolve the send, not hang");
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"an empty scrape is a lost turn, not a successful empty reply");
|
||||
assertTrue(waiter.getNow(null).text().contains("term_a"),
|
||||
"the failure names the member: " + waiter.getNow(null).text());
|
||||
assertTrue(waiter.getNow(null).text().toLowerCase().contains("empty"),
|
||||
"the failure says the scrape was empty: " + waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
// --- fleetd#164: a BUSY -> DONE transition inside the floor is a crash, not a fast answer ------
|
||||
|
||||
@Test
|
||||
void aBusyToDoneTransitionInsideTheFloorResolvesAsAFailure() {
|
||||
// The exact fleetd#164 scenario: the backend returned an HTTP 400 before the worker did
|
||||
// anything, and the member went BUSY -> DONE in ~1s. That transition alone is indistinguishable
|
||||
// from a genuine (if unusually fast) completion, so the resolver leans on the floor to catch it.
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ HTTP 400: invalid request\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
long[] clock = {10_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]); // delivered "now"
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS - 1; // 1ns inside the floor
|
||||
|
||||
resolver.resolve("term_a", turn);
|
||||
|
||||
assertTrue(waiter.isDone(), "a suspiciously fast turn must still resolve (as a failure)");
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
|
||||
assertTrue(waiter.getNow(null).text().contains("term_a"),
|
||||
"the failure names the member: " + waiter.getNow(null).text());
|
||||
assertTrue(waiter.getNow(null).text().contains("HTTP 400"),
|
||||
"the failure carries whatever was on screen: " + waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBusyToDoneTransitionJustOutsideTheFloorResolvesNormally() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ a real, if quick, answer\n❯ ");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
long[] clock = {10_000_000_000L};
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), () -> clock[0]);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
var turn = new CompletionResolver.InFlight(waiter, null, clock[0]); // delivered "now"
|
||||
clock[0] += CompletionResolver.MIN_TURN_NANOS + 1; // 1ns outside the floor
|
||||
|
||||
resolver.resolve("term_a", turn);
|
||||
|
||||
assertTrue(waiter.isDone());
|
||||
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(),
|
||||
"a turn that took longer than the floor resolves normally");
|
||||
assertEquals("a real, if quick, answer", waiter.getNow(null).text());
|
||||
}
|
||||
|
||||
// --- CB-115/CB-116 fail guard: an already-done or absent waiter is left alone ---------
|
||||
|
||||
@@ -8,6 +8,7 @@ import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.Agent;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.CharterReceipt;
|
||||
@@ -248,6 +249,112 @@ class CompositePeerLauncherTest {
|
||||
"stop routes to the spawning adapter and closes exactly that worker's pane");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopKeepsSameHerdrPaneIdSeparateByOwningAdapter() {
|
||||
// Separate herdr daemons can both issue w9:pRoot_1. The opaque handles identify their
|
||||
// spawning adapters, so each stop reaches only its recorded owner.
|
||||
FakeHerdr first = new FakeHerdr();
|
||||
FakeHerdr second = new FakeHerdr();
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
PeerHandle claude = composite.spawn(new SpawnRequest("claude", null, null));
|
||||
PeerHandle opencode = composite.spawn(new SpawnRequest("gemini", null, null));
|
||||
|
||||
assertNotEquals(claude.id(), opencode.id(), "each public paneId keeps its adapter owner");
|
||||
composite.stop(claude.id());
|
||||
assertTrue(first.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w9:pRoot_1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"the first daemon closes its own pane");
|
||||
assertFalse(second.called("pane.close"), "the matching pane on the second daemon stays live");
|
||||
|
||||
composite.stop(opencode.id());
|
||||
assertTrue(second.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w9:pRoot_1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"the second daemon then closes its own pane");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopAllowsALegacyBarePaneIdWithOneDaemon() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
PeerLauncher composite = new CompositePeerLauncher(List.of(claudeAdapter(herdr)), "claude");
|
||||
|
||||
composite.stop("w9:pRoot_1");
|
||||
|
||||
assertTrue(herdr.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w9:pRoot_1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"one daemon keeps the legacy bare-pane routing behaviour");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopAllowsAnUnownedPaneIdWithTwoAdaptersSharingOneDaemon() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
PeerLauncher composite = composite(herdr);
|
||||
|
||||
composite.stop("w9:pRoot_1");
|
||||
|
||||
assertTrue(herdr.calls.stream().anyMatch(c -> c.method().equals("pane.close")
|
||||
&& "w9:pRoot_1".equals(((Map<?, ?>) c.params()).get("pane_id"))),
|
||||
"two adapter kinds sharing one daemon keep the fallback route");
|
||||
}
|
||||
|
||||
@Test
|
||||
void listKeepsBothPanesWhenTwoDaemonsShareAPaneId() {
|
||||
// herdr pane ids are per-daemon counters, so two daemons really can both hold w1:p1 on
|
||||
// different panes. Deduplicating on the pane id alone dropped one of the two real agents.
|
||||
FakeHerdr first = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
|
||||
FakeHerdr second = new FakeHerdr().withAgent("y", "term_y", "w1:p1", "w1:t1");
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(second)), "claude");
|
||||
|
||||
List<Agent> agents = composite.list();
|
||||
|
||||
assertEquals(2, agents.stream().filter(a -> "w1:p1".equals(a.paneId())).count(),
|
||||
"one w1:p1 per daemon survives — the pane id alone is not a unique key");
|
||||
assertTrue(agents.stream().anyMatch(a -> "term_x".equals(a.terminalId())));
|
||||
assertTrue(agents.stream().anyMatch(a -> "term_y".equals(a.terminalId())));
|
||||
}
|
||||
|
||||
@Test
|
||||
void listStillDeduplicatesTwoAdaptersSharingOneDaemon() {
|
||||
// Both adapters ask the SAME daemon, so both see the same agent set. Without the dedupe this
|
||||
// would report every agent twice; the daemon key must not break that.
|
||||
FakeHerdr herdr = new FakeHerdr().withAgent("x", "term_x", "w1:p1", "w1:t1");
|
||||
CompositePeerLauncher composite = composite(herdr);
|
||||
|
||||
List<Agent> agents = composite.list();
|
||||
|
||||
assertEquals(1, agents.stream().filter(a -> "w1:p1".equals(a.paneId())).count(),
|
||||
"one daemon still reports each of its agents once");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopKeepsTheOwnerRecordWhenTheDelegateRefusesTheStop() {
|
||||
// Removing the record before the delegate accepted the stop lost the owner on failure: the
|
||||
// pane was still alive, but the retry landed in the ambiguous branch and refused it for good.
|
||||
FakeHerdr first = new FakeHerdr().paneCloseFailsWith("pane_busy");
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(first), opencodeAdapter(new FakeHerdr())), "claude");
|
||||
PeerHandle claude = composite.spawn(new SpawnRequest("claude", null, null));
|
||||
|
||||
assertThrows(HerdrException.class, () -> composite.stop(claude.id()));
|
||||
|
||||
// The retry must still know its owner — a HerdrException, never "ambiguous paneId".
|
||||
assertThrows(HerdrException.class, () -> composite.stop(claude.id()),
|
||||
"the owner record survives a failed stop, so the retry is not ambiguous");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopRejectsAnUnownedPaneIdWhenMultipleDaemonsCouldOwnIt() {
|
||||
PeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(claudeAdapter(new FakeHerdr()), opencodeAdapter(new FakeHerdr())), "claude");
|
||||
|
||||
IllegalArgumentException error = assertThrows(IllegalArgumentException.class,
|
||||
() -> composite.stop("w1:p1"));
|
||||
|
||||
assertEquals("ambiguous paneId 'w1:p1': no owning herdr daemon was recorded", error.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void opencodeContextResetIsANoOpAndWarnsOnlyOnce() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
+128
-1
@@ -1,5 +1,9 @@
|
||||
package dev.ltms.fleet.member;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
@@ -8,6 +12,7 @@ import dev.ltms.fleet.peer.Capability;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
@@ -108,6 +113,122 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null);
|
||||
}
|
||||
|
||||
/** Same as {@link #allowList()} but with an operator-configured {@code allow:} list. */
|
||||
private static Supplier<FleetConfig.MemberCredentials> allowListWithAllow(List<String> allow) {
|
||||
return () -> new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, allow, List.of(), null);
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up: a name that lives ONLY in {@code memberCredentials.allow:} — no profile
|
||||
* mentions it — must survive the scrub the real spawn path generates. Calling {@code
|
||||
* MemberEnvAllowList.derive} directly (as {@link MemberEnvAllowListTest} does) would pass even
|
||||
* if {@code HerdrPeerLauncher} never threaded {@code allow:} into the derivation at all; this
|
||||
* test goes through {@link HerdrPeerLauncher#spawn}, the method the daemon actually calls at
|
||||
* spawn time, so it proves the union is wired in, not just correct in isolation.
|
||||
*/
|
||||
@Test
|
||||
void spawningUnderAllowListPolicyIncludesAnOperatorConfiguredAllowName() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr,
|
||||
allowListWithAllow(List.of("OPERATOR_ONLY_NAME")));
|
||||
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
Path dir = Path.of(launcher.env.get("ZDOTDIR"));
|
||||
assertTrue(readAll(dir.resolve(EnvAllowListScrub.SCRUB_FILE)).contains("OPERATOR_ONLY_NAME"),
|
||||
"a name only in memberCredentials.allow: must reach the generated scrub through the "
|
||||
+ "real launcher spawn path");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code SSH_AUTH_SOCK} is a live ssh-agent handle, not a value — it must stay blocked under
|
||||
* {@code allow-list} even when the operator lists it under {@code allow:}, because {@code
|
||||
* sshAuthSock} defaults to blocked. Governed ONLY by {@code memberCredentials.sshAuthSock}.
|
||||
*/
|
||||
@Test
|
||||
void sshAuthSockStaysBlockedEvenWhenListedInMemberCredentialsAllow() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
WiringLauncher launcher = new WiringLauncher(herdr,
|
||||
allowListWithAllow(List.of("SSH_AUTH_SOCK")));
|
||||
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
Path dir = Path.of(launcher.env.get("ZDOTDIR"));
|
||||
String scrub = readAll(dir.resolve(EnvAllowListScrub.SCRUB_FILE));
|
||||
assertFalse(scrub.contains("'SSH_AUTH_SOCK'"),
|
||||
"SSH_AUTH_SOCK must not be on the derived allow-list just because the operator put "
|
||||
+ "it under allow: — sshAuthSock is unset here, so it defaults to block");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up criterion 3: on every allow-list spawn the daemon logs one INFO line, shaped
|
||||
* "member credentials: allowed N of M", with real counts — not constants. Real path: the count
|
||||
* is asserted after a real {@link HerdrPeerLauncher#spawn} call, reading the log the production
|
||||
* code actually emits.
|
||||
*/
|
||||
@Test
|
||||
void logsAnAllowedCountLineAgainstTheHostEnvironmentOnEverySpawn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
// INJECTED is a key of this launch's own env map, so it always survives; the other two are
|
||||
// neither derived from the profile nor configured anywhere, so they are blanked. Real
|
||||
// N=1 (INJECTED), real M=3 (all three names) — neither number is hardcoded in the assertion
|
||||
// by coincidence, they follow directly from this fixture.
|
||||
Set<String> hostEnvNames = Set.of(INJECTED, "SOME_UNRELATED_NAME", "ANOTHER_UNRELATED_NAME");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/zsh", () -> hostEnvNames);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
assertTrue(appender.list.stream()
|
||||
.anyMatch(e -> "member credentials: allowed 1 of 3".equals(e.getFormattedMessage())),
|
||||
"expected 'member credentials: allowed 1 of 3', got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
/**
|
||||
* Lead-review fix: on a NON-zsh shell no scrub ever runs (bash ignores {@code ZDOTDIR}), so the
|
||||
* "allowed N of M" line — which describes what the scrub does — must not be printed there either.
|
||||
* Before this fix the line was logged BEFORE the zsh gate, so a non-zsh host printed e.g.
|
||||
* "allowed 1 of 3" while blocking nothing at all, telling an operator a control ran when it did
|
||||
* not. Real path: goes through {@link HerdrPeerLauncher#spawn}, same as the sibling test above,
|
||||
* with the shell fixed to bash so the fallback branch is the one exercised.
|
||||
*/
|
||||
@Test
|
||||
void noAllowedCountLineIsEmittedOnTheNonZshFallbackPath() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Set<String> hostEnvNames = Set.of(INJECTED, "SOME_UNRELATED_NAME", "ANOTHER_UNRELATED_NAME");
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, allowList(), "/bin/bash", () -> hostEnvNames);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
assertFalse(appender.list.stream()
|
||||
.anyMatch(e -> e.getFormattedMessage().startsWith("member credentials: allowed ")),
|
||||
"no scrub runs on a non-zsh shell, so no 'allowed N of M' count may be printed — got: "
|
||||
+ appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList());
|
||||
}
|
||||
|
||||
private static String readAll(Path p) {
|
||||
try {
|
||||
return Files.readString(p);
|
||||
@@ -134,10 +255,16 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
}
|
||||
|
||||
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell) {
|
||||
this(herdr, creds, shell, null);
|
||||
}
|
||||
|
||||
/** Plus an injectable {@code hostEnvNames} source, for the "allowed N of M" log line test. */
|
||||
WiringLauncher(FakeHerdr herdr, Supplier<FleetConfig.MemberCredentials> creds, String shell,
|
||||
Supplier<Set<String>> hostEnvNames) {
|
||||
super("test", new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of("test", profile()), "test",
|
||||
name -> "SHELL".equals(name) ? shell : null,
|
||||
0, () -> 0L, () -> { }, null, creds);
|
||||
0, () -> 0L, () -> { }, null, creds, hostEnvNames);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -88,6 +88,39 @@ class MemberEnvAllowListTest {
|
||||
assertTrue(after.containsAll(Set.of("TOKEN_SECOND", "GIT_TOK", "SECOND_KEY")));
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-633 follow-up: a name that appears ONLY in {@code memberCredentials.allow:} — no profile
|
||||
* mentions it at all — must still survive the derivation. Before this fix {@code derive} never
|
||||
* saw {@code allow:}, so setting {@code policy: allow-list} silently blanked exactly this name.
|
||||
*/
|
||||
@Test
|
||||
void aNameOnlyInMemberCredentialsAllowSurvivesDerivation() {
|
||||
FleetConfig.Profile p = profile("p", "TOKEN_A", null, null, Map.of("KEY_A", "v"));
|
||||
|
||||
Set<String> derived = MemberEnvAllowList.derive(List.of(p), Set.of("OPERATOR_ONLY_NAME"));
|
||||
|
||||
assertTrue(derived.contains("OPERATOR_ONLY_NAME"),
|
||||
"memberCredentials.allow: must be unioned in, not ignored");
|
||||
// and the profile-derived half must still be present — this is a union, not a replacement.
|
||||
assertTrue(derived.contains("KEY_A"));
|
||||
assertTrue(derived.contains("TOKEN_A"));
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code SSH_AUTH_SOCK} is a live handle to the operator's ssh-agent, never a value — so it must
|
||||
* stay excluded from the derived set even when the operator lists it under {@code allow:} for an
|
||||
* unrelated reason. It is governed ONLY by {@code memberCredentials.sshAuthSock}, applied
|
||||
* separately by the caller ({@code HerdrPeerLauncher}).
|
||||
*/
|
||||
@Test
|
||||
void sshAuthSockInMemberCredentialsAllowIsStillExcluded() {
|
||||
Set<String> derived = MemberEnvAllowList.derive(List.of(), Set.of("SSH_AUTH_SOCK", "OTHER_NAME"));
|
||||
|
||||
assertFalse(derived.contains("SSH_AUTH_SOCK"),
|
||||
"SSH_AUTH_SOCK must never ride in on the generic allow: list");
|
||||
assertTrue(derived.contains("OTHER_NAME"), "other allow: names are unaffected");
|
||||
}
|
||||
|
||||
/** {@code LC_*} categories are infrastructure by prefix; everything else needs an exact match. */
|
||||
@Test
|
||||
void keepsMatchesExactlyPlusTheLocalePrefixRule() {
|
||||
|
||||
@@ -0,0 +1,149 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import com.rabbitmq.client.AMQP;
|
||||
import com.rabbitmq.client.Channel;
|
||||
import com.rabbitmq.client.Connection;
|
||||
import com.rabbitmq.client.DeliverCallback;
|
||||
import com.rabbitmq.client.Delivery;
|
||||
import com.rabbitmq.client.Envelope;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.lang.reflect.InvocationHandler;
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.ArrayDeque;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
|
||||
/**
|
||||
* Pins the manual-ack prefetch behaviour without a broker. The fake channel models a broker that
|
||||
* sends no more than its QoS window of unacked deliveries. If {@link AmqpReplyInbox} starts acking
|
||||
* messages while it adds them to {@code held}, this test drains the whole fake queue instead.
|
||||
*/
|
||||
class AmqpReplyInboxPrefetchTest {
|
||||
|
||||
@Test
|
||||
void unackedDeliveriesKeepTheHeldBacklogAtThePrefetchWindow() {
|
||||
int prefetch = 3;
|
||||
int published = 8;
|
||||
PrefetchBroker broker = new PrefetchBroker();
|
||||
for (int i = 0; i < published; i++) {
|
||||
broker.publish("m" + i, "payload " + i);
|
||||
}
|
||||
|
||||
try (AmqpReplyInbox inbox = new AmqpReplyInbox(connectionFor(broker.channel()), prefetch)) {
|
||||
inbox.own("worker");
|
||||
|
||||
assertEquals(prefetch, inbox.peek("worker").size(),
|
||||
"held messages must stop at the unacked prefetch window");
|
||||
assertEquals(published - prefetch, broker.queuedCount(),
|
||||
"messages beyond the window must remain on the broker");
|
||||
assertEquals(0, broker.ackCount(), "receipt must not ack a held message");
|
||||
|
||||
inbox.ack("worker", "m0");
|
||||
|
||||
assertEquals(prefetch, inbox.peek("worker").size(),
|
||||
"one caller ack frees exactly one slot for the broker");
|
||||
assertEquals(published - prefetch - 1, broker.queuedCount(),
|
||||
"only one queued message may enter after one caller ack");
|
||||
assertEquals(1, broker.ackCount(), "only the caller ack may reach the broker");
|
||||
}
|
||||
}
|
||||
|
||||
private static Connection connectionFor(Channel consumeChannel) {
|
||||
Channel publishChannel = (Channel) Proxy.newProxyInstance(
|
||||
AmqpReplyInboxPrefetchTest.class.getClassLoader(), new Class<?>[] {Channel.class},
|
||||
(proxy, method, args) -> defaultValue(method.getReturnType()));
|
||||
AtomicInteger channelCalls = new AtomicInteger();
|
||||
InvocationHandler handler = (proxy, method, args) -> {
|
||||
if (method.getName().equals("createChannel") && (args == null || args.length == 0)) {
|
||||
return channelCalls.getAndIncrement() == 0 ? consumeChannel : publishChannel;
|
||||
}
|
||||
return defaultValue(method.getReturnType());
|
||||
};
|
||||
return (Connection) Proxy.newProxyInstance(AmqpReplyInboxPrefetchTest.class.getClassLoader(),
|
||||
new Class<?>[] {Connection.class}, handler);
|
||||
}
|
||||
|
||||
private static final class PrefetchBroker implements InvocationHandler {
|
||||
private final ArrayDeque<Delivery> queued = new ArrayDeque<>();
|
||||
private final Map<Long, Delivery> unacked = new LinkedHashMap<>();
|
||||
private DeliverCallback consumer;
|
||||
private int prefetch;
|
||||
private int acks;
|
||||
private long nextTag = 1;
|
||||
|
||||
Channel channel() {
|
||||
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxPrefetchTest.class.getClassLoader(),
|
||||
new Class<?>[] {Channel.class}, this);
|
||||
}
|
||||
|
||||
void publish(String msgId, String content) {
|
||||
queued.add(new Delivery(new Envelope(nextTag++, false, "", ""),
|
||||
new AMQP.BasicProperties.Builder().messageId(msgId).build(),
|
||||
content.getBytes(StandardCharsets.UTF_8)));
|
||||
}
|
||||
|
||||
int queuedCount() {
|
||||
return queued.size();
|
||||
}
|
||||
|
||||
int ackCount() {
|
||||
return acks;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Object invoke(Object proxy, java.lang.reflect.Method method, Object[] args) throws IOException {
|
||||
switch (method.getName()) {
|
||||
case "basicQos" -> {
|
||||
prefetch = (int) args[0];
|
||||
return null;
|
||||
}
|
||||
case "basicConsume" -> {
|
||||
assertFalse((boolean) args[1], "the inbox consumer must use manual acknowledgements");
|
||||
consumer = (DeliverCallback) args[2];
|
||||
deliverAvailable();
|
||||
return "consumer";
|
||||
}
|
||||
case "basicAck" -> {
|
||||
unacked.remove((long) args[0]);
|
||||
acks++;
|
||||
deliverAvailable();
|
||||
return null;
|
||||
}
|
||||
default -> {
|
||||
return defaultValue(method.getReturnType());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private void deliverAvailable() throws IOException {
|
||||
while (consumer != null && unacked.size() < prefetch && !queued.isEmpty()) {
|
||||
Delivery delivery = queued.removeFirst();
|
||||
unacked.put(delivery.getEnvelope().getDeliveryTag(), delivery);
|
||||
consumer.handle("consumer", delivery);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static Object defaultValue(Class<?> type) {
|
||||
if (!type.isPrimitive() || type == void.class) {
|
||||
return null;
|
||||
}
|
||||
if (type == boolean.class) {
|
||||
return false;
|
||||
}
|
||||
if (type == long.class) {
|
||||
return 0L;
|
||||
}
|
||||
if (type == int.class) {
|
||||
return 0;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
@@ -33,8 +33,17 @@ class MessageServiceTest {
|
||||
private final FakeHerdr herdr = new FakeHerdr().readText("BUILD GREEN: 391 files");
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
private final Rendezvous rendezvous = new Rendezvous();
|
||||
/**
|
||||
* fleetd#164: this fixture drives a delivery and its completion back-to-back with no real time
|
||||
* between them, so the real clock would trip {@link CompletionResolver#MIN_TURN_NANOS} on every
|
||||
* completion-fallback test here. An ever-advancing fake clock stands in for the model-latency and
|
||||
* herdr round-trips a real turn would spend, so each delivery-then-resolve pair still lands
|
||||
* outside the floor.
|
||||
*/
|
||||
private final java.util.concurrent.atomic.AtomicLong resolverClock = new java.util.concurrent.atomic.AtomicLong();
|
||||
private final CompletionResolver completion =
|
||||
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
new CompletionResolver(agents, rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none(),
|
||||
() -> resolverClock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1));
|
||||
private final Injector injector = new Injector(agents, completion);
|
||||
private final InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
private final MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
||||
@@ -1081,4 +1090,147 @@ class MessageServiceTest {
|
||||
assertEquals(phase, view.phase());
|
||||
return view;
|
||||
}
|
||||
|
||||
// --- CB-640: fleet health evidence accessors --------------------------------------------
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryIsFalseForAnUnknownTarget() {
|
||||
assertFalse(messages.hasQueuedDelivery("nobody-ever-sent-here"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryIsFalseBeforeAnyTimeout() {
|
||||
assertFalse(messages.hasQueuedDelivery(T));
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryIsTrueAfterAnUndeliveredSendTimesOut() {
|
||||
// Nothing ever delivers the message (never goes IDLE/BLOCKED), so the send times out with
|
||||
// TIMED_OUT_QUEUED — same setup as sendTimesOutBeforeDeliveryIsQueuedNotWorking above.
|
||||
MessageService.Reply r = messages.send(T, "never delivered", 50);
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome());
|
||||
|
||||
assertTrue(messages.hasQueuedDelivery(T),
|
||||
"a TIMED_OUT_QUEUED send leaves the message still queued in the injector");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
|
||||
assertTrue(messages.hasQueuedDelivery(T));
|
||||
|
||||
// A fresh send accepts delivery (opens its own waiter) — the stale queued fact is cleared.
|
||||
CompletableFuture<MessageService.Reply> second = sendAsync();
|
||||
awaitWaiting();
|
||||
assertFalse(messages.hasQueuedDelivery(T),
|
||||
"a fresh accepted delivery supersedes the earlier queued fact");
|
||||
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
assertTrue(rendezvous.resolve(T, "second done"));
|
||||
second.get(5, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasQueuedDeliveryClearsOnAbandon() {
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, messages.send(T, "first", 50).outcome());
|
||||
assertTrue(messages.hasQueuedDelivery(T));
|
||||
|
||||
messages.abandon(T, "session released");
|
||||
assertFalse(messages.hasQueuedDelivery(T), "a torn-down target has nothing left queued for it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyIsFalseForAnUnknownTarget() {
|
||||
assertFalse(messages.hasStrandedReply("nobody-ever-sent-here"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyIsFalseWhenTheReplyResolvedALiveSend() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
awaitWaiting();
|
||||
|
||||
assertTrue(messages.reply(T, "resolved-live"));
|
||||
assertFalse(messages.hasStrandedReply(T), "a reply that resolved an open send is not stranded");
|
||||
|
||||
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.REPLIED, r.outcome());
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyIsTrueWhenNoSendWasWaiting() {
|
||||
// No send is open for T — the reply queues into the inbox and is recorded as stranded.
|
||||
assertTrue(messages.reply(T, "nobody was waiting"));
|
||||
assertTrue(messages.hasStrandedReply(T),
|
||||
"a reply with no open send strands, even though it is safely queued in the inbox");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyClearsOnceTheTargetsNextDeliveryIsAccepted() throws Exception {
|
||||
assertTrue(messages.reply(T, "stray"));
|
||||
assertTrue(messages.hasStrandedReply(T));
|
||||
|
||||
// The next accepted delivery for T clears the stale stranding fact — the one case the
|
||||
// ticket calls out as the one that matters.
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
awaitWaiting();
|
||||
assertFalse(messages.hasStrandedReply(T),
|
||||
"a stranded reply must clear once the target's delivery is accepted again");
|
||||
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
send.get(5, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasStrandedReplyClearsOnAbandon() {
|
||||
assertTrue(messages.reply(T, "stray"));
|
||||
assertTrue(messages.hasStrandedReply(T));
|
||||
|
||||
messages.abandon(T, "session released");
|
||||
assertFalse(messages.hasStrandedReply(T), "a torn-down target has nothing left to strand");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasOrphanedDelegationIsFalseForAnUnknownTarget() {
|
||||
assertFalse(messages.hasOrphanedDelegation("nobody-ever-sent-here"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasOrphanedDelegationIsFalseWhilePendingTicketsHaveAnAcceptedDelivery() throws Exception {
|
||||
// Mirrors abandonFailsEveryPendingAsyncTicketForTheReleasedTarget above: "first" holds the
|
||||
// session lock and its waiter is open, so the target genuinely has something in flight even
|
||||
// though "second" and "third" are themselves parked (PENDING) behind the lock.
|
||||
messages.sendAsync(T, "first task");
|
||||
awaitWaiting(); // first task owns the target lock and rendezvous waiter
|
||||
String second = messages.sendAsync(T, "second task");
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(second).phase());
|
||||
|
||||
assertFalse(messages.hasOrphanedDelegation(T),
|
||||
"the target has an accepted delivery in flight (first), so nothing here is orphaned");
|
||||
|
||||
assertTrue(messages.abandon(T, "session released")); // release the lock and the parked tickets
|
||||
}
|
||||
|
||||
@Test
|
||||
void hasOrphanedDelegationIsTrueOnceAnUnansweredAskLapsesBackToPending() throws Exception {
|
||||
// Same setup as unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget above:
|
||||
// once the ask lapses, the ticket goes back to PENDING but send() already closed the
|
||||
// forward waiter the instant the question surfaced — nothing is left in flight for T.
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT,
|
||||
messages.ask(T, "which config?", 200).outcome());
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
|
||||
|
||||
assertFalse(messages.hasAcceptedDelivery(T), "the forward waiter closed when the question surfaced");
|
||||
assertFalse(messages.hasQueuedDelivery(T), "this ticket never timed out as queued");
|
||||
assertTrue(messages.hasOrphanedDelegation(T),
|
||||
"a PENDING ticket with no accepted or queued delivery for its target is orphaned");
|
||||
|
||||
assertTrue(messages.abandon(T, "session released")); // clean up the still-open ticket
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,6 +32,7 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
import java.util.function.Predicate;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -64,6 +65,11 @@ class FleetAppTest {
|
||||
}
|
||||
|
||||
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement, Worktrees worktrees) {
|
||||
return start(herdr, workerBaseUrl, allow, placement, worktrees, ignored -> false);
|
||||
}
|
||||
|
||||
private int start(FakeHerdr herdr, String workerBaseUrl, Set<String> allow, String placement,
|
||||
Worktrees worktrees, Predicate<String> deliverable) {
|
||||
FleetConfig.Profile wcfg = new FleetConfig.Profile(
|
||||
"ltms-local", workerBaseUrl, "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
placement, "fleet", "worker: {profile} #{n}", null, null, null);
|
||||
@@ -84,7 +90,8 @@ class FleetAppTest {
|
||||
// it directly so the inbox contract holds for those endpoints.
|
||||
inbox.own("term_a");
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous, inbox);
|
||||
app = new FleetApp(herdr, workers, sessions, messages, this.presence, null)
|
||||
app = new FleetApp(herdr, workers, sessions, messages, this.presence, null,
|
||||
null, null, id -> this.presence.isPresent(id) || deliverable.test(id))
|
||||
.build().start("127.0.0.1", 0);
|
||||
return app.port();
|
||||
}
|
||||
@@ -484,6 +491,16 @@ class FleetAppTest {
|
||||
assertTrue(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean());
|
||||
}
|
||||
|
||||
@Test
|
||||
void sessionStatusReportsRegisteredLeadAsReady() throws Exception {
|
||||
Map<String, String> leads = Map.of("term_lead", "terra");
|
||||
int port = start(new FakeHerdr(), "http://gx00.gw:8000", Set.of("gx00.gw"), "tab",
|
||||
new GitWorktrees(), leads::containsKey);
|
||||
|
||||
JsonNode body = mapper.readTree(req(port, "GET", "/sessions/term_lead/status").body());
|
||||
assertTrue(body.get("ready").asBoolean());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-582: a lead polling {@code GET /sessions/{id}/status} on its normal cadence — not the
|
||||
* ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code fleet_ask}
|
||||
|
||||
@@ -55,12 +55,21 @@ class GitWorktreesTest {
|
||||
}
|
||||
|
||||
private static void git(Path cwd, String... args) throws Exception {
|
||||
gitOutput(cwd, args);
|
||||
}
|
||||
|
||||
private static String gitOutput(Path cwd, String... args) throws Exception {
|
||||
List<String> cmd = new java.util.ArrayList<>(List.of("git"));
|
||||
cmd.addAll(List.of(args));
|
||||
Process p = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true).start();
|
||||
ProcessBuilder pb = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true);
|
||||
pb.environment().put("GIT_CONFIG_GLOBAL", "/dev/null");
|
||||
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
|
||||
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
|
||||
Process p = pb.start();
|
||||
String out = new String(p.getInputStream().readAllBytes());
|
||||
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: " + String.join(" ", cmd));
|
||||
assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out);
|
||||
return out;
|
||||
}
|
||||
|
||||
/** Pending changes to {@code file} in {@code cwd}, empty when git considers it unmodified. */
|
||||
@@ -205,6 +214,144 @@ class GitWorktreesTest {
|
||||
"expected an explicitly empty server map, got:\n" + body);
|
||||
}
|
||||
|
||||
/**
|
||||
* The worktree command shares the primary checkout's config, so this checks the URL git actually
|
||||
* reads after {@link GitWorktrees#add}, rather than checking only a URL formatting helper.
|
||||
*/
|
||||
@Test
|
||||
void aProvisionedWorktreeUsesACleanHttpsOrigin(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://synthetic-test-token@git.ltms.dev/akb/kb.git");
|
||||
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "fleetd-157-safe-origin", "HEAD");
|
||||
|
||||
Path worktree = Path.of(wt);
|
||||
String origin = gitOutput(worktree, "config", "--get", "remote.origin.url").trim();
|
||||
assertEquals("https://git.ltms.dev/akb/kb.git", origin);
|
||||
assertFalse(origin.contains("synthetic-test-token"), "provisioned worktree kept user info");
|
||||
assertFalse(gitOutput(worktree, "remote", "-v").contains("synthetic-test-token"),
|
||||
"git remote -v exposed user info");
|
||||
assertFalse(gitOutput(worktree, "config", "--list").contains("synthetic-test-token"),
|
||||
"git config --list exposed user info");
|
||||
String helper = gitOutput(worktree, "config", "--worktree", "--get", "credential.helper");
|
||||
assertTrue(helper.contains("WORKER_GITEA_TOKEN"), "credential helper does not read the member environment");
|
||||
assertFalse(helper.contains("synthetic-test-token"), "credential helper stored user info");
|
||||
}
|
||||
|
||||
@Test
|
||||
void worktreeCredentialHelperCompletesWithoutUsingAnInheritedHelper(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "fleetd-157-helper", "HEAD");
|
||||
Path globalConfig = tmp.resolve("global.gitconfig");
|
||||
Files.writeString(globalConfig, """
|
||||
[credential]
|
||||
helper = !f() { printf 'username=%s\\npassword=%s\\n\\n' operator operator-secret; }; f
|
||||
""");
|
||||
|
||||
ProcessBuilder pb = new ProcessBuilder("git", "credential", "fill")
|
||||
.directory(Path.of(wt).toFile()).redirectErrorStream(true);
|
||||
pb.environment().put("GIT_CONFIG_GLOBAL", globalConfig.toString());
|
||||
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
|
||||
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
|
||||
pb.environment().put("WORKER_GITEA_TOKEN", "synthetic-worker-value");
|
||||
Process p = pb.start();
|
||||
p.getOutputStream().write("protocol=https\nhost=git.ltms.dev\n\n".getBytes(StandardCharsets.UTF_8));
|
||||
p.getOutputStream().close();
|
||||
String credential = new String(p.getInputStream().readAllBytes(), StandardCharsets.UTF_8);
|
||||
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git credential fill timed out");
|
||||
assertEquals(0, p.exitValue(), "git credential fill failed");
|
||||
assertTrue(credential.contains("username=git"), "helper did not return its fixed username");
|
||||
assertTrue(credential.contains("password=synthetic-worker-value"),
|
||||
"helper did not return the worker token as the password");
|
||||
assertFalse(credential.contains("operator-secret"), "Git used the inherited global helper");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #157 follow-up. An SSH origin never consults {@code credential.helper} — the environment
|
||||
* credential helper set by {@link GitWorktrees#add} is therefore useless when a member's origin is
|
||||
* SSH, which is exactly this repo's shape. The worktree must instead get a worktree-scoped
|
||||
* {@code url.<https>.insteadOf <ssh>} rewrite so both fetch and push resolve to HTTPS, while the
|
||||
* parent checkout — sharing the same repo-level origin config — must resolve the original SSH URL
|
||||
* completely unchanged. The host/port here (a synthetic {@code forge.example.test:2222}, not
|
||||
* {@code git.ltms.dev}) proves the rewrite is derived from the origin, not a hardcoded constant.
|
||||
*/
|
||||
@Test
|
||||
void aProvisionedWorktreeRewritesAnSshOriginToHttpsWorktreeScoped(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "ssh://git@forge.example.test:2222/acme/proj.git");
|
||||
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "fleetd-157-ssh-rewrite", "HEAD");
|
||||
|
||||
Path worktree = Path.of(wt);
|
||||
assertEquals("https://forge.example.test/acme/proj.git",
|
||||
gitOutput(worktree, "remote", "get-url", "origin").trim(),
|
||||
"worktree fetch URL was not rewritten to HTTPS");
|
||||
assertEquals("https://forge.example.test/acme/proj.git",
|
||||
gitOutput(worktree, "remote", "get-url", "--push", "origin").trim(),
|
||||
"worktree push URL was not rewritten to HTTPS");
|
||||
// The raw config value is unchanged — only the resolved URL is rewritten, via insteadOf.
|
||||
assertEquals("ssh://git@forge.example.test:2222/acme/proj.git",
|
||||
gitOutput(worktree, "config", "--get", "remote.origin.url").trim());
|
||||
|
||||
assertEquals("ssh://git@forge.example.test:2222/acme/proj.git",
|
||||
gitOutput(repo, "remote", "get-url", "origin").trim(),
|
||||
"the parent checkout's fetch URL must be untouched");
|
||||
assertEquals("ssh://git@forge.example.test:2222/acme/proj.git",
|
||||
gitOutput(repo, "remote", "get-url", "--push", "origin").trim(),
|
||||
"the parent checkout's push URL must be untouched");
|
||||
}
|
||||
|
||||
/** An origin already on HTTPS is left alone — the environment credential helper already covers it. */
|
||||
@Test
|
||||
void aProvisionedWorktreeLeavesAnHttpsOriginAlone(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
|
||||
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "fleetd-157-https-noop", "HEAD");
|
||||
|
||||
Path worktree = Path.of(wt);
|
||||
assertEquals("https://git.ltms.dev/akb/kb.git",
|
||||
gitOutput(worktree, "remote", "get-url", "origin").trim());
|
||||
assertEquals(1, exitCode("git", "-C", wt, "config", "--worktree", "--get-regexp", "^url\\."),
|
||||
"no url.*.insteadOf rewrite should be added for an already-HTTPS origin");
|
||||
}
|
||||
|
||||
/** Test-local exit-code probe, mirroring {@link GitWorktrees#exitCode} for an assertion the
|
||||
* production class does not expose. */
|
||||
private static int exitCode(String... command) throws Exception {
|
||||
ProcessBuilder pb = new ProcessBuilder(command).redirectErrorStream(true);
|
||||
pb.environment().put("GIT_CONFIG_GLOBAL", "/dev/null");
|
||||
pb.environment().put("GIT_CONFIG_SYSTEM", "/dev/null");
|
||||
pb.environment().put("GIT_TERMINAL_PROMPT", "0");
|
||||
Process p = pb.start();
|
||||
p.getInputStream().readAllBytes();
|
||||
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "command timed out: " + String.join(" ", command));
|
||||
return p.exitValue();
|
||||
}
|
||||
|
||||
@Test
|
||||
void provisioningRefusesAWorktreeWhoseOriginStillHasHttpsUserInfo(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
git(repo, "remote", "add", "origin", "https://git.ltms.dev/akb/kb.git");
|
||||
GitWorktrees worktrees = new GitWorktrees(tmp.resolve("wts").toString(), worktreePath -> {
|
||||
try {
|
||||
git(Path.of(worktreePath), "remote", "set-url", "origin",
|
||||
"https://synthetic-test-token@git.ltms.dev/akb/kb.git");
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
});
|
||||
|
||||
WorktreeException error = assertThrows(WorktreeException.class,
|
||||
() -> worktrees.add(repo.toString(), "fleetd-157-refuse-origin", "HEAD"));
|
||||
assertEquals("worktree origin contains HTTPS user info; refusing provision", error.getMessage());
|
||||
}
|
||||
|
||||
/** Neutralizing must not look like work in progress, or a worker would commit it into its PR. */
|
||||
@Test
|
||||
void theNeutralizedConfigIsNotAPendingLocalModification(@TempDir Path tmp) throws Exception {
|
||||
|
||||
Reference in New Issue
Block a user