Compare commits
66 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| f188947750 | |||
| 26f380a00b | |||
| b8182c96c2 | |||
| 3da44eed63 | |||
| 93a9ed3f83 | |||
| 0b032f5a1a | |||
| fad99c4c5e | |||
| 87871eaefb | |||
| a476a14f1c | |||
| 5eb4267a4a | |||
| 7611b69667 | |||
| cc302fe4af | |||
| 4a8a780274 | |||
| 343ce0f4c0 | |||
| cec3e191d4 | |||
| 1850a5f324 | |||
| e4eb3dbed4 | |||
| 57cd96f5e6 | |||
| 202e37e3b3 | |||
| 90253f832d | |||
| f1640f5dcc | |||
| c7903c1efe | |||
| 7d711942fe | |||
| bec87f987c | |||
| 7f8a8829f9 | |||
| af9589783e | |||
| 190436c9cf | |||
| 8ea5c2bb1f | |||
| 8335b12562 | |||
| d25c863118 | |||
| 7c34e8f4f9 | |||
| be07ed2033 | |||
| a6415f3e52 | |||
| bb6fc9e0d7 | |||
| 3366590dbe | |||
| 01adc841fa | |||
| 08771e270b | |||
| b5843ab43f | |||
| d8985719eb | |||
| 5b1e13ca3d | |||
| c89a375e5d | |||
| 37dcefa834 | |||
| 6a7342b1f0 | |||
| a5ad7c6561 | |||
| c71ac231e5 | |||
| 33720c42b3 | |||
| 3833d8e52b | |||
| b37def9238 | |||
| 8f02576df6 | |||
| 525bc1c5f4 | |||
| d59ece6dec | |||
| 32408d1e64 | |||
| 6e23bf8309 | |||
| aa4c0b84c3 | |||
| 40c593cd09 | |||
| 979adf82eb | |||
| 36870836aa | |||
| 136312fb11 | |||
| 4f28da62a3 | |||
| 708f1795ad | |||
| 599419f9e6 | |||
| ac351ee1de | |||
| 274afafde6 | |||
| 8f59019305 | |||
| b17f37a683 | |||
| dcd505286f |
@@ -39,6 +39,10 @@ jobs:
|
||||
# that runs does so against the fake UDS herdr and fake ccs/claude stubs.
|
||||
run: mvn -B clean install
|
||||
|
||||
- name: javadoc reference lint
|
||||
working-directory: fleetd
|
||||
run: mvn -B -DskipTests javadoc:javadoc -Ddoclint=reference
|
||||
|
||||
# Deliberately NOT actions/upload-artifact: this Gitea instance presents as GHES, and
|
||||
# @actions/artifact v2+ (i.e. upload-artifact@v4) refuses to run there —
|
||||
# "GHESNotSupportedError ... not currently supported on GHES", which red-Xes an otherwise
|
||||
@@ -54,6 +58,23 @@ jobs:
|
||||
done
|
||||
exit 0
|
||||
|
||||
# fleetd #550 — nothing ran scripts/test-redeploy-fleetd.sh in CI before this, on any platform,
|
||||
# so it had run only on macOS by hand and two Linux-only bugs (this issue's items 1 and 2)
|
||||
# survived undetected: shasum is a macOS-only tool (it ships with Perl; GNU coreutils, i.e. every
|
||||
# mainstream Linux distro including this runner's ubuntu-latest, does not have it and ships
|
||||
# sha256sum instead). The gate here is the step's own exit code, nothing else: a `run:` step in
|
||||
# Gitea/GitHub Actions already fails the job on a non-zero exit with no extra scripting needed,
|
||||
# so this deliberately does NOT grep the output for a `FAIL:` count. That is the #550 item-2
|
||||
# lesson one level up — a suite that dies before it runs a single test prints zero FAIL lines,
|
||||
# which is exactly what a clean pass also prints, so counting FAIL lines can never be the gate.
|
||||
shell-tests:
|
||||
runs-on: ubuntu-latest
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
|
||||
- name: redeploy-fleetd.sh shell suite
|
||||
run: bash scripts/test-redeploy-fleetd.sh
|
||||
|
||||
# CB-521 — actually run the AMQP contract test in CI, against a REAL broker. The broker is a
|
||||
# RabbitMQ SERVICE CONTAINER, not Testcontainers-with-Docker: the runner image has no Docker, so
|
||||
# AmqpReplyInboxContractTest reads AMQP_URI (set below to the service's network alias) and binds
|
||||
|
||||
@@ -96,8 +96,10 @@ below are the procedure — run them in order, every task, not only the big ones
|
||||
that answers it. **A worker's ask waits ~55 seconds, and no nudge makes that longer** — so never
|
||||
brief a worker to "ask me". Decide before you delegate, or give it an explicit default.
|
||||
6. **Verify yourself.** Re-run the build and the checks. A worker cannot run your IDE tooling, any
|
||||
forge tools it appears to have hold a blocked credential and fail, and a piped command
|
||||
(`… | tail`) hides failures behind a zero exit — never promote a worker's "clean" to a fact.
|
||||
forge MCP server it appears to have holds a blocked credential and fails every call, and a piped
|
||||
command (`… | tail`) hides failures behind a zero exit — never promote a worker's "clean" to a
|
||||
fact. Its injected repo-scoped `GITEA_TOKEN` is a different credential and does work, so a worker
|
||||
reporting that it opened its own PR is reporting something it really can do.
|
||||
7. **Review — fan out.** Spawn reviewers against the diff, one per dimension or per file, with
|
||||
`wait:false`. Never the implementer of the scope it reviews, and brief them from the diff — not
|
||||
from the implementer's rationale, which carries its own blind spot. Dispatch each PR's reviewers
|
||||
@@ -200,8 +202,11 @@ simply complies has thrown away the reason there are two of you.
|
||||
assume them.** What you mount depends on your backend: an opencode member gets the bridge and
|
||||
nothing else, while a Claude Code member also inherits the operator's user-scope MCP servers,
|
||||
which the bridge never chose for you. Two rules follow. The primary's IDE tooling is still not
|
||||
yours, whatever you see. And **a mounted tool is not a working tool** — the forge server you may
|
||||
find there holds a deliberately blocked credential and fails every call, by design.
|
||||
yours, whatever you see. And **a mounted tool is not a working tool** — the forge MCP server you
|
||||
may find there holds a deliberately blocked credential and fails every call, by design. That is
|
||||
not your only forge route, and the two must not be confused: the repo-scoped `GITEA_TOKEN` the
|
||||
daemon injects into your environment does work, and using it to open your own PR is part of the
|
||||
job. A blocked MCP tool is never a reason to skip that step.
|
||||
6. **Never merge.** Stage files explicitly — never `git add -A` — and leave alone anything the
|
||||
project marks as not-yours-to-commit.
|
||||
|
||||
|
||||
@@ -80,6 +80,7 @@ import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Predicate;
|
||||
import java.util.function.Supplier;
|
||||
import java.util.regex.Pattern;
|
||||
@@ -270,15 +271,12 @@ public final class Fleetd {
|
||||
// the first thing that actually talks to herdr, so without this wait a boot-order race
|
||||
// would crash the daemon into a restart loop. Wait, then degrade rather than die: serving
|
||||
// with /healthz reporting "degraded" is strictly more useful than exiting.
|
||||
boolean herdrUp = awaitHerdr(herdr);
|
||||
HerdrAwaitOutcome herdrOutcome = awaitHerdr(herdr, System::nanoTime, Fleetd::sleepHerdrPoll);
|
||||
boolean herdrUp = logHerdrWaitOutcomeAndShouldReap(herdrOutcome);
|
||||
if (herdrUp) {
|
||||
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died
|
||||
// with the previous process — reap those leaked orphans now, before we start serving.
|
||||
workers.reapOrphanWorkers();
|
||||
} else {
|
||||
log.warn("herdr did not answer within {}s — starting anyway; /healthz will report "
|
||||
+ "degraded until it comes up. Orphaned worker panes (if any) were NOT reaped.",
|
||||
HERDR_WAIT_SECONDS);
|
||||
}
|
||||
|
||||
// CB-301: authoritative session registry + lifecycle FSM on top of ClaudeCodeLauncher.
|
||||
@@ -696,13 +694,9 @@ public final class Fleetd {
|
||||
}, outagePolicy);
|
||||
|
||||
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
|
||||
primaryRegistry, callers, metrics,
|
||||
primaryRegistry, callers, FleetMcp.AuthorizationMode.ENFORCED, metrics,
|
||||
capacitySource(config, cfg, profile -> liveCountRef.get().apply(profile)),
|
||||
new FleetMcp.HealthCoverageSource(() -> {
|
||||
var health = config.get().health();
|
||||
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
|
||||
health != null && health.notifications() != null && health.notifications().configured());
|
||||
}),
|
||||
healthCoverageSource(config),
|
||||
quarantineSource,
|
||||
leadMailbox,
|
||||
outageSource,
|
||||
@@ -1040,6 +1034,38 @@ public final class Fleetd {
|
||||
cfg.profiles()::keySet, System::nanoTime);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #426: package-private factory for {@code fleet_list}'s {@code healthCoverage} source,
|
||||
* extracted out of {@code main} for the same reason {@link #capacitySource} and {@link
|
||||
* #quarantineSource} were — and the same reason {@link #exhaustedPatternCoverageLine}/{@link
|
||||
* #errorPatternCoverageLine} exist: {@link FleetHealthMonitor#coverage}'s three-branch method
|
||||
* is easy to pin directly (a plain {@code (boolean, boolean) -> String} call), but that proves
|
||||
* nothing about whether <em>this call site</em> pairs the right boolean with the right meaning.
|
||||
* fleetd #415's measured lesson is the reason this matters here — swapping the two arguments at
|
||||
* a call site like this one compiled clean and left the full suite green, because every existing
|
||||
* test exercised the method in both directions without ever exercising the pairing.
|
||||
*
|
||||
* <p>{@code #407}'s "keep the config invalid, assert on the log line before the throw" option
|
||||
* does not apply to this call site: the five reporters #407 covers all run in {@code main}
|
||||
* <em>before</em> {@code cfg.validateAll()} (line ~171), so an invalid config still exercises
|
||||
* them. This call site is built during {@code FleetMcp} construction, which runs only after
|
||||
* {@code UnixSocketHerdrClient.connect} has already opened a real herdr socket (line ~188) —
|
||||
* reaching it at all means main already performed real I/O, which the no-socket constraint on
|
||||
* this ticket rules out. So the pairing is pinned by extracting it to this directly-callable
|
||||
* factory instead, the same shape {@link #capacitySource}/{@link #quarantineSource} already use.
|
||||
*
|
||||
* <p>Reads {@code config.get().health()} live (health.notifications is a {@code SPLIT_KEYS}
|
||||
* entry — see {@link ConfigRef#SPLIT_KEYS}), so a hot-reloaded notifications block changes what
|
||||
* {@code fleet_list} reports without a restart, exactly like {@link #capacitySource}'s maxLoad.
|
||||
*/
|
||||
static FleetMcp.HealthCoverageSource healthCoverageSource(ConfigRef config) {
|
||||
return new FleetMcp.HealthCoverageSource(() -> {
|
||||
var health = config.get().health();
|
||||
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
|
||||
health != null && health.notifications() != null && health.notifications().configured());
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248: package-private factory for the member worktree/branch lookup {@link
|
||||
* CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241).
|
||||
@@ -1696,12 +1722,70 @@ public final class Fleetd {
|
||||
}
|
||||
|
||||
/**
|
||||
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
|
||||
*
|
||||
* @return true if herdr answered, false if it never did
|
||||
* How {@link #awaitHerdr} ended (fleetd #498). The old code returned a bare {@code boolean},
|
||||
* which collapsed two different facts onto the same {@code false}: the configured wait budget
|
||||
* genuinely running out, and the waiting thread being interrupted possibly milliseconds in.
|
||||
* Those need different operator messages — see {@link #logHerdrWaitOutcomeAndShouldReap} — so
|
||||
* this is a third state, not a better number (the same shape fleetd #497 named). Never treat
|
||||
* {@link #INTERRUPTED} as if it were {@link #DEADLINE_PASSED}: only the latter means herdr was
|
||||
* actually given the full {@link #HERDR_WAIT_SECONDS} and still failed to answer.
|
||||
*/
|
||||
private static boolean awaitHerdr(HerdrClient herdr) {
|
||||
long deadline = System.nanoTime() + HERDR_WAIT_SECONDS * 1_000_000_000L;
|
||||
enum HerdrWaitResult {
|
||||
/** herdr answered {@code ping} before the deadline. */
|
||||
ANSWERED,
|
||||
/** the configured {@link #HERDR_WAIT_SECONDS} budget elapsed with no answer. */
|
||||
DEADLINE_PASSED,
|
||||
/**
|
||||
* the waiting thread was interrupted before the budget ran out — a different event from
|
||||
* {@link #DEADLINE_PASSED} and must never be reported as "did not answer within Ns".
|
||||
*/
|
||||
INTERRUPTED
|
||||
}
|
||||
|
||||
/**
|
||||
* The outcome of one {@link #awaitHerdr} call, carrying the MEASURED elapsed wait time
|
||||
* alongside {@link #result}. {@code elapsedNanos} is always measured against the {@code nanos}
|
||||
* supplier passed to {@link #awaitHerdr} — never assume it equals the configured budget, the
|
||||
* same defect fleetd #494 already fixed once in {@code LeadRollover}.
|
||||
*/
|
||||
record HerdrAwaitOutcome(HerdrWaitResult result, long elapsedNanos) {}
|
||||
|
||||
/**
|
||||
* The real per-poll wait {@link #main} passes to {@link #awaitHerdr}: sleep
|
||||
* {@link #HERDR_WAIT_POLL_MILLIS}, and on interruption re-set the thread's interrupt flag
|
||||
* rather than throwing — {@link #awaitHerdr} detects an interruption by checking {@link
|
||||
* Thread#isInterrupted()} right after this returns, so a poller that swallowed the flag
|
||||
* instead of restoring it would make that check silently miss the interruption.
|
||||
*/
|
||||
private static void sleepHerdrPoll() {
|
||||
try {
|
||||
Thread.sleep(HERDR_WAIT_POLL_MILLIS);
|
||||
} catch (InterruptedException ie) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Poll herdr's {@code ping} until it answers, the configured {@link #HERDR_WAIT_SECONDS}
|
||||
* budget elapses, or the waiting thread is interrupted (CB-504, fleetd #498).
|
||||
*
|
||||
* <p>{@code nanos} and {@code poller} are required parameters with no defaulted overload
|
||||
* (fleetd #415's shape: a defaulted overload is a silent survivor a green suite would vouch
|
||||
* for) — the previous version read {@link System#nanoTime()} and called {@link Thread#sleep}
|
||||
* directly, so nothing could drive it from a test. The one production call site in {@link
|
||||
* #main} passes {@code System::nanoTime} and {@link #sleepHerdrPoll}.
|
||||
*
|
||||
* @param nanos a monotonic elapsed-time clock, e.g. {@code System::nanoTime} — never a
|
||||
* wall-clock source, since only elapsed time (not a timestamp) is measured here
|
||||
* @param poller called once per failed ping while the budget remains; must, on an
|
||||
* {@link InterruptedException}, re-set the thread's interrupt flag rather than
|
||||
* throw or swallow it — this method's interruption check reads that flag right
|
||||
* after {@code poller.run()} returns
|
||||
* @return the outcome and the measured elapsed wait time — see {@link HerdrAwaitOutcome}
|
||||
*/
|
||||
static HerdrAwaitOutcome awaitHerdr(HerdrClient herdr, LongSupplier nanos, Runnable poller) {
|
||||
long start = nanos.getAsLong();
|
||||
long deadline = start + HERDR_WAIT_SECONDS * 1_000_000_000L;
|
||||
boolean waited = false;
|
||||
while (true) {
|
||||
try {
|
||||
@@ -1709,25 +1793,57 @@ public final class Fleetd {
|
||||
if (waited) {
|
||||
log.info("herdr is up");
|
||||
}
|
||||
return true;
|
||||
return new HerdrAwaitOutcome(HerdrWaitResult.ANSWERED, nanos.getAsLong() - start);
|
||||
} catch (HerdrException e) {
|
||||
if (System.nanoTime() >= deadline) {
|
||||
return false;
|
||||
if (nanos.getAsLong() >= deadline) {
|
||||
return new HerdrAwaitOutcome(HerdrWaitResult.DEADLINE_PASSED, nanos.getAsLong() - start);
|
||||
}
|
||||
if (!waited) {
|
||||
log.info("waiting up to {}s for the herdr socket…", HERDR_WAIT_SECONDS);
|
||||
waited = true;
|
||||
}
|
||||
try {
|
||||
Thread.sleep(HERDR_WAIT_POLL_MILLIS);
|
||||
} catch (InterruptedException ie) {
|
||||
Thread.currentThread().interrupt();
|
||||
return false;
|
||||
poller.run();
|
||||
if (Thread.currentThread().isInterrupted()) {
|
||||
return new HerdrAwaitOutcome(HerdrWaitResult.INTERRUPTED, nanos.getAsLong() - start);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Log the right message for {@code outcome} — never the configured {@link #HERDR_WAIT_SECONDS}
|
||||
* budget alone, always the measured elapsed time next to it — and say whether {@link #main}
|
||||
* should now reap orphan worker panes (fleetd #498).
|
||||
*
|
||||
* <p>Extracted out of {@link #main} so this decision is drivable from a test: {@link #main}
|
||||
* boots the whole daemon and cannot itself be run in a unit test, but this is the exact,
|
||||
* unmodified code {@link #main} calls for the decision, not a re-derivation of it.
|
||||
*
|
||||
* @return true only for {@link HerdrWaitResult#ANSWERED} — orphan workers are reaped only
|
||||
* then, exactly as before this ticket
|
||||
*/
|
||||
static boolean logHerdrWaitOutcomeAndShouldReap(HerdrAwaitOutcome outcome) {
|
||||
long elapsedMillis = TimeUnit.NANOSECONDS.toMillis(outcome.elapsedNanos());
|
||||
if (outcome.result() == HerdrWaitResult.ANSWERED) {
|
||||
return true;
|
||||
}
|
||||
if (outcome.result() == HerdrWaitResult.DEADLINE_PASSED) {
|
||||
log.warn("herdr did not answer within the configured wait (configured={}s elapsed={}ms) "
|
||||
+ "— starting anyway; /healthz will report degraded until it comes up. Orphaned "
|
||||
+ "worker panes (if any) were NOT reaped.",
|
||||
HERDR_WAIT_SECONDS, elapsedMillis);
|
||||
return false;
|
||||
}
|
||||
// HerdrWaitResult.INTERRUPTED — a different fact from DEADLINE_PASSED (fleetd #498): the
|
||||
// wait was cut short, not exhausted, and must never be reported as "did not answer within
|
||||
// Ns" — that claim would be false and would send an operator to debug herdr for nothing.
|
||||
log.warn("herdr wait was interrupted before the configured wait ran out (configured={}s "
|
||||
+ "elapsed={}ms) — starting anyway; /healthz will report degraded until it comes "
|
||||
+ "up. Orphaned worker panes (if any) were NOT reaped.",
|
||||
HERDR_WAIT_SECONDS, elapsedMillis);
|
||||
return false;
|
||||
}
|
||||
|
||||
private Fleetd() {
|
||||
}
|
||||
}
|
||||
|
||||
@@ -244,7 +244,15 @@ public final class CallerResolver {
|
||||
// already names what happens if that case is handed the primary role: a worker→primary
|
||||
// escalation. So an unresolved caller is refused (ANONYMOUS — the same clean, already-tested
|
||||
// "authenticated as nothing" outcome used everywhere else in this method), never promoted.
|
||||
return isLoopback(remoteAddr) && c.resolved() ? Principal.primary(c.pid()) : Principal.anonymous();
|
||||
//
|
||||
// fleetd #505: the OTHER way a real pid can wrongly reach here with a null terminal — not a
|
||||
// failed lsof lookup, but a herdr error partway through PaneLocator's pane scan. c.resolved()
|
||||
// says nothing about that; it only tests the lsof sentinel (by design — see
|
||||
// ConnectionIdentity.Caller#resolved). c.scanComplete() is the separate signal: a scan that
|
||||
// could not check every pane must not be read as "checked everywhere, no match" — the pane it
|
||||
// could not check might have been the caller's own. So both must hold before this promotes.
|
||||
return isLoopback(remoteAddr) && c.resolved() && c.scanComplete()
|
||||
? Principal.primary(c.pid()) : Principal.anonymous();
|
||||
}
|
||||
|
||||
private boolean presentedTokenMatches(String authorizationHeader) {
|
||||
|
||||
@@ -208,7 +208,7 @@ import java.util.function.Supplier;
|
||||
* that is correctly hot and for a key nobody triaged. Three times now — {@code worktreeGroup} (#323),
|
||||
* {@code primary}/{@code configReload} (#326), and {@code fleet.leaders} sitting in the escape hatch
|
||||
* (#333) — the second kind hid among the first. A top-level coverage checker in the
|
||||
* {@link ConfigRefProfileCoverageTest} shape (one level up, over {@code FleetConfig} itself rather
|
||||
* {@code ConfigRefProfileCoverageTest} shape (one level up, over {@code FleetConfig} itself rather
|
||||
* than {@code FleetConfig.Profile}) proves this file's four classes exhaust the record's components
|
||||
* — see {@code ConfigRefTopLevelCoverageTest}. That test proves the record's <em>shape</em> is fully
|
||||
* triaged; it does NOT prove a {@code SPLIT_KEYS}/{@code COLD_KEYS}/{@code DEFERRED_KEYS} member has
|
||||
|
||||
@@ -1670,7 +1670,7 @@ public record FleetConfig(
|
||||
* <p>A herdr pane runs a login shell that re-sources the operator's own secret store, so a
|
||||
* member inherits every credential the operator's shell holds — measured at 31 names on this
|
||||
* host, of which only one ({@code GITEA_ACCESS_TOKEN}) used to be blocked, and that block was a
|
||||
* single name hardcoded in {@link HerdrPeerLauncher} rather than driven by config (gitea issue
|
||||
* single name hardcoded in {@link dev.ltms.fleet.member.HerdrPeerLauncher} rather than driven by config (gitea issue
|
||||
* #82). This record replaces that hardcoded shadow with a config-driven one.
|
||||
*
|
||||
* <p><b>deny-by-default, not a deny-list.</b> A deny-list (block these specific names, let
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package dev.ltms.fleet.herdr;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
@@ -36,6 +38,8 @@ import java.util.Set;
|
||||
*/
|
||||
public final class PaneLocator {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(PaneLocator.class);
|
||||
|
||||
/**
|
||||
* Bound on how many ancestor generations {@link #ancestorsOf} walks. This runs on every MCP
|
||||
* call, so a cycle or a pathologically deep process tree must not hang identity resolution;
|
||||
@@ -73,22 +77,46 @@ public final class PaneLocator {
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@code terminal_id} of the agent pane whose process tree contains {@code pid}, or
|
||||
* {@code null} if no agent pane on any searched daemon owns it (e.g. the caller is the
|
||||
* primary, or off-host).
|
||||
* The outcome of a {@link #terminalForPid} scan: the {@code terminal_id} of the agent pane
|
||||
* whose process tree contains the pid ({@link #terminal} is {@code null} if none matched),
|
||||
* and whether the scan that produced that answer ran to completion on every daemon searched.
|
||||
*
|
||||
* <p>{@link #complete} is {@code false} exactly when some {@code pane.process_info} call
|
||||
* failed and, despite that, no pane was ever found to own the pid. In that case a {@code null}
|
||||
* {@link #terminal} means "could not tell", not "definitely not a worker" — fleetd #505: a
|
||||
* transient herdr error on the very pane that <em>does</em> own the caller's pid must not read
|
||||
* as a clean negative and fall through to {@code Principal.primary}, the same way #317's
|
||||
* {@code Caller.resolved()} already guards a failed lsof lookup. Callers ({@code
|
||||
* ConnectionIdentity}, {@code CallerResolver}) must refuse rather than promote on an incomplete
|
||||
* scan.
|
||||
*
|
||||
* <p>When a pane genuinely owns the pid, {@link #complete} is {@code true} regardless of
|
||||
* whether some other, unrelated pane failed to answer earlier in the same scan — a positive
|
||||
* match is definitive and does not need every pane to have been checked (a pane that "vanished
|
||||
* mid-scan" but was never the match is still a clean, complete result).
|
||||
*/
|
||||
public String terminalForPid(long pid) {
|
||||
public record Lookup(String terminal, boolean complete) {
|
||||
private static final Lookup NOT_FOUND = new Lookup(null, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve {@code pid} to the agent pane whose process tree contains it, across every searched
|
||||
* herdr daemon. See {@link Lookup} for how to read a {@code null} terminal.
|
||||
*/
|
||||
public Lookup terminalForPid(long pid) {
|
||||
if (pid <= 0) {
|
||||
return null;
|
||||
return Lookup.NOT_FOUND;
|
||||
}
|
||||
Set<Long> ancestry = ancestorsOf(pid);
|
||||
for (HerdrClient herdr : herdrs) {
|
||||
String terminal = terminalForPid(herdr, ancestry);
|
||||
if (terminal != null) {
|
||||
return terminal;
|
||||
boolean complete = true;
|
||||
for (int i = 0; i < herdrs.size(); i++) {
|
||||
Lookup outcome = scan(herdrs.get(i), i, herdrs.size(), ancestry);
|
||||
if (outcome.terminal() != null) {
|
||||
return outcome; // a definite match — no need to finish checking other clients
|
||||
}
|
||||
complete = complete && outcome.complete();
|
||||
}
|
||||
return null;
|
||||
return new Lookup(null, complete);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -117,31 +145,51 @@ public final class PaneLocator {
|
||||
return ancestry;
|
||||
}
|
||||
|
||||
private static String terminalForPid(HerdrClient herdr, Set<Long> ancestry) {
|
||||
/** Whether a pane owns one of the scanned pid's ancestors, or the check of it failed outright. */
|
||||
private enum Ownership { OWNS, DOES_NOT_OWN, UNKNOWN }
|
||||
|
||||
private static Lookup scan(HerdrClient herdr, int clientIndex, int clientCount, Set<Long> ancestry) {
|
||||
boolean complete = true;
|
||||
for (JsonNode pane : herdr.call("pane.list", Map.of()).path("panes")) {
|
||||
String paneId = pane.path("pane_id").asText(null);
|
||||
if (paneId != null && paneOwnsAnyOf(herdr, paneId, ancestry)) {
|
||||
return pane.path("terminal_id").asText(null);
|
||||
if (paneId == null) {
|
||||
continue;
|
||||
}
|
||||
Ownership owns = paneOwnsAnyOf(herdr, clientIndex, clientCount, paneId, ancestry);
|
||||
if (owns == Ownership.OWNS) {
|
||||
return new Lookup(pane.path("terminal_id").asText(null), true);
|
||||
}
|
||||
if (owns == Ownership.UNKNOWN) {
|
||||
complete = false;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
return new Lookup(null, complete);
|
||||
}
|
||||
|
||||
private static boolean paneOwnsAnyOf(HerdrClient herdr, String paneId, Set<Long> ancestry) {
|
||||
private static Ownership paneOwnsAnyOf(HerdrClient herdr, int clientIndex, int clientCount,
|
||||
String paneId, Set<Long> ancestry) {
|
||||
JsonNode info;
|
||||
try {
|
||||
info = herdr.call("pane.process_info", Map.of("pane_id", paneId)).path("process_info");
|
||||
} catch (HerdrException e) {
|
||||
return false; // pane vanished mid-scan — just skip it
|
||||
// fleetd #505: this used to be read as a clean "does not own it" (the pane vanished
|
||||
// mid-scan, just skip it) — one boolean carrying two different facts. It is UNKNOWN
|
||||
// now: if THIS pane is the one that owns the pid, the caller must not be told "no pane
|
||||
// owns it", because that reads as a real primary and is promoted under loopback-trust.
|
||||
log.warn("pane.process_info failed for pane {} on herdr client {} of {} during a "
|
||||
+ "pid-owner scan — treating it as \"could not tell\", not a clean "
|
||||
+ "negative (fleetd #505): {}",
|
||||
paneId, clientIndex + 1, clientCount, e.getMessage());
|
||||
return Ownership.UNKNOWN;
|
||||
}
|
||||
if (ancestry.contains(info.path("shell_pid").asLong(-1))) {
|
||||
return true;
|
||||
return Ownership.OWNS;
|
||||
}
|
||||
for (JsonNode p : info.path("foreground_processes")) {
|
||||
if (ancestry.contains(p.path("pid").asLong(-1))) {
|
||||
return true;
|
||||
return Ownership.OWNS;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
return Ownership.DOES_NOT_OWN;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,6 +15,7 @@ import java.util.Set;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Predicate;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@@ -94,6 +95,14 @@ public final class Injector {
|
||||
private final TurnListener turnListener;
|
||||
private final Predicate<String> ready; // CB-113: a target is deliverable only when available
|
||||
private final Consumer<String> forget; // CB-114: clear a gone worker's readiness/presence
|
||||
/**
|
||||
* Wall-clock source for the readiness-grace elapsed time logged in {@link #onStatus} (fleetd
|
||||
* #501). Production constructors default this to {@code System::currentTimeMillis}; the
|
||||
* package-private constructors below take it explicitly so a test can supply a stub whose
|
||||
* advance does not track {@link #POLL_INTERVAL_MILLIS} — copying the shape {@code LeadRollover}
|
||||
* already uses for the same purpose.
|
||||
*/
|
||||
private final LongSupplier nowMillis;
|
||||
private final ConcurrentHashMap<String, Target> targets = new ConcurrentHashMap<>();
|
||||
|
||||
/** Delivery only; completion signalling is a no-op and every target is treated as available. */
|
||||
@@ -125,20 +134,34 @@ public final class Injector {
|
||||
*/
|
||||
public Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
|
||||
Consumer<String> forget) {
|
||||
this(agents, turnListener, ready, forget, System::currentTimeMillis);
|
||||
}
|
||||
|
||||
/** Full constructor — for tests: an injectable wall-clock supplier (fleetd #501). */
|
||||
Injector(AgentControl agents, TurnListener turnListener, Predicate<String> ready,
|
||||
Consumer<String> forget, LongSupplier nowMillis) {
|
||||
this.agents = agents;
|
||||
this.router = null;
|
||||
this.turnListener = turnListener;
|
||||
this.ready = ready;
|
||||
this.forget = forget;
|
||||
this.nowMillis = nowMillis;
|
||||
}
|
||||
|
||||
public Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
|
||||
Consumer<String> forget) {
|
||||
this(router, turnListener, ready, forget, System::currentTimeMillis);
|
||||
}
|
||||
|
||||
/** Full constructor — for tests: an injectable wall-clock supplier (fleetd #501). */
|
||||
Injector(HerdrRouter router, TurnListener turnListener, Predicate<String> ready,
|
||||
Consumer<String> forget, LongSupplier nowMillis) {
|
||||
this.agents = null;
|
||||
this.router = router;
|
||||
this.turnListener = turnListener;
|
||||
this.ready = ready;
|
||||
this.forget = forget;
|
||||
this.nowMillis = nowMillis;
|
||||
}
|
||||
|
||||
private AgentControl agentsFor(String target) {
|
||||
@@ -209,6 +232,7 @@ public final class Injector {
|
||||
int unknownSinceTurn; // consecutive `unknown` samples while a delegation is outstanding (CB-109)
|
||||
int unknownSincePostTurn; // the same, for the post-turn housekeeping phase (fleetd #306)
|
||||
int notReadySincePoll; // consecutive injectable samples a queued message waited on the readiness gate (CB-114)
|
||||
long notReadySinceMillis; // wall-clock time of the FIRST non-ready sample in the current notReadySincePoll streak (fleetd #501); reset alongside it
|
||||
boolean postTurnPending; // completion observed; adapter housekeeping has not started yet
|
||||
boolean awaitingPostTurnPickup;
|
||||
boolean postTurnObserved;
|
||||
@@ -282,7 +306,7 @@ public final class Injector {
|
||||
if (t == null) return;
|
||||
|
||||
Pending sent = null;
|
||||
RuntimeException sendError = null;
|
||||
Throwable sendError = null;
|
||||
boolean turnCompleted = false;
|
||||
boolean turnFailed = false;
|
||||
boolean resubmit = false;
|
||||
@@ -302,6 +326,7 @@ public final class Injector {
|
||||
t.unknownSinceTurn = 0;
|
||||
t.unknownSincePostTurn = 0;
|
||||
t.notReadySincePoll = 0;
|
||||
t.notReadySinceMillis = 0;
|
||||
if (t.awaitingCompletion) t.turnObserved = true;
|
||||
} else if (status.injectable()) { // IDLE or BLOCKED
|
||||
t.unknownSinceTurn = 0;
|
||||
@@ -353,6 +378,7 @@ public final class Injector {
|
||||
Pending p = t.queue.peek();
|
||||
if (p != null && ready.test(target)) {
|
||||
t.notReadySincePoll = 0;
|
||||
t.notReadySinceMillis = 0;
|
||||
try {
|
||||
agentsFor(target).send(target, p.text());
|
||||
t.queue.poll();
|
||||
@@ -362,31 +388,65 @@ public final class Injector {
|
||||
t.turnObserved = false;
|
||||
t.injectableSincePickup = 0;
|
||||
sent = p;
|
||||
} catch (RuntimeException e) {
|
||||
} catch (Throwable e) {
|
||||
// Delivery failed at herdr; drop the poisoned message and surface it
|
||||
// rather than blocking the queue behind it.
|
||||
// rather than blocking the queue behind it. Catches Throwable, not just
|
||||
// RuntimeException: fleetd #546 — an Error escaping this send (e.g. a
|
||||
// NoClassDefFoundError, see #413) would otherwise leave the entry QUEUED
|
||||
// at the head of t.queue. Line :378 peeks rather than polls, so the next
|
||||
// onStatus round would re-enter this try and send the same text again,
|
||||
// typing the same brief into the member's pane a second time.
|
||||
t.queue.poll();
|
||||
p.state = Pending.State.NOT_DELIVERED;
|
||||
sent = p;
|
||||
sendError = e;
|
||||
}
|
||||
} else if (p != null && ++t.notReadySincePoll >= READINESS_GRACE_POLLS) {
|
||||
// The worker has been idle-but-not-ready for the whole grace: its Claude
|
||||
// never connected the bridge MCP (crashed during boot, or wedged on a
|
||||
// startup prompt). The readiness gate would hold this message forever, so
|
||||
// fail every queued message and release the target (CB-114) instead of
|
||||
// polling it indefinitely with the caller's future never completing.
|
||||
notReady = new ArrayList<>(t.queue);
|
||||
for (Pending pending : notReady) {
|
||||
pending.state = Pending.State.NOT_DELIVERED;
|
||||
} else if (p != null) {
|
||||
// fleetd #501: stamp the wall-clock time of the FIRST non-ready sample in
|
||||
// this streak, so the expiry log below can print how long the target
|
||||
// actually sat non-ready — not just how many polls that took.
|
||||
if (t.notReadySincePoll == 0) {
|
||||
t.notReadySinceMillis = nowMillis.getAsLong();
|
||||
}
|
||||
if (++t.notReadySincePoll >= READINESS_GRACE_POLLS) {
|
||||
// The worker has been idle-but-not-ready for the whole grace: its Claude
|
||||
// never connected the bridge MCP (crashed during boot, or wedged on a
|
||||
// startup prompt). The readiness gate would hold this message forever, so
|
||||
// fail every queued message and release the target (CB-114) instead of
|
||||
// polling it indefinitely with the caller's future never completing.
|
||||
notReady = new ArrayList<>(t.queue);
|
||||
for (Pending pending : notReady) {
|
||||
pending.state = Pending.State.NOT_DELIVERED;
|
||||
}
|
||||
// fleetd #501: t.notReadySincePoll — the loop's own counter, already in
|
||||
// scope — is printed here instead of the READINESS_GRACE_POLLS constant.
|
||||
// On this branch the counter has JUST reached the threshold, so the two
|
||||
// agree by construction and no test can tell them apart. Printed anyway:
|
||||
// it gives this line one source of truth instead of two, so a later
|
||||
// change to the loop above cannot leave this message reporting a number
|
||||
// the loop no longer produces.
|
||||
//
|
||||
// elapsedMillis is a different case: it is NOT equal-by-construction to
|
||||
// the truth. notReadySincePoll only increments on a sample that reaches
|
||||
// this branch (p != null, not ready) — a poll that misses that condition
|
||||
// advances real time without advancing the counter — and this loop's real
|
||||
// period is not guaranteed to equal POLL_INTERVAL_MILLIS (load, or a host
|
||||
// sleep, can widen the real gap far past it). READINESS_GRACE_POLLS *
|
||||
// POLL_INTERVAL_MILLIS / 1000 is arithmetic on two constants, not a
|
||||
// measurement, so it stays here only as the labelled CONFIGURED budget,
|
||||
// never presented as elapsed time.
|
||||
long elapsedMillis = nowMillis.getAsLong() - t.notReadySinceMillis;
|
||||
log.warn("readiness grace for {} expired after {} polls (configured={} "
|
||||
+ "polls/{}s elapsed={}ms): target never became "
|
||||
+ "deliverable, so failing {} queued message(s) that "
|
||||
+ "never reached its pane",
|
||||
target, t.notReadySincePoll, READINESS_GRACE_POLLS,
|
||||
READINESS_GRACE_POLLS * POLL_INTERVAL_MILLIS / 1000, elapsedMillis,
|
||||
notReady.size());
|
||||
t.queue.clear();
|
||||
t.notReadySincePoll = 0;
|
||||
t.notReadySinceMillis = 0;
|
||||
}
|
||||
log.warn("readiness grace for {} expired after {} polls ({}s): target never "
|
||||
+ "became deliverable, so failing {} queued message(s) that never "
|
||||
+ "reached its pane",
|
||||
target, READINESS_GRACE_POLLS,
|
||||
READINESS_GRACE_POLLS * POLL_INTERVAL_MILLIS / 1000, notReady.size());
|
||||
t.queue.clear();
|
||||
t.notReadySincePoll = 0;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -64,34 +64,39 @@ public final class StatusPoller {
|
||||
}
|
||||
|
||||
private void loop() {
|
||||
while (running) {
|
||||
Set<String> active = injector.activeTargets();
|
||||
for (String target : active) {
|
||||
if (!running) return;
|
||||
try {
|
||||
// herdr's agent_status can misreport a settled worker as `unknown`; refine it
|
||||
// against the pane content before it drives delivery/completion (CB-115).
|
||||
// CB-185: refine THROUGH the same control the raw status came from — a router
|
||||
// splits lead/member targets across two herdr daemons, and reading a lead's pane
|
||||
// through the (fixed) member refiner never finds it, wedging that lead at UNKNOWN.
|
||||
AgentControl control = router != null ? router.agentsFor(target) : agents;
|
||||
AgentStatus status = refiner.refine(target, control.status(target), control);
|
||||
injector.onStatus(target, status);
|
||||
} catch (HerdrException e) {
|
||||
// The worker's agent is gone — stop trying and unblock its waiters.
|
||||
if (e.code() != null && e.code().endsWith("_not_found")) {
|
||||
log.debug("target {} gone; dropping its queue", target);
|
||||
injector.drop(target, e);
|
||||
} else {
|
||||
log.debug("status poll for {} failed (will retry): {}", target, e.getMessage());
|
||||
try {
|
||||
while (running) {
|
||||
Set<String> active = injector.activeTargets();
|
||||
for (String target : active) {
|
||||
if (!running) return;
|
||||
try {
|
||||
// herdr's agent_status can misreport a settled worker as `unknown`; refine it
|
||||
// against the pane content before it drives delivery/completion (CB-115).
|
||||
// CB-185: refine THROUGH the same control the raw status came from — a router
|
||||
// splits lead/member targets across two herdr daemons, and reading a lead's pane
|
||||
// through the (fixed) member refiner never finds it, wedging that lead at UNKNOWN.
|
||||
AgentControl control = router != null ? router.agentsFor(target) : agents;
|
||||
AgentStatus status = refiner.refine(target, control.status(target), control);
|
||||
injector.onStatus(target, status);
|
||||
} catch (HerdrException e) {
|
||||
// The worker's agent is gone — stop trying and unblock its waiters.
|
||||
if (e.code() != null && e.code().endsWith("_not_found")) {
|
||||
log.debug("target {} gone; dropping its queue", target);
|
||||
injector.drop(target, e);
|
||||
} else {
|
||||
log.debug("status poll for {} failed (will retry): {}", target, e.getMessage());
|
||||
}
|
||||
} catch (Throwable e) {
|
||||
log.error("unexpected failure polling {}; skipping this round", target, e);
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
// Never let one target's unexpected error (e.g. an odd agent.get shape) kill
|
||||
// the single poller thread and stall injection for every worker.
|
||||
log.warn("unexpected error polling {}; skipping this round", target, e);
|
||||
}
|
||||
sleep();
|
||||
}
|
||||
sleep();
|
||||
} finally {
|
||||
if (running) {
|
||||
log.error("status poller loop exited unexpectedly; it can be restarted");
|
||||
}
|
||||
running = false;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -33,9 +33,10 @@ public final class ConnectionIdentity {
|
||||
|
||||
/**
|
||||
* The caller resolved from the connection: its worker {@code terminal} (or {@code null} for the
|
||||
* primary / an off-host client) and its {@code pid} (or {@code -1} if not resolvable).
|
||||
* primary / an off-host client), its {@code pid} (or {@code -1} if not resolvable), and whether
|
||||
* the pane scan behind {@code terminal} ran to completion ({@link #scanComplete}).
|
||||
*/
|
||||
public record Caller(String terminal, long pid) {
|
||||
public record Caller(String terminal, long pid, boolean scanComplete) {
|
||||
|
||||
/**
|
||||
* Whether the OS peer-PID lookup actually succeeded — {@code false} means {@code pid} is
|
||||
@@ -51,6 +52,10 @@ public final class ConnectionIdentity {
|
||||
* {@link ConnectionIdentity#isLoopback} is centralised rather than left for each caller to
|
||||
* reimplement: a raw {@code pid > 0} check duplicated at every call site is precisely the
|
||||
* "one rule, two copies" shape that let #305 drift.
|
||||
*
|
||||
* <p>This method is deliberately NOT widened for fleetd #505's failure (a herdr error
|
||||
* during the pane scan, not a failed lsof lookup) — it still tests only the sentinel it is
|
||||
* named for. #505 is a different axis, carried separately in {@link #scanComplete}.
|
||||
*/
|
||||
public boolean resolved() {
|
||||
return pid > 0;
|
||||
@@ -60,10 +65,11 @@ public final class ConnectionIdentity {
|
||||
/** Resolve the caller's terminal and PID from one peer-PID lookup. */
|
||||
public Caller resolve(String remoteAddr, int remotePort) {
|
||||
if (!isLoopback(remoteAddr)) {
|
||||
return new Caller(null, -1); // only same-host callers can be workers
|
||||
return new Caller(null, -1, true); // only same-host callers can be workers
|
||||
}
|
||||
long pid = pids.pidForLocalPort(remotePort);
|
||||
return new Caller(panes.terminalForPid(pid), pid);
|
||||
PaneLocator.Lookup lookup = panes.terminalForPid(pid);
|
||||
return new Caller(lookup.terminal(), pid, lookup.complete());
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -93,7 +93,18 @@ public final class FleetMcp {
|
||||
|
||||
private final HttpServletStreamableServerTransportProvider transport;
|
||||
private final McpSyncServer server;
|
||||
private final CallerResolver authz; // CB-501: null → authorization not enforced (legacy)
|
||||
/**
|
||||
* fleetd #518: whether {@link #denyFor} enforces the CB-505 policy table at all. Replaces the
|
||||
* old {@code CallerResolver authz} field, whose null-ness used to decide BOTH this AND which
|
||||
* principal-resolution code path {@link #contextExtractor} ran — reaching "authorization off"
|
||||
* by simply not passing a {@link CallerResolver} also meant the resolved {@link Principal}
|
||||
* came from a second, separately-maintained heuristic ({@code legacyPrincipal}, now deleted)
|
||||
* that nothing ever exercised. There is now exactly one resolution path ({@code callers},
|
||||
* required and non-null below) and a separate, explicitly-chosen {@link AuthorizationMode}
|
||||
* for this flag — so a caller can turn enforcement off without silently swapping in a second,
|
||||
* untested identity heuristic.
|
||||
*/
|
||||
private final boolean authorizationEnforced;
|
||||
private final Metrics metrics; // CB-502: null → auth failures not counted
|
||||
private final CapacitySource capacity;
|
||||
private final HealthCoverageSource healthCoverage;
|
||||
@@ -250,6 +261,17 @@ public final class FleetMcp {
|
||||
public static CoordinationSource none() { return new CoordinationSource(null, List.of()); }
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #518: whether {@link #denyFor} enforces the CB-505 policy table. A required
|
||||
* constructor parameter with no default, so "authorization is off" can only be reached by a
|
||||
* caller explicitly saying so — never by omitting a {@link CallerResolver} the way the old
|
||||
* {@code callers == null} idiom allowed. {@code callers} itself is required either way: even
|
||||
* under {@link #UNENFORCED}, the one real {@link CallerResolver} still resolves every caller's
|
||||
* {@link Principal} (so {@code markSpawnedMemberPresent}/{@code recordPrimarySingleton} see a
|
||||
* real identity), and {@link #denyFor} is the only thing that changes.
|
||||
*/
|
||||
public enum AuthorizationMode { ENFORCED, UNENFORCED }
|
||||
|
||||
/**
|
||||
* The only constructor (fleetd #480 Unit C correction round). Every field below used to have
|
||||
* its own defaulting overload — {@code leadChannel}/{@code outage}/{@code leadSeats}/
|
||||
@@ -268,10 +290,16 @@ public final class FleetMcp {
|
||||
* {@link OutageSource#none()}, {@link LeadSeatSource#none()}, {@code List.of()} are all still
|
||||
* perfectly fine values, just never an implicit default reached by omission.
|
||||
*
|
||||
* @param callers resolves each call's {@link Principal}; {@code null} disables
|
||||
* authorization. This surface needs its own enforcement: {@code /mcp} is a
|
||||
* raw servlet on Jetty's context handler and never passes through
|
||||
* Javalin's {@code before} filter, so the REST guard does not cover it.
|
||||
* @param callers resolves each call's {@link Principal}. Required, never {@code null} —
|
||||
* fleetd #518: use {@link AuthorizationMode#UNENFORCED} to disable
|
||||
* enforcement, not a missing resolver. This surface needs its own
|
||||
* enforcement: {@code /mcp} is a raw servlet on Jetty's context handler and
|
||||
* never passes through Javalin's {@code before} filter, so the REST guard
|
||||
* does not cover it.
|
||||
* @param authorizationMode fleetd #518: whether {@link #denyFor} enforces the CB-505 policy
|
||||
* table ({@link AuthorizationMode#ENFORCED}) or leaves the gate open
|
||||
* ({@link AuthorizationMode#UNENFORCED}, for the pre-CB-513 test suite that
|
||||
* does not exercise authorization). Required, with no default.
|
||||
* @param metrics registry for auth-failure counting; may be {@code null}
|
||||
* @param quarantine CB-578 stage B facts for {@code fleet_profiles}; pass
|
||||
* {@link QuarantineSource#none()} for a caller that does not want the
|
||||
@@ -301,9 +329,13 @@ public final class FleetMcp {
|
||||
*/
|
||||
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
|
||||
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
|
||||
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
CallerResolver callers, AuthorizationMode authorizationMode, Metrics metrics,
|
||||
CapacitySource capacity, HealthCoverageSource healthCoverage,
|
||||
QuarantineSource quarantine, LeadChannel leadChannel, OutageSource outage,
|
||||
LeadSeatSource leadSeats, List<String> peers, LeadRollover leadRollover) {
|
||||
Objects.requireNonNull(callers, "callers");
|
||||
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
|
||||
== AuthorizationMode.ENFORCED;
|
||||
this.leadChannel = leadChannel;
|
||||
this.peers = peers == null ? List.of() : List.copyOf(peers);
|
||||
this.capacity = capacity;
|
||||
@@ -322,11 +354,11 @@ public final class FleetMcp {
|
||||
// (CB-113) — its MCP initialize is the reliable "the agent is up" signal.
|
||||
.contextExtractor(req -> {
|
||||
// One resolution per call, shared with the REST surface via CallerResolver so
|
||||
// the two paths cannot drift on who a caller is.
|
||||
Principal p = callers != null
|
||||
? callers.resolve(req.getRemoteAddr(), req.getRemotePort(),
|
||||
req.getHeader("Authorization"))
|
||||
: legacyPrincipal(identity, req.getRemoteAddr(), req.getRemotePort());
|
||||
// the two paths cannot drift on who a caller is. fleetd #518: callers is
|
||||
// required (never null) so there is no second, untested resolution path to
|
||||
// fall back to here — AuthorizationMode governs enforcement, not identity.
|
||||
Principal p = callers.resolve(req.getRemoteAddr(), req.getRemotePort(),
|
||||
req.getHeader("Authorization"));
|
||||
// CB-532: guard on the ROLE, not on the terminal being null. This excludes a
|
||||
// lead, which carries its pane too, while including every spawned member role.
|
||||
// Enrolling a lead would count it as an available member in the roster.
|
||||
@@ -443,7 +475,7 @@ public final class FleetMcp {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_list", Map.of()), null);
|
||||
if (denied != null) return denied;
|
||||
return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, outage,
|
||||
leadSeats, callers == null ? Map.of() : callers.leads(),
|
||||
leadSeats, callers.leads(),
|
||||
callerTerminal(exchange),
|
||||
new CoordinationSource(leadChannel, peers),
|
||||
coordinatorVisibleTo(principal(exchange)));
|
||||
@@ -522,21 +554,9 @@ public final class FleetMcp {
|
||||
.toolCall(fleetWhoami, whoamiHandler)
|
||||
.toolCall(fleetHandover, handoverHandler)
|
||||
.build();
|
||||
this.authz = callers;
|
||||
this.metrics = metrics;
|
||||
}
|
||||
|
||||
/**
|
||||
* Pre-CB-501 identity: worker if the connection maps to a pane, otherwise the primary. Used
|
||||
* only by the legacy constructor, where authorization is not enforced anyway.
|
||||
*/
|
||||
private static Principal legacyPrincipal(ConnectionIdentity identity, String addr, int port) {
|
||||
ConnectionIdentity.Caller c = identity.resolve(addr, port);
|
||||
return c.terminal() != null
|
||||
? Principal.worker(c.terminal(), c.pid())
|
||||
: Principal.primary(c.pid());
|
||||
}
|
||||
|
||||
/** The caller reconstructed from the transport context. */
|
||||
private static Principal principal(McpSyncServerExchange exchange) {
|
||||
return principalFrom(exchange.transportContext().get(CALLER_ROLE),
|
||||
@@ -591,8 +611,8 @@ public final class FleetMcp {
|
||||
McpSchema.CallToolResult denyFor(Principal caller, Authz.Action action, String target) {
|
||||
// The enforcement switch lives HERE rather than in the exchange-facing wrapper: any future
|
||||
// tool that calls this directly must not be able to skip the gate by accident.
|
||||
if (authz == null) {
|
||||
return null; // legacy constructor: authorization not enforced
|
||||
if (!authorizationEnforced) {
|
||||
return null; // AuthorizationMode.UNENFORCED: authorization not enforced (fleetd #518)
|
||||
}
|
||||
if (Authz.permits(caller, action, target)) {
|
||||
if (action != Authz.Action.READ) {
|
||||
|
||||
@@ -246,7 +246,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
* neither. So before this method existed with a requeue step, it dropped {@link #held}'s entries
|
||||
* for {@code target} while the broker still considered them outstanding: never acked, never
|
||||
* nacked, never requeued, and no longer reachable by {@link #peek} — permanently invisible. This
|
||||
* is unlike {@link #handleRecovery} and {@link #close()}, whose bare {@code held.clear()} is
|
||||
* is unlike {@link RecoveryListener#handleRecovery(Recoverable)} and {@link #close()}, whose bare {@code held.clear()} is
|
||||
* correct because each has already made the broker requeue (a real connection drop, or
|
||||
* {@code channel.close()} respectively) before clearing local state.
|
||||
*
|
||||
|
||||
@@ -713,7 +713,7 @@ public final class MessageService {
|
||||
* failure.
|
||||
*
|
||||
* <p><strong>Without {@code sweepAsking} on the release path, a target torn down while
|
||||
* genuinely {@code ASKING} was unrecoverable.</strong> {@link #resolveQuestion} had already
|
||||
* genuinely {@code ASKING} was unrecoverable.</strong> {@link Rendezvous#resolveQuestion(String, String, String)} had already
|
||||
* closed the forward waiter the instant the question surfaced (so the {@code waiter} branch
|
||||
* below finds nothing to fail), the {@code question == null} guard excluded the task from
|
||||
* {@code matching} (so the loop below skipped it too), and the worker's own {@code fleet_ask}
|
||||
|
||||
@@ -44,7 +44,7 @@ import java.util.stream.Collectors;
|
||||
* still holds an unacked message ({@link #pendingReplies}), tickets not yet collected
|
||||
* ({@link #pendingTickets}), and open questions not yet answered or lapsed
|
||||
* ({@link #pendingQuestions}) — and sends at most one combined nudge per tick
|
||||
* ({@link #injectNudge(String, int, int, int)}). Work that arrives while the lead is busy is
|
||||
* ({@link #injectNudge(String, int, int, int, int, int)}). Work that arrives while the lead is busy is
|
||||
* never lost: it is re-read fresh on every tick until the lead is injectable or its own reminder
|
||||
* cap ({@link #maxReminders}) is reached — each source spends from its own budget, so one source
|
||||
* exhausting its cap does not stop nudges about the others (post-CB-590 regression fix; see
|
||||
|
||||
@@ -302,9 +302,10 @@ public final class SessionManager implements TurnListener {
|
||||
* with no copy and no error. Do NOT fuse these back together; the cost of an orphaned worktree
|
||||
* is a logged path an operator can reclaim, the cost of a deleted one is unrecoverable work.
|
||||
*/
|
||||
private void release(String paneId, ReleaseCause cause) {
|
||||
private MemberSession release(String paneId, ReleaseCause cause) {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
|
||||
return removed;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1064,16 +1065,62 @@ public final class SessionManager implements TurnListener {
|
||||
* drain (see above), and a straggler must not buy the drain more time than the flag it lost the
|
||||
* race against would have. In the ordinary case the sweep finds nothing and costs one empty
|
||||
* {@link #roster()} call.
|
||||
*
|
||||
* <p>fleetd #512: a drain that releases every session cleanly used to log nothing at all — the
|
||||
* only log calls in this method and {@link #drainSnapshot} sit on abnormal paths, so "nothing
|
||||
* logged" was indistinguishable from "died on the first session". The {@code log.info} at the
|
||||
* end below is a positive assertion that the drain actually finished, on the normal path,
|
||||
* every time — including the all-zero case, which is a common and legitimate outcome (no
|
||||
* members were live) and must still produce the line. Both {@link #drainSnapshot} passes (the
|
||||
* main snapshot and the straggler sweep) are folded into the one line: a caller reading two
|
||||
* lines could not tell a two-pass drain from two separate drains.
|
||||
*
|
||||
* <p><strong>Non-goal: this line must never move into a {@code finally} block, and this method
|
||||
* must never grow one around it.</strong> "Every time" above means every time the drain
|
||||
* <em>finishes</em>, not every time this method exits. The absence of the line is the signal
|
||||
* that the drain died, so a {@code finally} would destroy the signal and print confident
|
||||
* partial counts in the same edit — the line would appear after a drain that threw, carrying
|
||||
* whatever {@code tally} it had reached. Both halves of the value are lost at once. The line
|
||||
* has to be the last statement of the successful path and reachable only from it.
|
||||
*
|
||||
* <p>This is written down because it is the obvious review comment ("shouldn't we always log
|
||||
* the drain result?"), it sounds like thoroughness, and the paragraph above reads as an
|
||||
* invitation to it. Raised by the fleet01 lead on 2026-09-12, from their 2026-09-10 incident:
|
||||
* we only know that drain died on that host because it <em>threw</em>, and a
|
||||
* {@code NoClassDefFoundError} reached the JVM's uncaught handler. A drain that hung on one
|
||||
* session, or returned early on a condition rather than an exception, would leave no stack
|
||||
* trace, no {@code ERROR} token and no priority — only a missing line. That makes the loud
|
||||
* variant the one we have seen and the quiet variants the ones this line exists to catch.
|
||||
*
|
||||
* <p>Related: {@code released} and {@code abandoned} are counted incrementally inside {@link
|
||||
* #drainSnapshot}'s loop and folded with {@link DrainTally#plus}, rather than derived from a
|
||||
* collection read at the end, for the same reason. If a partial report is ever wanted it must
|
||||
* be a different line with a different verb. One line must not serve both, or a reader cannot
|
||||
* tell a finished drain from an interrupted one by its wording.
|
||||
*/
|
||||
void drainAll(long timeoutNanos) {
|
||||
long deadline = System.nanoTime() + timeoutNanos;
|
||||
draining.set(true);
|
||||
drainSnapshot(roster(), deadline);
|
||||
DrainTally tally = drainSnapshot(roster(), deadline);
|
||||
List<MemberSession> stragglers = roster();
|
||||
if (!stragglers.isEmpty()) {
|
||||
log.warn("drain sweep found {} session(s) registered after the drain snapshot was "
|
||||
+ "taken (raced past the shutdown guard); draining them too", stragglers.size());
|
||||
drainSnapshot(stragglers, deadline);
|
||||
tally = tally.plus(drainSnapshot(stragglers, deadline));
|
||||
}
|
||||
log.info("drain complete: released={} abandoned={} (still BUSY at the shutdown deadline)",
|
||||
tally.released(), tally.abandoned());
|
||||
}
|
||||
|
||||
/**
|
||||
* Running count for one {@link #drainAll} invocation, folded across both {@link #drainSnapshot}
|
||||
* passes (fleetd #512). {@code abandoned} counts sessions that were still {@code BUSY} at the
|
||||
* moment they were released — i.e. the whole-drain deadline passed before they left {@code BUSY}
|
||||
* on their own (see {@link #drainSnapshot}) — a subset of {@code released}, not additional to it.
|
||||
*/
|
||||
private record DrainTally(int released, int abandoned) {
|
||||
private DrainTally plus(DrainTally other) {
|
||||
return new DrainTally(released + other.released, abandoned + other.abandoned);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1081,8 +1128,12 @@ public final class SessionManager implements TurnListener {
|
||||
* Drain exactly the sessions in {@code snapshot}, waiting out a {@code BUSY} one against the
|
||||
* shared whole-drain {@code deadline} before releasing it. Shared by {@link #drainAll}'s main
|
||||
* pass and its post-loop straggler sweep (fleetd #308) so both honor the same one budget.
|
||||
* Returns how many sessions this pass released, and how many of those were still {@code BUSY}
|
||||
* (abandoned mid-turn) at the moment of release.
|
||||
*/
|
||||
private void drainSnapshot(List<MemberSession> snapshot, long deadline) {
|
||||
private DrainTally drainSnapshot(List<MemberSession> snapshot, long deadline) {
|
||||
int released = 0;
|
||||
int abandoned = 0;
|
||||
for (MemberSession s : snapshot) {
|
||||
try {
|
||||
if (s.state() == MemberSession.State.BUSY) {
|
||||
@@ -1100,11 +1151,16 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
}
|
||||
}
|
||||
release(s.paneId(), ReleaseCause.SHUTDOWN);
|
||||
MemberSession removed = release(s.paneId(), ReleaseCause.SHUTDOWN);
|
||||
released++;
|
||||
if (removed != null && removed.state() == MemberSession.State.BUSY) {
|
||||
abandoned++;
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("drain failed for pane={}; continuing with remaining sessions", s.paneId(), e);
|
||||
}
|
||||
}
|
||||
return new DrainTally(released, abandoned);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -60,14 +60,21 @@ public final class SessionReaper {
|
||||
}
|
||||
|
||||
private void loop() {
|
||||
while (running) {
|
||||
try {
|
||||
sessions.reapIdle(idleTtlNanos);
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("session reaper iteration failed; continuing", e);
|
||||
try {
|
||||
while (running) {
|
||||
try {
|
||||
sessions.reapIdle(idleTtlNanos);
|
||||
} catch (Throwable e) {
|
||||
log.error("session reaper iteration failed; continuing", e);
|
||||
}
|
||||
maybeSweepWipRefs();
|
||||
sleep();
|
||||
}
|
||||
maybeSweepWipRefs();
|
||||
sleep();
|
||||
} finally {
|
||||
if (running) {
|
||||
log.error("session reaper loop exited unexpectedly; it can be restarted");
|
||||
}
|
||||
running = false;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -90,8 +97,8 @@ public final class SessionReaper {
|
||||
log.info("refs/wip retention sweep deleted {} snapshot ref(s) older than 24h whose "
|
||||
+ "content was already reachable from main", deleted);
|
||||
}
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("refs/wip retention sweep failed; continuing", e);
|
||||
} catch (Throwable e) {
|
||||
log.error("refs/wip retention sweep failed; continuing", e);
|
||||
}
|
||||
// Set even when the sweep threw, so a broken repo is retried on the slow cadence rather
|
||||
// than hammering git on every 5-second iteration.
|
||||
|
||||
@@ -0,0 +1,192 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
import static org.junit.jupiter.api.Assertions.fail;
|
||||
|
||||
/**
|
||||
* fleetd #498: {@code Fleetd.awaitHerdr} used to return a bare {@code boolean}, collapsing "the
|
||||
* configured wait budget genuinely ran out" and "the waiting thread was interrupted, possibly
|
||||
* milliseconds in" onto the same {@code false} — and the caller's log line printed only the
|
||||
* configured budget, never how long the wait actually ran. This class covers both halves of the
|
||||
* fix:
|
||||
* <ul>
|
||||
* <li>the seam — {@link Fleetd#awaitHerdr} itself, driven with an injected clock and a stub
|
||||
* {@link HerdrClient}, one test per {@link Fleetd.HerdrWaitResult};</li>
|
||||
* <li>the call site — {@link Fleetd#logHerdrWaitOutcomeAndShouldReap}, the exact decision {@code
|
||||
* main} calls (extracted here because {@code main} itself boots the whole daemon and cannot
|
||||
* be driven from a unit test), pinning the three distinct log messages it emits.</li>
|
||||
* </ul>
|
||||
* Every expected message below is a plain literal, not built from {@code HERDR_WAIT_SECONDS} or
|
||||
* any other production constant — a test that derives its expectation the way the code does
|
||||
* cannot see a change to either (fleetd #496's identical trap).
|
||||
*/
|
||||
class FleetdAwaitHerdrTest {
|
||||
|
||||
// ---- the seam: Fleetd.awaitHerdr ----------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void answeredReturnsImmediatelyWithZeroElapsedAndNeverPolls() {
|
||||
HerdrStub herdr = new HerdrStub(0); // succeeds on the very first call
|
||||
LongSupplier clock = fixedClock(1_000L);
|
||||
AtomicBoolean polled = new AtomicBoolean(false);
|
||||
Runnable poller = () -> polled.set(true);
|
||||
|
||||
Fleetd.HerdrAwaitOutcome outcome = Fleetd.awaitHerdr(herdr, clock, poller);
|
||||
|
||||
assertEquals(Fleetd.HerdrWaitResult.ANSWERED, outcome.result());
|
||||
assertEquals(0L, outcome.elapsedNanos(), "a fixed clock must measure zero elapsed time");
|
||||
assertFalse(polled.get(), "herdr answering on the first try must never poll");
|
||||
}
|
||||
|
||||
@Test
|
||||
void deadlinePassedIsMeasuredNotAssumed() {
|
||||
HerdrStub herdr = new HerdrStub(-1); // never succeeds
|
||||
// call order inside awaitHerdr: start, then per failed attempt: deadline-check, elapsed-calc
|
||||
ScriptedClock clock = new ScriptedClock(0L, 30_500_000_000L, 30_500_000_000L);
|
||||
Runnable poller = () -> fail("the deadline was already exceeded on the first attempt — must not poll");
|
||||
|
||||
Fleetd.HerdrAwaitOutcome outcome = Fleetd.awaitHerdr(herdr, clock, poller);
|
||||
|
||||
assertEquals(Fleetd.HerdrWaitResult.DEADLINE_PASSED, outcome.result());
|
||||
assertEquals(30_500_000_000L, outcome.elapsedNanos(),
|
||||
"elapsed must be the MEASURED clock delta, not the configured budget");
|
||||
}
|
||||
|
||||
@Test
|
||||
void interruptedIsDistinctFromDeadlinePassedAndPreservesTheInterruptFlag() {
|
||||
HerdrStub herdr = new HerdrStub(-1); // never succeeds
|
||||
// start=0, deadline-check returns 500ms (well under the 30s budget) -> not deadline-passed,
|
||||
// then the poller interrupts, and the elapsed-calc call returns 750ms.
|
||||
ScriptedClock clock = new ScriptedClock(0L, 500_000_000L, 750_000_000L);
|
||||
Runnable poller = () -> Thread.currentThread().interrupt();
|
||||
|
||||
try {
|
||||
Fleetd.HerdrAwaitOutcome outcome = Fleetd.awaitHerdr(herdr, clock, poller);
|
||||
|
||||
assertEquals(Fleetd.HerdrWaitResult.INTERRUPTED, outcome.result());
|
||||
assertEquals(750_000_000L, outcome.elapsedNanos(),
|
||||
"elapsed must be measured even when the wait ends via interruption, not the deadline");
|
||||
assertTrue(Thread.currentThread().isInterrupted(),
|
||||
"the interrupt flag the old code re-set must still be set on return");
|
||||
} finally {
|
||||
Thread.interrupted(); // clear it so it cannot leak into another test on this thread
|
||||
}
|
||||
}
|
||||
|
||||
// ---- the call site: Fleetd.logHerdrWaitOutcomeAndShouldReap -------------------------------
|
||||
|
||||
@Test
|
||||
void answeredLogsNothingAndSaysReap() {
|
||||
try (CapturedLog log = attach()) {
|
||||
boolean shouldReap = Fleetd.logHerdrWaitOutcomeAndShouldReap(
|
||||
new Fleetd.HerdrAwaitOutcome(Fleetd.HerdrWaitResult.ANSWERED, 0L));
|
||||
|
||||
assertTrue(shouldReap, "only ANSWERED should tell main to reap orphan workers");
|
||||
assertEquals(0, log.events().size(), "the answered path logs nothing itself");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void deadlinePassedLogsConfiguredAndMeasuredElapsedTogether() {
|
||||
try (CapturedLog log = attach()) {
|
||||
boolean shouldReap = Fleetd.logHerdrWaitOutcomeAndShouldReap(
|
||||
new Fleetd.HerdrAwaitOutcome(Fleetd.HerdrWaitResult.DEADLINE_PASSED, 30_500_000_000L));
|
||||
|
||||
assertFalse(shouldReap, "a deadline-passed wait must not tell main to reap");
|
||||
assertEquals(1, log.events().size());
|
||||
ILoggingEvent event = log.events().getFirst();
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
assertEquals("herdr did not answer within the configured wait (configured=30s "
|
||||
+ "elapsed=30500ms) — starting anyway; /healthz will report degraded until it "
|
||||
+ "comes up. Orphaned worker panes (if any) were NOT reaped.",
|
||||
event.getFormattedMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void interruptedLogsItsOwnMessageAndNeverClaimsTheBudgetElapsed() {
|
||||
try (CapturedLog log = attach()) {
|
||||
// 3ms: the ticket's own example of "a few milliseconds in", not the 30s budget.
|
||||
boolean shouldReap = Fleetd.logHerdrWaitOutcomeAndShouldReap(
|
||||
new Fleetd.HerdrAwaitOutcome(Fleetd.HerdrWaitResult.INTERRUPTED, 3_000_000L));
|
||||
|
||||
assertFalse(shouldReap, "an interrupted wait must not tell main to reap");
|
||||
assertEquals(1, log.events().size());
|
||||
ILoggingEvent event = log.events().getFirst();
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
assertEquals("herdr wait was interrupted before the configured wait ran out "
|
||||
+ "(configured=30s elapsed=3ms) — starting anyway; /healthz will report "
|
||||
+ "degraded until it comes up. Orphaned worker panes (if any) were NOT reaped.",
|
||||
message);
|
||||
assertFalse(message.contains("did not answer"),
|
||||
"an interrupted wait must not be reported as if herdr failed to answer within the budget");
|
||||
}
|
||||
}
|
||||
|
||||
// ---- fixtures --------------------------------------------------------------------------
|
||||
|
||||
/** Always returns the same value, i.e. a clock that measures zero elapsed time. */
|
||||
private static LongSupplier fixedClock(long value) {
|
||||
return () -> value;
|
||||
}
|
||||
|
||||
/** Returns each value in order, then repeats the last one for any call beyond the list. */
|
||||
private static final class ScriptedClock implements LongSupplier {
|
||||
private final long[] values;
|
||||
private int index;
|
||||
|
||||
ScriptedClock(long... values) {
|
||||
this.values = values;
|
||||
}
|
||||
|
||||
@Override
|
||||
public long getAsLong() {
|
||||
long v = values[Math.min(index, values.length - 1)];
|
||||
if (index < values.length - 1) {
|
||||
index++;
|
||||
}
|
||||
return v;
|
||||
}
|
||||
}
|
||||
|
||||
/** Fails {@code failuresBeforeSuccess} times, then succeeds forever; {@code -1} never succeeds. */
|
||||
private static final class HerdrStub implements HerdrClient {
|
||||
private final int failuresBeforeSuccess;
|
||||
private int calls;
|
||||
|
||||
HerdrStub(int failuresBeforeSuccess) {
|
||||
this.failuresBeforeSuccess = failuresBeforeSuccess;
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) throws HerdrException {
|
||||
calls++;
|
||||
if (failuresBeforeSuccess < 0 || calls <= failuresBeforeSuccess) {
|
||||
throw new HerdrException("herdr not up yet");
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
private static CapturedLog attach() {
|
||||
return CapturedLog.at(Fleetd.class, Level.DEBUG);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,160 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #426: {@code FleetHealthMonitor.coverage} had zero references anywhere in the test
|
||||
* tree — not the method, not either output string, not the field it populates. {@code
|
||||
* FleetHealthMonitorCoverageTest} (package {@code dev.ltms.fleet.health}) pins the three-branch
|
||||
* method itself; that is the easy half.
|
||||
*
|
||||
* <p>The half that actually matters is this one: {@link Fleetd#healthCoverageSource} is the exact
|
||||
* call site {@code Fleetd.main} wires into {@code FleetMcp}'s constructor, and it is what feeds
|
||||
* {@code fleet_list}'s {@code healthCoverage} field (see {@code FleetMcp#listFleet}'s {@code
|
||||
* result.put("healthCoverage", healthCoverage.value().get())}). Measured precedent on fleetd #423
|
||||
* (for #415): swapping two arguments at a call site like this one — recreating #415's defect with
|
||||
* the keys exchanged — compiled with 0 errors and ran the ENTIRE suite (1506 tests) green. A
|
||||
* thoroughly-tested method proves nothing about whether the call site pairs its arguments correctly;
|
||||
* only a test that drives the call site itself can catch that.
|
||||
*
|
||||
* <p><strong>Why this is not driven through a real {@code Fleetd.main} the way #407 drives its five
|
||||
* reporters</strong> (keep the config invalid, assert on the log line emitted before {@code
|
||||
* cfg.validateAll()} throws): all five of #407's reporters run in {@code main} before {@code
|
||||
* validateAll()} (line ~171). {@link Fleetd#healthCoverageSource} is built during {@code FleetMcp}
|
||||
* construction, which happens only after {@code UnixSocketHerdrClient.connect} has already opened a
|
||||
* real herdr socket (line ~188) and after {@code sessions}/{@code workers} are constructed. Reaching
|
||||
* this call site by actually running {@code main} would require a real socket connect — banned by
|
||||
* this ticket's hard constraints — so #407's option 1 does not apply here. Instead {@link
|
||||
* Fleetd#healthCoverageSource} is extracted to a directly-callable package-private factory, the same
|
||||
* shape {@link Fleetd#capacitySource} and {@link Fleetd#quarantineSource} already use for the same
|
||||
* reason (see {@code FleetdCapacitySourceWiringTest}, the direct precedent this test follows).
|
||||
*
|
||||
* <p>Uses a real {@link FleetConfig#load} + {@link ConfigRef} (no socket, no port bind, no spawn,
|
||||
* nothing written outside {@code @TempDir}) so the fixture goes through the actual YAML parser and
|
||||
* {@code FleetConfig.Health}/{@code Notifications} records, not a hand-built stand-in that could
|
||||
* silently drift from what the parser actually produces.
|
||||
*/
|
||||
class FleetdHealthCoverageSourceWiringTest {
|
||||
|
||||
private static final String BASE = """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: ~/.config/herdr/herdr.sock
|
||||
""";
|
||||
|
||||
private static final String HEALTH_DISABLED_WITH_WEBHOOK = BASE + """
|
||||
health:
|
||||
enabled: false
|
||||
notifications:
|
||||
mode: webhook
|
||||
""";
|
||||
|
||||
private static final String HEALTH_DETECTION_ONLY = BASE + """
|
||||
health:
|
||||
enabled: true
|
||||
""";
|
||||
|
||||
private static final String HEALTH_FULL = BASE + """
|
||||
health:
|
||||
enabled: true
|
||||
notifications:
|
||||
mode: webhook
|
||||
""";
|
||||
|
||||
private static final String HEALTH_ABSENT = BASE;
|
||||
|
||||
@Test
|
||||
@DisplayName("enabled: false reports off, even with a webhook configured")
|
||||
void disabledHealthReportsOff(@TempDir Path dir) throws Exception {
|
||||
FleetMcp.HealthCoverageSource source = sourceFor(dir, HEALTH_DISABLED_WITH_WEBHOOK);
|
||||
|
||||
assertEquals("off", source.value().get(),
|
||||
"health.enabled: false must report 'off' regardless of notifications — flipping "
|
||||
+ "the 'enabled' argument at the HealthCoverageSource call site would report "
|
||||
+ "'full' here instead");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("enabled with no notifications reports detection-only")
|
||||
void enabledWithoutNotificationsReportsDetectionOnly(@TempDir Path dir) throws Exception {
|
||||
FleetMcp.HealthCoverageSource source = sourceFor(dir, HEALTH_DETECTION_ONLY);
|
||||
|
||||
assertEquals("detection-only", source.value().get());
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("enabled with a webhook configured reports full")
|
||||
void enabledWithNotificationsReportsFull(@TempDir Path dir) throws Exception {
|
||||
FleetMcp.HealthCoverageSource source = sourceFor(dir, HEALTH_FULL);
|
||||
|
||||
assertEquals("full", source.value().get(),
|
||||
"health.enabled: true with notifications.mode: webhook must report 'full' — "
|
||||
+ "swapping 'full' and 'detection-only' at the call site, or breaking the "
|
||||
+ "enabled/notificationConfigured argument pairing, would report "
|
||||
+ "'detection-only' here instead");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("an absent health: block reports off")
|
||||
void absentHealthBlockReportsOff(@TempDir Path dir) throws Exception {
|
||||
FleetMcp.HealthCoverageSource source = sourceFor(dir, HEALTH_ABSENT);
|
||||
|
||||
assertEquals("off", source.value().get());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #426, the live-wiring half: {@code health.notifications} is a {@code
|
||||
* ConfigRef.SPLIT_KEYS} entry, and {@link Fleetd#healthCoverageSource} reads {@code
|
||||
* config.get().health()} live (not the frozen startup {@code cfg}) — exactly like {@link
|
||||
* Fleetd#capacitySource}'s {@code maxLoad} ({@code FleetdCapacitySourceWiringTest}'s {@code
|
||||
* reloadedMaxLoadStillChangesWhatFleetListReports}). A hot-reloaded notifications block must
|
||||
* change what {@code fleet_list} reports without a restart; a fix that froze the whole source
|
||||
* against the startup snapshot would silently break that and every other test above would stay
|
||||
* green, since none of them reload.
|
||||
*/
|
||||
@Test
|
||||
@DisplayName("a hot notifications reload still changes what fleet_list reports")
|
||||
void reloadedNotificationsStillChangeWhatFleetListReports(@TempDir Path dir) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, HEALTH_DETECTION_ONLY);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = new ConfigRef(file, cfg);
|
||||
|
||||
FleetMcp.HealthCoverageSource source = Fleetd.healthCoverageSource(config);
|
||||
assertEquals("detection-only", source.value().get(),
|
||||
"sanity: detection-only before any reload");
|
||||
|
||||
Files.writeString(file, HEALTH_FULL);
|
||||
assertTrue(config.reload().applied());
|
||||
// The live snapshot now carries the webhook — proves the reload really happened and this
|
||||
// test is not accidentally passing because nothing changed.
|
||||
assertTrue(config.get().health().notifications() != null
|
||||
&& config.get().health().notifications().configured(),
|
||||
"sanity: the reloaded config really carries a configured webhook");
|
||||
|
||||
assertEquals("full", source.value().get(),
|
||||
"the SAME HealthCoverageSource instance must reflect a reloaded notifications "
|
||||
+ "block without a restart — health.notifications is read live off "
|
||||
+ "config.get(), exactly like capacitySource's maxLoad");
|
||||
}
|
||||
|
||||
private static FleetMcp.HealthCoverageSource sourceFor(Path dir, String yaml) throws Exception {
|
||||
Path file = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(file, yaml);
|
||||
FleetConfig cfg = FleetConfig.load(file);
|
||||
ConfigRef config = new ConfigRef(file, cfg);
|
||||
return Fleetd.healthCoverageSource(config);
|
||||
}
|
||||
}
|
||||
@@ -1,15 +1,12 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
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.msg.LeadMailbox;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
@@ -52,29 +49,22 @@ class FleetdLeadMailboxSelectionTest {
|
||||
}
|
||||
}
|
||||
|
||||
private static ListAppender<ILoggingEvent> captureFleetdLogs() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static String joined(ListAppender<ILoggingEvent> appender, Level level) {
|
||||
return appender.list.stream().filter(e -> e.getLevel() == level)
|
||||
private static String joined(CapturedLog captured, Level level) {
|
||||
return captured.events().stream().filter(e -> e.getLevel() == level)
|
||||
.map(ILoggingEvent::getFormattedMessage).reduce("", (a, b) -> a + "\n" + b);
|
||||
}
|
||||
|
||||
@Test
|
||||
void noCoordinatorBlockLeavesTheFeatureOffSilently() {
|
||||
var appender = captureFleetdLogs();
|
||||
var opener = new RecordingOpener();
|
||||
try (var captured = CapturedLog.of(Fleetd.class)) {
|
||||
var opener = new RecordingOpener();
|
||||
|
||||
assertNull(Fleetd.openLeadMailbox(null, Map.of(), opener));
|
||||
assertNull(Fleetd.openLeadMailbox(null, Map.of(), opener));
|
||||
|
||||
assertNull(opener.offeredUri, "nothing configured means nothing is opened");
|
||||
assertEquals("", joined(appender, Level.WARN),
|
||||
"an opt-in feature nobody asked for must not warn on every boot");
|
||||
assertNull(opener.offeredUri, "nothing configured means nothing is opened");
|
||||
assertEquals("", joined(captured, Level.WARN),
|
||||
"an opt-in feature nobody asked for must not warn on every boot");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -114,31 +104,33 @@ class FleetdLeadMailboxSelectionTest {
|
||||
|
||||
@Test
|
||||
void warnsAndStaysOffWhenSelfIdIsMissing() {
|
||||
var appender = captureFleetdLogs();
|
||||
var opener = new RecordingOpener();
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null, null);
|
||||
try (var captured = CapturedLog.of(Fleetd.class)) {
|
||||
var opener = new RecordingOpener();
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null, null);
|
||||
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener));
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener));
|
||||
|
||||
assertNull(opener.offeredUri, "a mailbox with no owning coord-id has no queue to declare");
|
||||
String warns = joined(appender, Level.WARN);
|
||||
assertTrue(warns.contains("coordinator.selfId"), () -> "say which key is missing: " + warns);
|
||||
assertFalse(warns.contains(SECRET), () -> "the URI's password must never be logged: " + warns);
|
||||
assertNull(opener.offeredUri, "a mailbox with no owning coord-id has no queue to declare");
|
||||
String warns = joined(captured, Level.WARN);
|
||||
assertTrue(warns.contains("coordinator.selfId"), () -> "say which key is missing: " + warns);
|
||||
assertFalse(warns.contains(SECRET), () -> "the URI's password must never be logged: " + warns);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void warnsAndStaysOffWhenTheBrokerIsUnreachableAtBoot() {
|
||||
var appender = captureFleetdLogs();
|
||||
var opener = new RecordingOpener();
|
||||
opener.unreachable = true;
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
|
||||
try (var captured = CapturedLog.of(Fleetd.class)) {
|
||||
var opener = new RecordingOpener();
|
||||
opener.unreachable = true;
|
||||
var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null, null);
|
||||
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener),
|
||||
"a down coordination broker turns the feature off; it must never take the daemon down");
|
||||
assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener),
|
||||
"a down coordination broker turns the feature off; it must never take the daemon down");
|
||||
|
||||
String warns = joined(appender, Level.WARN);
|
||||
assertTrue(warns.contains("coord.example"), () -> "name the host that failed: " + warns);
|
||||
assertFalse(warns.contains(SECRET), () -> "with credentials stripped: " + warns);
|
||||
assertTrue(warns.contains("Connection refused"), () -> "and the real reason: " + warns);
|
||||
String warns = joined(captured, Level.WARN);
|
||||
assertTrue(warns.contains("coord.example"), () -> "name the host that failed: " + warns);
|
||||
assertFalse(warns.contains(SECRET), () -> "with credentials stripped: " + warns);
|
||||
assertTrue(warns.contains("Connection refused"), () -> "and the real reason: " + warns);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,15 +1,13 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
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.msg.AmqpReplyInbox;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.net.ServerSocket;
|
||||
import java.util.List;
|
||||
@@ -50,43 +48,34 @@ class FleetdReplyInboxSelectionTest {
|
||||
}
|
||||
}
|
||||
|
||||
private static ListAppender<ILoggingEvent> attach() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
private static CapturedLog attach() {
|
||||
// logback-test.xml pins dev.ltms.fleet to WARN; raise it so INFO selection lines are captured.
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
return CapturedLog.at(Fleetd.class, Level.INFO);
|
||||
}
|
||||
|
||||
private static void detach(ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
|
||||
}
|
||||
|
||||
private static void assertNoLogContains(ListAppender<ILoggingEvent> appender, String secret) {
|
||||
assertTrue(appender.list.stream().noneMatch(e -> e.getFormattedMessage().contains(secret)),
|
||||
private static void assertNoLogContains(List<ILoggingEvent> events, String secret) {
|
||||
assertTrue(events.stream().noneMatch(e -> e.getFormattedMessage().contains(secret)),
|
||||
"no log line may contain the resolved URI's password");
|
||||
}
|
||||
|
||||
@Test
|
||||
void uriEnvSetAndPresentSelectsAmqpWithTheResolvedUri() {
|
||||
FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), false, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), false, (opener, events, inbox) -> {
|
||||
assertEquals(RESOLVED_URI, opener.offeredUri,
|
||||
"the daemon must connect with the value resolved from uriEnv — selection, not just parse");
|
||||
assertEquals(opener.inbox, inbox, "the AMQP opener's inbox is what is selected");
|
||||
assertNoLogContains(appender, SECRET);
|
||||
assertNoLogContains(events, SECRET);
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void uriEnvSetButVariableMissingFallsBackToInMemoryAndWarns() {
|
||||
FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
|
||||
recording(broker, Map.of(), false, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of(), false, (opener, events, inbox) -> {
|
||||
assertInstanceOf(InMemoryReplyInbox.class, inbox);
|
||||
assertNull(opener.offeredUri, "AMQP must never be attempted when the variable is missing");
|
||||
assertTrue(hasWarnContaining(appender, "LAVINMQ_URI") && hasWarnContaining(appender, "DISABLED"),
|
||||
assertTrue(hasWarnContaining(events, "LAVINMQ_URI") && hasWarnContaining(events, "DISABLED"),
|
||||
"a missing uriEnv variable must warn loudly, not fail silently");
|
||||
});
|
||||
}
|
||||
@@ -94,10 +83,10 @@ class FleetdReplyInboxSelectionTest {
|
||||
@Test
|
||||
void uriEnvSetButVariableBlankFallsBackToInMemoryAndWarns() {
|
||||
FleetConfig.Broker broker = new FleetConfig.Broker("amqp://user:lame@old:5672/", "LAVINMQ_URI", null);
|
||||
recording(broker, Map.of("LAVINMQ_URI", " "), false, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of("LAVINMQ_URI", " "), false, (opener, events, inbox) -> {
|
||||
assertInstanceOf(InMemoryReplyInbox.class, inbox);
|
||||
assertNull(opener.offeredUri, "a blank env value must not select AMQP, not even via the literal uri");
|
||||
assertTrue(hasWarnContaining(appender, "LAVINMQ_URI"),
|
||||
assertTrue(hasWarnContaining(events, "LAVINMQ_URI"),
|
||||
"a blank uriEnv value must warn, and must not fall back to the literal uri");
|
||||
});
|
||||
}
|
||||
@@ -106,22 +95,22 @@ class FleetdReplyInboxSelectionTest {
|
||||
void bothUriAndUriEnvSetUriEnvWinsDeterministically() {
|
||||
FleetConfig.Broker broker
|
||||
= new FleetConfig.Broker("amqp://user:oldpw@old.example:5672/", "LAVINMQ_URI", null);
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), false, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), false, (opener, events, inbox) -> {
|
||||
assertEquals(RESOLVED_URI, opener.offeredUri,
|
||||
"uriEnv must win over uri, deterministically, every run");
|
||||
assertTrue(appender.list.stream().anyMatch(e -> e.getFormattedMessage().contains("broker.uri is ignored")),
|
||||
assertTrue(events.stream().anyMatch(e -> e.getFormattedMessage().contains("broker.uri is ignored")),
|
||||
"must log that the literal uri is ignored when uriEnv is set");
|
||||
assertNoLogContains(appender, SECRET);
|
||||
assertNoLogContains(appender, "oldpw");
|
||||
assertNoLogContains(events, SECRET);
|
||||
assertNoLogContains(events, "oldpw");
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void unreachableBrokerStartsDaemonWithInMemoryInboxAndLoudWarning() {
|
||||
FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), true, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), true, (opener, events, inbox) -> {
|
||||
assertInstanceOf(InMemoryReplyInbox.class, inbox, "an unreachable broker must NOT stop the daemon");
|
||||
String warn = appender.list.stream()
|
||||
String warn = events.stream()
|
||||
.filter(e -> e.getLevel() == Level.WARN)
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.reduce("", (a, b) -> a + "\n" + b)
|
||||
@@ -129,16 +118,16 @@ class FleetdReplyInboxSelectionTest {
|
||||
assertTrue(warn.contains("durable") && warn.contains("soft-state"),
|
||||
"the warning must say exactly what was lost: durable delivery off, replies soft-state");
|
||||
assertTrue(!warn.contains(SECRET), "the failing URI must be logged with credentials stripped");
|
||||
assertNoLogContains(appender, SECRET);
|
||||
assertNoLogContains(events, SECRET);
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void noBrokerConfiguredStaysQuietInMemory() {
|
||||
FleetConfig.Broker broker = null;
|
||||
recording(broker, Map.of(), false, (opener, appender, inbox) -> {
|
||||
recording(broker, Map.of(), false, (opener, events, inbox) -> {
|
||||
assertInstanceOf(InMemoryReplyInbox.class, inbox);
|
||||
assertTrue(appender.list.stream().noneMatch(e -> e.getLevel() == Level.WARN),
|
||||
assertTrue(events.stream().noneMatch(e -> e.getLevel() == Level.WARN),
|
||||
"no broker configured must keep the existing QUIET in-memory path — no warning");
|
||||
});
|
||||
}
|
||||
@@ -152,21 +141,18 @@ class FleetdReplyInboxSelectionTest {
|
||||
closedPort = s.getLocalPort();
|
||||
}
|
||||
FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
try (CapturedLog log = attach()) {
|
||||
ReplyInbox inbox = Fleetd.selectReplyInbox(
|
||||
broker, Map.of("LAVINMQ_URI", "amqp://user:" + SECRET + "@127.0.0.1:" + closedPort + "/vh"),
|
||||
AmqpReplyInbox::open);
|
||||
assertInstanceOf(InMemoryReplyInbox.class, inbox,
|
||||
"a genuinely unreachable broker (real AmqpReplyInbox::open) must fall back to in-memory");
|
||||
} finally {
|
||||
detach(appender);
|
||||
assertNoLogContains(log.events(), SECRET);
|
||||
}
|
||||
assertNoLogContains(appender, SECRET);
|
||||
}
|
||||
|
||||
private boolean hasWarnContaining(ListAppender<ILoggingEvent> appender, String fragment) {
|
||||
return appender.list.stream().anyMatch(e ->
|
||||
private boolean hasWarnContaining(List<ILoggingEvent> events, String fragment) {
|
||||
return events.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.WARN && e.getFormattedMessage().contains(fragment));
|
||||
}
|
||||
|
||||
@@ -174,18 +160,14 @@ class FleetdReplyInboxSelectionTest {
|
||||
private void recording(FleetConfig.Broker broker, Map<String, String> env, boolean unreachable, Check check) {
|
||||
RecordingAmqp opener = new RecordingAmqp();
|
||||
opener.unreachable = unreachable;
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
ReplyInbox inbox;
|
||||
try {
|
||||
inbox = Fleetd.selectReplyInbox(broker, env, opener);
|
||||
} finally {
|
||||
detach(appender);
|
||||
try (CapturedLog log = attach()) {
|
||||
ReplyInbox inbox = Fleetd.selectReplyInbox(broker, env, opener);
|
||||
check.run(opener, log.events(), inbox);
|
||||
}
|
||||
check.run(opener, appender, inbox);
|
||||
}
|
||||
|
||||
@FunctionalInterface
|
||||
private interface Check {
|
||||
void run(RecordingAmqp opener, ListAppender<ILoggingEvent> appender, ReplyInbox inbox);
|
||||
void run(RecordingAmqp opener, List<ILoggingEvent> events, ReplyInbox inbox);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,15 +1,12 @@
|
||||
package dev.ltms.fleet.auth;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -24,28 +21,21 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
class AuditLogTest {
|
||||
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
private ListAppender<ILoggingEvent> appender;
|
||||
private ch.qos.logback.classic.Logger auditLogger;
|
||||
private CapturedLog auditLog;
|
||||
|
||||
@BeforeEach
|
||||
void attach() {
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
auditLogger = ctx.getLogger("audit");
|
||||
appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
auditLogger.addAppender(appender);
|
||||
auditLogger.setLevel(Level.INFO);
|
||||
auditLog = CapturedLog.at("audit", Level.INFO);
|
||||
}
|
||||
|
||||
@AfterEach
|
||||
void detach() {
|
||||
auditLogger.detachAppender(appender);
|
||||
auditLog.close();
|
||||
}
|
||||
|
||||
private JsonNode onlyRecord() throws Exception {
|
||||
assertEquals(1, appender.list.size(), "exactly one audit line expected");
|
||||
String line = appender.list.getFirst().getFormattedMessage();
|
||||
assertEquals(1, auditLog.events().size(), "exactly one audit line expected");
|
||||
String line = auditLog.events().getFirst().getFormattedMessage();
|
||||
return mapper.readTree(line); // throws if the line is not valid JSON
|
||||
}
|
||||
|
||||
|
||||
@@ -186,6 +186,43 @@ class CallerResolverTest {
|
||||
assertEquals(Role.PRIMARY, r.resolve("127.0.0.1", 99, "BEARER s3cret").role());
|
||||
}
|
||||
|
||||
// ── fleetd #505: a herdr error DURING THE SCAN must not be conflated with "not a worker" ──────
|
||||
// #317 (above) covers a failed lsof lookup. This is the other input to the same decision: the
|
||||
// lsof lookup succeeds (a real pid), but PaneLocator's own pane scan hits a herdr error on the
|
||||
// pane that owns that pid — so c.resolved() is true and c.terminal() is null, exactly like a
|
||||
// real primary. c.scanComplete() is what tells them apart.
|
||||
|
||||
/**
|
||||
* The discriminating case named in the ticket: the error must land on the pane that DOES own
|
||||
* the caller's pid, or the test proves nothing (any other pane's failure is invisible to the
|
||||
* scan's outcome, since a match found elsewhere is definitive regardless).
|
||||
*/
|
||||
@Test
|
||||
void aHerdrErrorOnTheOwningPaneDuringTheScanIsRefusedNotPromotedToPrimary() {
|
||||
FakeHerdr failing = new FakeHerdr().processInfoFailsForPane("w2:p7", "transient");
|
||||
ConnectionIdentity incomplete = new ConnectionIdentity(new PaneLocator(failing), _ -> FakeHerdr.WORKER_PID);
|
||||
|
||||
Principal p = new CallerResolver(incomplete).resolve("127.0.0.1", 55555, null);
|
||||
|
||||
assertEquals(Role.ANONYMOUS, p.role(),
|
||||
"an incomplete pane scan must never be read as a clean negative and promoted to primary");
|
||||
}
|
||||
|
||||
/**
|
||||
* The companion invariant: a herdr error on a DIFFERENT, non-owning pane must not turn every
|
||||
* mid-scan teardown into a refusal — the real match is still found and resolves as a worker.
|
||||
*/
|
||||
@Test
|
||||
void aHerdrErrorOnANonOwningPaneStillResolvesTheRealWorker() {
|
||||
FakeHerdr vanishedElsewhere = new FakeHerdr().processInfoFailsForPane("w2:p9", "pane_not_found");
|
||||
ConnectionIdentity id = new ConnectionIdentity(new PaneLocator(vanishedElsewhere), _ -> FakeHerdr.WORKER_PID);
|
||||
|
||||
Principal p = new CallerResolver(id).resolve("127.0.0.1", 55555, null);
|
||||
|
||||
assertEquals(Role.WORKER, p.role());
|
||||
assertEquals("term_a", p.terminal());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNonLoopbackCallerIsNeverThePrimaryUnderLoopbackTrust() {
|
||||
// Defence in depth: startup already refuses this pairing (validateAuthExposure), but if a
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
package dev.ltms.fleet.health;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #426: {@link FleetHealthMonitor#coverage} had zero references anywhere in the test tree —
|
||||
* not the method, not either output string, not the field it populates. This pins the method's own
|
||||
* three branches directly.
|
||||
*
|
||||
* <p>This is the easy half. It proves the method words each combination correctly, but it proves
|
||||
* nothing about whether {@code Fleetd.java}'s two call sites pass the right argument in the right
|
||||
* position — see {@code FleetdHealthCoverageSourceWiringTest} (package {@code dev.ltms.fleet}) for
|
||||
* the half that actually guards the call site, following the measured fleetd #415 lesson that a
|
||||
* thoroughly-tested method and an untested argument pairing at its call site are different risks.
|
||||
*
|
||||
* <p><strong>The three output strings are load-bearing and must not change here.</strong> {@code
|
||||
* "detection-only"} is read live off a running daemon's {@code fleet_list} today (measured
|
||||
* 2026-09-12) — this test intentionally asserts the exact literal strings so a future edit to the
|
||||
* wording trips it here first.
|
||||
*/
|
||||
class FleetHealthMonitorCoverageTest {
|
||||
|
||||
@Test
|
||||
void disabledIsOffRegardlessOfNotificationConfig() {
|
||||
assertEquals("off", FleetHealthMonitor.coverage(false, false));
|
||||
assertEquals("off", FleetHealthMonitor.coverage(false, true));
|
||||
}
|
||||
|
||||
@Test
|
||||
void enabledWithoutNotificationsIsDetectionOnly() {
|
||||
assertEquals("detection-only", FleetHealthMonitor.coverage(true, false));
|
||||
}
|
||||
|
||||
@Test
|
||||
void enabledWithNotificationsIsFull() {
|
||||
assertEquals("full", FleetHealthMonitor.coverage(true, true));
|
||||
}
|
||||
}
|
||||
@@ -44,6 +44,7 @@ public final class FakeHerdr implements HerdrClient {
|
||||
private int workerTabPaneCount = 1;
|
||||
private String paneCloseErrorCode = null;
|
||||
private final Map<String, String> paneCloseErrorCodeFor = new ConcurrentHashMap<>();
|
||||
private final Map<String, String> processInfoErrorCodeFor = new ConcurrentHashMap<>();
|
||||
private String tabCloseErrorCode = null;
|
||||
private final Map<String, String> tabCloseErrorCodeFor = new ConcurrentHashMap<>();
|
||||
private String agentSendErrorCode = null;
|
||||
@@ -144,6 +145,18 @@ public final class FakeHerdr implements HerdrClient {
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Make {@code pane.process_info} fail with this herdr error code, but only for the given
|
||||
* {@code pane_id} — every other pane's {@code pane.process_info} still succeeds. Models a
|
||||
* transient herdr failure partway through a {@link PaneLocator} pid→pane scan (fleetd #505):
|
||||
* the scan must be able to tell "this pane does not own the pid" apart from "the scan could
|
||||
* not check this pane at all", instead of collapsing both into one {@code false}.
|
||||
*/
|
||||
public FakeHerdr processInfoFailsForPane(String paneId, String code) {
|
||||
this.processInfoErrorCodeFor.put(paneId, code);
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Set the {@code agent_status} that {@code agent.get} reports (drives the injector). */
|
||||
public FakeHerdr agentStatus(String status) {
|
||||
this.agentStatus = status;
|
||||
@@ -407,6 +420,13 @@ public final class FakeHerdr implements HerdrClient {
|
||||
{"pane_id":"w2:p9","terminal_id":"term_shell","workspace_id":"w2","tab_id":"w2:t8"}]}""");
|
||||
case "pane.process_info" -> {
|
||||
Object paneId = params instanceof java.util.Map<?, ?> m ? m.get("pane_id") : null;
|
||||
String failCode = paneId == null ? null
|
||||
: processInfoErrorCodeFor.get(String.valueOf(paneId));
|
||||
if (failCode != null) {
|
||||
throw new HerdrException(
|
||||
"herdr error [" + failCode + "]: pane.process_info failed",
|
||||
failCode, null);
|
||||
}
|
||||
yield "w2:p7".equals(paneId)
|
||||
? mapper.readTree(("""
|
||||
{"type":"pane_process_info","process_info":{"pane_id":"w2:p7","shell_pid":%d,
|
||||
|
||||
@@ -37,7 +37,7 @@ class PaneLocatorContractTest {
|
||||
.path("pane").path("terminal_id").asText(null);
|
||||
assertNotNull(terminalId, "seed pane should carry a terminal_id");
|
||||
|
||||
assertEquals(terminalId, new PaneLocator(herdr).terminalForPid(shellPid),
|
||||
assertEquals(terminalId, new PaneLocator(herdr).terminalForPid(shellPid).terminal(),
|
||||
"a real PID must resolve back to its own pane's terminal_id");
|
||||
} finally {
|
||||
spaces.closeTab(tab.tab().tabId());
|
||||
|
||||
@@ -15,18 +15,20 @@ class PaneLocatorTest {
|
||||
|
||||
@Test
|
||||
void resolvesTerminalForAForegroundPid() {
|
||||
assertEquals("term_a", loc.terminalForPid(FakeHerdr.WORKER_PID));
|
||||
assertEquals("term_a", loc.terminalForPid(FakeHerdr.WORKER_PID).terminal());
|
||||
}
|
||||
|
||||
@Test
|
||||
void nullForAPidInNoPane() {
|
||||
assertNull(loc.terminalForPid(999_999));
|
||||
PaneLocator.Lookup outcome = loc.terminalForPid(999_999);
|
||||
assertNull(outcome.terminal());
|
||||
assertTrue(outcome.complete(), "a full, error-free scan that finds no match is complete");
|
||||
}
|
||||
|
||||
@Test
|
||||
void nullForNonPositivePid() {
|
||||
assertNull(loc.terminalForPid(0));
|
||||
assertNull(loc.terminalForPid(-1));
|
||||
assertNull(loc.terminalForPid(0).terminal());
|
||||
assertNull(loc.terminalForPid(-1).terminal());
|
||||
}
|
||||
|
||||
// --- two-daemon fallback (CB-185) -----------------------------------------
|
||||
@@ -38,7 +40,7 @@ class PaneLocatorTest {
|
||||
HerdrClient lead = new FakeHerdr().withNoPanes();
|
||||
HerdrClient member = new FakeHerdr();
|
||||
PaneLocator two = new PaneLocator(lead, member);
|
||||
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID));
|
||||
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID).terminal());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -48,13 +50,13 @@ class PaneLocatorTest {
|
||||
HerdrClient lead = new FakeHerdr();
|
||||
HerdrClient member = new FakeHerdr().withNoPanes();
|
||||
PaneLocator two = new PaneLocator(lead, member);
|
||||
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID));
|
||||
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID).terminal());
|
||||
}
|
||||
|
||||
@Test
|
||||
void nullWhenNeitherClientHasTheMatch() {
|
||||
PaneLocator two = new PaneLocator(new FakeHerdr().withNoPanes(), new FakeHerdr().withNoPanes());
|
||||
assertNull(two.terminalForPid(FakeHerdr.WORKER_PID));
|
||||
assertNull(two.terminalForPid(FakeHerdr.WORKER_PID).terminal());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -63,7 +65,7 @@ class PaneLocatorTest {
|
||||
// must behave exactly like the one-arg constructor, including making only one herdr call.
|
||||
FakeHerdr shared = new FakeHerdr();
|
||||
PaneLocator two = new PaneLocator(shared, shared);
|
||||
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID));
|
||||
assertEquals("term_a", two.terminalForPid(FakeHerdr.WORKER_PID).terminal());
|
||||
long paneListCalls = shared.calls.stream().filter(c -> c.method().equals("pane.list")).count();
|
||||
assertEquals(1, paneListCalls, "same-object lead/member must scan exactly once, not twice");
|
||||
}
|
||||
@@ -75,7 +77,7 @@ class PaneLocatorTest {
|
||||
// Regression: a pid with no parent chain at all — no ancestry walk is needed to match it.
|
||||
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
|
||||
PaneLocator loc = new PaneLocator(pane, new FakeParentResolver());
|
||||
assertEquals("term_x", loc.terminalForPid(5000));
|
||||
assertEquals("term_x", loc.terminalForPid(5000).terminal());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -83,7 +85,7 @@ class PaneLocatorTest {
|
||||
// Regression: same as above, but matching via the foreground-processes list.
|
||||
OnePaneHerdr pane = new OnePaneHerdr("term_x", "pX", 5000, 6000);
|
||||
PaneLocator loc = new PaneLocator(pane, new FakeParentResolver());
|
||||
assertEquals("term_x", loc.terminalForPid(6000));
|
||||
assertEquals("term_x", loc.terminalForPid(6000).terminal());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -97,7 +99,7 @@ class PaneLocatorTest {
|
||||
.parent(7002, 7001) // grandchild -> child
|
||||
.parent(7001, 5000); // child -> shell (the pane's shell_pid)
|
||||
PaneLocator loc = new PaneLocator(pane, parents);
|
||||
assertEquals("term_x", loc.terminalForPid(7002));
|
||||
assertEquals("term_x", loc.terminalForPid(7002).terminal());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -110,7 +112,7 @@ class PaneLocatorTest {
|
||||
.parent(9002, 9001)
|
||||
.parent(9001, 9000); // chain never reaches 5000 or 6000
|
||||
PaneLocator loc = new PaneLocator(pane, parents);
|
||||
assertNull(loc.terminalForPid(9002));
|
||||
assertNull(loc.terminalForPid(9002).terminal());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -122,7 +124,7 @@ class PaneLocatorTest {
|
||||
.parent(100, 101)
|
||||
.parent(101, 100); // cycle, never reaches the pane's pids
|
||||
PaneLocator loc = new PaneLocator(pane, parents);
|
||||
assertNull(loc.terminalForPid(100));
|
||||
assertNull(loc.terminalForPid(100).terminal());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -141,11 +143,69 @@ class PaneLocatorTest {
|
||||
};
|
||||
HerdrClient noPanes = new FakeHerdr().withNoPanes();
|
||||
PaneLocator two = new PaneLocator(noPanes, pane, counting);
|
||||
assertEquals("term_x", two.terminalForPid(7002));
|
||||
assertEquals("term_x", two.terminalForPid(7002).terminal());
|
||||
assertEquals(3, calls.get(), "ancestry must be walked once (3 lookups: 7002, 7001, 5000), "
|
||||
+ "not re-walked per herdr client");
|
||||
}
|
||||
|
||||
// --- fleetd #505: a herdr error during the scan must not read as a clean negative ---------
|
||||
|
||||
@Test
|
||||
void anErrorOnThePaneThatOwnsThePidMakesTheScanIncompleteNotAClearNegative() {
|
||||
// The discriminating case: pane.process_info fails for exactly the pane that DOES own the
|
||||
// caller's pid ("w2:p7", term_a). Before the fix, that failure was swallowed into a plain
|
||||
// "does not own it" and the scan finished with a clean-looking null — indistinguishable
|
||||
// from a real primary. It must now report incomplete, not a definite null.
|
||||
FakeHerdr herdr = new FakeHerdr().processInfoFailsForPane("w2:p7", "transient");
|
||||
PaneLocator loc = new PaneLocator(herdr);
|
||||
|
||||
PaneLocator.Lookup outcome = loc.terminalForPid(FakeHerdr.WORKER_PID);
|
||||
|
||||
assertNull(outcome.terminal(), "the failing pane's ownership could not be confirmed");
|
||||
assertFalse(outcome.complete(),
|
||||
"a scan that could not check the owning pane must not report as complete");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aVanishedPaneThatIsNotTheMatchLeavesAnOtherwiseSuccessfulScanComplete() {
|
||||
// The companion invariant: a DIFFERENT pane (not the caller's own) failing mid-scan must
|
||||
// not turn every mid-scan teardown into a refusal — the real match is still found, and the
|
||||
// scan is still reported complete.
|
||||
FakeHerdr herdr = new FakeHerdr().processInfoFailsForPane("w2:p9", "pane_not_found");
|
||||
PaneLocator loc = new PaneLocator(herdr);
|
||||
|
||||
PaneLocator.Lookup outcome = loc.terminalForPid(FakeHerdr.WORKER_PID);
|
||||
|
||||
assertEquals("term_a", outcome.terminal());
|
||||
assertTrue(outcome.complete(), "a positive match elsewhere in the scan is definitive");
|
||||
}
|
||||
|
||||
// --- fleetd #509: the completeness fold across clients must not collapse to "last wins" ----
|
||||
|
||||
@Test
|
||||
void anEarlierClientsErrorSurvivesALaterClientsCleanNegative() {
|
||||
// terminalForPid folds each client's Lookup.complete() with
|
||||
// complete = complete && outcome.complete();
|
||||
// (PaneLocator.java:117). With a SINGLE client, a fold that keeps only the last outcome
|
||||
// (dropping the "complete &&" prefix) agrees with the real fold — which is why 14 of the
|
||||
// 15 pre-existing tests never catch that mutation: none of them vary the number of clients.
|
||||
// Here the LEAD client errors on exactly the pane that would have owned the pid (so its
|
||||
// scan is incomplete AND finds no match), and the MEMBER client cleanly reports no panes
|
||||
// at all (a complete, negative scan). The real fold ANDs the two into false. A fold that
|
||||
// just keeps the last client's outcome would read this as a clean true — the earlier
|
||||
// error is erased, and CallerResolver.java:254 would read scanComplete() as true and
|
||||
// promote an unverified caller to the primary.
|
||||
HerdrClient lead = new FakeHerdr().processInfoFailsForPane("w2:p7", "transient");
|
||||
HerdrClient member = new FakeHerdr().withNoPanes();
|
||||
PaneLocator two = new PaneLocator(lead, member);
|
||||
|
||||
PaneLocator.Lookup outcome = two.terminalForPid(FakeHerdr.WORKER_PID);
|
||||
|
||||
assertNull(outcome.terminal(), "the pane that could have owned the pid was never checked");
|
||||
assertFalse(outcome.complete(),
|
||||
"an earlier client's error must survive a later client's clean negative");
|
||||
}
|
||||
|
||||
/** Minimal single-pane {@link HerdrClient} fake, purpose-built for the ancestry tests above. */
|
||||
private static final class OnePaneHerdr implements HerdrClient {
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
|
||||
@@ -1,9 +1,7 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
@@ -12,8 +10,8 @@ import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TestTurnTokens;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.Set;
|
||||
import java.util.regex.Pattern;
|
||||
@@ -470,15 +468,7 @@ class CompletionResolverTest {
|
||||
// CB-564: this used to be a bare DEBUG "failed send to X via turn-stall fallback" — a symptom
|
||||
// with no cause, and below the level anyone watching for member health would see. A fail that
|
||||
// resolves a caller's blocked send is at least WARN and must carry the reason.
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger resolverLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(CompletionResolver.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
resolverLog.addAppender(appender);
|
||||
resolverLog.setLevel(Level.WARN);
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.at(CompletionResolver.class, Level.WARN)) {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous, ExhaustedPatternLookup.none(), ExhaustionSink.none());
|
||||
@@ -486,7 +476,7 @@ class CompletionResolverTest {
|
||||
|
||||
resolver.fail("term_a", null);
|
||||
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
@@ -494,8 +484,6 @@ class CompletionResolverTest {
|
||||
assertTrue(warn.contains("term_a"), "the log names the target: " + warn);
|
||||
assertTrue(warn.contains("stuck on an error screen"), "the log carries the reason: " + warn);
|
||||
assertTrue(waiter.isDone());
|
||||
} finally {
|
||||
resolverLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1,16 +1,16 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.msg.TestTurnTokens;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
@@ -19,6 +19,8 @@ import java.util.Set;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -491,23 +493,15 @@ class InjectorTest {
|
||||
void readinessGraceExpiryIsLogged() {
|
||||
// CB-562: the grace-expiry path used to clear the queue silently, so a message that never
|
||||
// reached the worker's pane surfaced elsewhere as an unrelated turn-stall failure. Assert the
|
||||
// expiry now names the real cause. (ListAppender capture pattern mirrors AuditLogTest.)
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger injectorLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(Injector.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
injectorLog.addAppender(appender);
|
||||
injectorLog.setLevel(Level.WARN);
|
||||
try {
|
||||
// expiry now names the real cause. (CapturedLog pattern mirrors AuditLogTest.)
|
||||
try (CapturedLog log = CapturedLog.at(Injector.class, Level.WARN)) {
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> {
|
||||
});
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
@@ -516,8 +510,77 @@ class InjectorTest {
|
||||
assertTrue(warn.contains("never reached"), "the log names the real cause: " + warn);
|
||||
assertTrue(warn.contains("1 queued message"),
|
||||
"the log carries the failed message count: " + warn);
|
||||
} finally {
|
||||
injectorLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void readinessGraceExpiryLogsTheMeasuredPollCountNextToTheConfiguredBudget() {
|
||||
// fleetd #501, defect 1: READINESS_GRACE_POLLS (240) used to be printed twice — once as
|
||||
// "after {} polls" and once inside the parenthesised budget — even though the loop's own
|
||||
// counter (Target.notReadySincePoll) was in scope at the same call site. On THIS branch the
|
||||
// counter has just reached the threshold, so it equals the constant by construction and this
|
||||
// test cannot tell the two apart — it only pins that the message still carries both a poll
|
||||
// count and a labelled configured budget, using literal numbers (240, 60), never
|
||||
// READINESS_GRACE_POLLS or POLL_INTERVAL_MILLIS, so the assertion can't silently track a
|
||||
// constant change instead of catching a real regression.
|
||||
try (CapturedLog log = CapturedLog.at(Injector.class, Level.WARN)) {
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> {
|
||||
});
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
.orElse("no grace-expiry WARN logged");
|
||||
assertTrue(warn.contains("after 240 polls"),
|
||||
"must print the measured poll count as a plain number: " + warn);
|
||||
assertTrue(warn.contains("configured=240 polls/60s"),
|
||||
"must print the configured budget, clearly labelled: " + warn);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void readinessGraceExpiryLogsTheMeasuredElapsedTimeNotArithmeticOnConstants() {
|
||||
// fleetd #501, defect 2: the old line computed "({}s)" as READINESS_GRACE_POLLS *
|
||||
// POLL_INTERVAL_MILLIS / 1000 — arithmetic on two constants, never a measurement, and wrong
|
||||
// in the direction that says everything ran on schedule. This stub clock returns two FIXED
|
||||
// values (1_000ms at the first non-ready sample, 318_412ms at the poll that trips the grace)
|
||||
// whose difference — 317_412ms — does NOT equal 240 * POLL_INTERVAL_MILLIS (=60_000ms).
|
||||
// Asserting on that literal, non-derived number is what makes this test able to fail if the
|
||||
// production code goes back to printing the constant-arithmetic value instead of the
|
||||
// injected clock's measurement.
|
||||
long[] readings = {1_000L, 318_412L};
|
||||
AtomicInteger call = new AtomicInteger(0);
|
||||
LongSupplier stubClock = () -> {
|
||||
int i = call.getAndIncrement();
|
||||
if (i >= readings.length) {
|
||||
throw new AssertionError("nowMillis read more times than this fixture expects (" + i
|
||||
+ "); the readiness-not-ready branch should read the clock exactly twice — "
|
||||
+ "once to stamp the first non-ready sample, once at grace expiry");
|
||||
}
|
||||
return readings[i];
|
||||
};
|
||||
|
||||
try (CapturedLog log = CapturedLog.at(Injector.class, Level.WARN)) {
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> {
|
||||
}, stubClock);
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
.orElse("no grace-expiry WARN logged");
|
||||
assertTrue(warn.contains("elapsed=317412ms"), "must print the MEASURED elapsed time from "
|
||||
+ "the injected clock (318412 - 1000 = 317412), not an arithmetic value: " + warn);
|
||||
assertFalse(warn.contains("elapsed=60000ms"), "must not print "
|
||||
+ "READINESS_GRACE_POLLS * POLL_INTERVAL_MILLIS (240 * 250 = 60000ms) as if it "
|
||||
+ "were the measured elapsed time: " + warn);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -554,21 +617,13 @@ class InjectorTest {
|
||||
// CB-564: a vanished worker used to drop its queue with no log at all — the only trace was
|
||||
// whatever failed downstream (e.g. a caller's send timing out with no clue why). Assert the
|
||||
// drop itself now names the cause and the number of messages it failed.
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger injectorLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(Injector.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
injectorLog.addAppender(appender);
|
||||
injectorLog.setLevel(Level.WARN);
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.at(Injector.class, Level.WARN)) {
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> true, _ -> {
|
||||
});
|
||||
inj.enqueue(T, "orphan", TestTurnTokens.inert(T));
|
||||
inj.drop(T, new HerdrException("worker gone", "pane_not_found", null));
|
||||
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
@@ -576,8 +631,6 @@ class InjectorTest {
|
||||
assertTrue(warn.contains(T), "the log names the target terminal: " + warn);
|
||||
assertTrue(warn.contains("1 message"), "the log carries the failed message count: " + warn);
|
||||
assertTrue(warn.contains("worker gone"), "the log carries the real cause: " + warn);
|
||||
} finally {
|
||||
injectorLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -613,4 +666,111 @@ class InjectorTest {
|
||||
ExecutionException ex = assertThrows(ExecutionException.class, f::get);
|
||||
assertInstanceOf(HerdrException.class, ex.getCause());
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@link HerdrClient} that throws a non-{@link RuntimeException} {@link Error} from {@code
|
||||
* agent.prompt} instead of delegating — the fleetd #546 case: a stray non-RuntimeException
|
||||
* throwable (e.g. a {@code NoClassDefFoundError}, fleetd #413) escaping the send seam at
|
||||
* {@code Injector.java:383}. Records every call it sees itself, including the ones it throws
|
||||
* for, since the delegate's own recording is never reached for {@code agent.prompt} — so a test
|
||||
* can assert on exactly what this fake actually received.
|
||||
*/
|
||||
private static final class ErrorOnPrompt implements HerdrClient {
|
||||
private final FakeHerdr delegate;
|
||||
private final List<FakeHerdr.Call> calls = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
|
||||
private ErrorOnPrompt(FakeHerdr delegate) {
|
||||
this.delegate = delegate;
|
||||
}
|
||||
|
||||
List<FakeHerdr.Call> calls() {
|
||||
return calls;
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
calls.add(new FakeHerdr.Call(method, params));
|
||||
if (method.equals("agent.prompt")) {
|
||||
throw new AssertionError("simulated non-RuntimeException send failure (fleetd #546)");
|
||||
}
|
||||
return delegate.call(method, params);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
delegate.close();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Mirrors StatusPoller's own per-target {@code catch (Throwable)} (fleetd #538 / PR #543): the
|
||||
* production polling loop already swallows whatever escapes one target's round and comes back
|
||||
* for the next one. A unit test that calls {@code onStatus} directly (bypassing StatusPoller)
|
||||
* needs the same survival so it can observe what a SECOND round does, regardless of whether
|
||||
* fleetd #546's fix is present.
|
||||
*/
|
||||
private static void pollOnceSurviving(Injector inj, String target, AgentStatus status) {
|
||||
try {
|
||||
inj.onStatus(target, status);
|
||||
} catch (Throwable ignored) {
|
||||
// matches StatusPoller.loop's own catch (Throwable) added by PR #543
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anErrorFromSendRemovesTheMessageAndMarksItNotDelivered() {
|
||||
// fleetd #546, acceptance test 1: an Error (not a RuntimeException) escaping the send seam
|
||||
// at Injector.java:383 must still be caught, the message dropped from the queue, and its
|
||||
// state set to NOT_DELIVERED — never left QUEUED at the head.
|
||||
ErrorOnPrompt throwing = new ErrorOnPrompt(new FakeHerdr());
|
||||
Injector inj = new Injector(new AgentControl(throwing));
|
||||
Injector.Delivery delivery = inj.enqueue(T, "brief", TestTurnTokens.inert(T));
|
||||
|
||||
assertDoesNotThrow(() -> inj.onStatus(T, AgentStatus.IDLE),
|
||||
"fleetd #546: an Error from the send seam must be caught inside onStatus, not "
|
||||
+ "escape it");
|
||||
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"the delivery's future must surface the send failure");
|
||||
assertEquals(Injector.Cancellation.NOT_DELIVERED, inj.cancel(delivery),
|
||||
"the message must be dropped and marked NOT_DELIVERED, not left QUEUED at the "
|
||||
+ "head of the queue");
|
||||
}
|
||||
|
||||
@Test
|
||||
void anErrorFromSendDoesNotRedeliverOnASecondRound() {
|
||||
// fleetd #546, acceptance test 2 — the test this ticket exists for. Before the fix, an
|
||||
// Error at the send seam left the message QUEUED (Injector.java:378 peeks, not polls), so a
|
||||
// second onStatus round re-entered the same try and sent the same text again: the member's
|
||||
// pane got the same brief typed into it twice.
|
||||
ErrorOnPrompt throwing = new ErrorOnPrompt(new FakeHerdr());
|
||||
Injector inj = new Injector(new AgentControl(throwing));
|
||||
inj.enqueue(T, "brief", TestTurnTokens.inert(T));
|
||||
|
||||
pollOnceSurviving(inj, T, AgentStatus.IDLE);
|
||||
pollOnceSurviving(inj, T, AgentStatus.IDLE);
|
||||
|
||||
long promptCalls = throwing.calls().stream().filter(c -> c.method().equals("agent.prompt")).count();
|
||||
assertEquals(1, promptCalls, "fleetd #546: the poisoned text must be sent exactly once — a "
|
||||
+ "second onStatus round must not re-enter send for the same message");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aHerdrExceptionFromSendStillProducesNotDeliveredUnchanged() {
|
||||
// fleetd #546, acceptance test 3 (control): HerdrException extends RuntimeException, so it
|
||||
// was already caught before this ticket's widening. This pins that the ordinary path is
|
||||
// unchanged — still dropped, still NOT_DELIVERED, still surfaced to the caller.
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
Injector.Delivery delivery = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertTrue(delivery.completion().isCompletedExceptionally(),
|
||||
"a HerdrException at the send seam must still surface to the caller, unchanged by "
|
||||
+ "fleetd #546's widening");
|
||||
assertEquals(Injector.Cancellation.NOT_DELIVERED, inj.cancel(delivery),
|
||||
"a HerdrException must still be dropped and marked NOT_DELIVERED, unchanged by the "
|
||||
+ "wider Throwable catch");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.TestTurnTokens;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotSame;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class StatusPollerResilienceTest {
|
||||
|
||||
@Test
|
||||
void anErrorForOneTargetDoesNotStopPollingTheNextTarget() throws Exception {
|
||||
FakeHerdr fake = new FakeHerdr().withAgent("worker", "term_b", "w2:p8", "w2:t8");
|
||||
CountDownLatch errorThrown = new CountDownLatch(1);
|
||||
AgentControl agents = new AgentControl(new ErrorOnceForFirstTarget(fake, errorThrown));
|
||||
Injector injector = new Injector(agents);
|
||||
StatusPoller poller = new StatusPoller(agents, injector, 1);
|
||||
poller.start();
|
||||
try {
|
||||
injector.enqueue("term_a", "first", TestTurnTokens.inert("term_a"));
|
||||
assertTrue(errorThrown.await(2, TimeUnit.SECONDS),
|
||||
"the first target must throw its test Error");
|
||||
|
||||
CompletableFuture<Void> delivered =
|
||||
injector.enqueue("term_b", "second", TestTurnTokens.inert("term_b")).completion();
|
||||
delivered.get(2, TimeUnit.SECONDS);
|
||||
} finally {
|
||||
poller.stop();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anAbnormalExitClearsRunningSoStartCreatesANewLoop() throws Exception {
|
||||
StatusPoller poller = new StatusPoller(new AgentControl(new FakeHerdr()), new Injector(new AgentControl(new FakeHerdr())), -1);
|
||||
poller.start();
|
||||
Thread first = threadOf(poller);
|
||||
first.join(2000);
|
||||
assertFalse(runningOf(poller), "an abnormal loop exit must clear running");
|
||||
|
||||
poller.start();
|
||||
Thread restarted = threadOf(poller);
|
||||
try {
|
||||
assertNotSame(first, restarted, "start() must create a new loop after an abnormal exit");
|
||||
restarted.join(2000);
|
||||
} finally {
|
||||
poller.stop();
|
||||
}
|
||||
}
|
||||
|
||||
private static Thread threadOf(StatusPoller poller) throws ReflectiveOperationException {
|
||||
Field field = StatusPoller.class.getDeclaredField("thread");
|
||||
field.setAccessible(true);
|
||||
return (Thread) field.get(poller);
|
||||
}
|
||||
|
||||
private static boolean runningOf(StatusPoller poller) throws ReflectiveOperationException {
|
||||
Field field = StatusPoller.class.getDeclaredField("running");
|
||||
field.setAccessible(true);
|
||||
return field.getBoolean(poller);
|
||||
}
|
||||
|
||||
private static final class ErrorOnceForFirstTarget implements HerdrClient {
|
||||
private final FakeHerdr delegate;
|
||||
private final CountDownLatch errorThrown;
|
||||
private final AtomicBoolean first = new AtomicBoolean(true);
|
||||
|
||||
private ErrorOnceForFirstTarget(FakeHerdr delegate, CountDownLatch errorThrown) {
|
||||
this.delegate = delegate;
|
||||
this.errorThrown = errorThrown;
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if (method.equals("agent.get") && params instanceof Map<?, ?> map
|
||||
&& "w2:p7".equals(map.get("target")) && first.compareAndSet(true, false)) {
|
||||
errorThrown.countDown();
|
||||
throw new AssertionError("test Error from the first poll target");
|
||||
}
|
||||
return delegate.call(method, params);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
delegate.close();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,23 +1,22 @@
|
||||
package dev.ltms.fleet.lead;
|
||||
|
||||
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 com.fasterxml.jackson.databind.JsonNode;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.HerdrException;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
@@ -862,25 +861,16 @@ class LeadRolloverTest {
|
||||
|
||||
// ---- fleetd #494: the log lines must print MEASURED values, never the configured budget --
|
||||
|
||||
private static ListAppender<ILoggingEvent> attachLog() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(LeadRollover.class);
|
||||
logger.setLevel(Level.DEBUG);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
private static CapturedLog attachLog() {
|
||||
return CapturedLog.at(LeadRollover.class, Level.DEBUG);
|
||||
}
|
||||
|
||||
private static void detachLog(ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(LeadRollover.class)).detachAppender(appender);
|
||||
}
|
||||
|
||||
private static ILoggingEvent lastEventContaining(ListAppender<ILoggingEvent> events, String substring) {
|
||||
return events.list.stream()
|
||||
private static ILoggingEvent lastEventContaining(List<ILoggingEvent> events, String substring) {
|
||||
return events.stream()
|
||||
.filter(e -> e.getFormattedMessage().contains(substring))
|
||||
.reduce((_, b) -> b)
|
||||
.orElseThrow(() -> new AssertionError("no log event contained \"" + substring
|
||||
+ "\"; got: " + events.list.stream().map(ILoggingEvent::getFormattedMessage).toList()));
|
||||
+ "\"; got: " + events.stream().map(ILoggingEvent::getFormattedMessage).toList()));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -914,15 +904,14 @@ class LeadRolloverTest {
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(flipsAfterClear, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
try (CapturedLog log = attachLog()) {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
assertTrue(decision.accepted(), "every synchronous gate passes; the refusal is logged "
|
||||
+ "only, deep inside the deferred continuation");
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "NOT sending bootstrapText");
|
||||
ILoggingEvent event = lastEventContaining(log.events(), "NOT sending bootstrapText");
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
assertTrue(message.contains("configured=1s"), "must label the configured budget: " + message);
|
||||
@@ -934,8 +923,6 @@ class LeadRolloverTest {
|
||||
+ message);
|
||||
assertFalse(message.contains("within 1s"), "must not present the configured budget as if "
|
||||
+ "it were the measured wait duration: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -954,15 +941,14 @@ class LeadRolloverTest {
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
try (CapturedLog log = attachLog()) {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason()
|
||||
+ " / " + decision.detail());
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "releasing rather than wedging the roll");
|
||||
ILoggingEvent event = lastEventContaining(log.events(), "releasing rather than wedging the roll");
|
||||
assertEquals(Level.WARN, event.getLevel(), "the grace-limit release must be WARN, not "
|
||||
+ "INFO — it is exactly the case that reported false success in the real incident "
|
||||
+ "this fix comes from (a roll that 'succeeded' after 438ms of a 20s budget)");
|
||||
@@ -982,8 +968,6 @@ class LeadRolloverTest {
|
||||
assertTrue(message.contains("after 8 consecutive IDLE/DONE polls (7 of those were nudged)"),
|
||||
"must print the measured poll count and nudge count as plain numbers, not the "
|
||||
+ "PICKUP_GRACE_POLLS constant standing in for either: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1000,15 +984,14 @@ class LeadRolloverTest {
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
try (CapturedLog log = attachLog()) {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
assertTrue(decision.accepted(), "expected approval; got: " + decision.reason()
|
||||
+ " / " + decision.detail());
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "lead-rollover: rolled");
|
||||
ILoggingEvent event = lastEventContaining(log.events(), "lead-rollover: rolled");
|
||||
assertEquals(Level.INFO, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
// fleetd #494 follow-up: waitUntilAtTurnBoundary now also reads the injected clock one
|
||||
@@ -1018,8 +1001,6 @@ class LeadRolloverTest {
|
||||
assertTrue(message.contains("elapsedMs=7000"), "must print the MEASURED elapsed time for "
|
||||
+ "the whole roll — with this fixture's advancing clock, the full roll (turn-settle "
|
||||
+ "wait + /clear wait + bootstrapText) took 7000ms: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1040,8 +1021,7 @@ class LeadRolloverTest {
|
||||
AtomicLong clock = new AtomicLong(1_000);
|
||||
LeadRollover rollover = newRollover(herdr, config, () -> clock.addAndGet(500));
|
||||
|
||||
ListAppender<ILoggingEvent> events = attachLog();
|
||||
try {
|
||||
try (CapturedLog log = attachLog()) {
|
||||
LeadRollover.PendingRollover pending = rollover.open(LEAD, "context is full");
|
||||
LeadRollover.RollDecision decision = rollover.confirm(LEAD, pending.token(), true);
|
||||
|
||||
@@ -1049,7 +1029,7 @@ class LeadRolloverTest {
|
||||
+ "only inside the deferred continuation, which this test's synchronous runner "
|
||||
+ "has already run to completion by the time confirm() returns");
|
||||
|
||||
ILoggingEvent event = lastEventContaining(events, "refusing to send /clear at all");
|
||||
ILoggingEvent event = lastEventContaining(log.events(), "refusing to send /clear at all");
|
||||
assertEquals(Level.WARN, event.getLevel());
|
||||
String message = event.getFormattedMessage();
|
||||
assertTrue(message.contains("configured=1s"), "must label the configured budget: " + message);
|
||||
@@ -1058,8 +1038,6 @@ class LeadRolloverTest {
|
||||
+ "1s(=1000ms) configured budget: " + message);
|
||||
assertFalse(message.contains("within 1s"), "must not present the configured budget as if "
|
||||
+ "it were the measured wait duration: " + message);
|
||||
} finally {
|
||||
detachLog(events);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -57,6 +57,23 @@ class ConnectionIdentityTest {
|
||||
// must read as "resolved" — the distinction #317 turns on.
|
||||
ConnectionIdentity.Caller c = with(_ -> 999_999).resolve("127.0.0.1", 55555);
|
||||
assertTrue(c.resolved());
|
||||
assertTrue(c.scanComplete(), "no herdr error happened, so the scan is complete");
|
||||
}
|
||||
|
||||
@Test
|
||||
void scanIsIncompleteWhenHerdrErrorsOnThePaneThatOwnsThePid() {
|
||||
// fleetd #505: a transient herdr error on exactly the pane that DOES own the caller's pid
|
||||
// must be visible as an incomplete scan, distinct from a real primary (resolved(), null
|
||||
// terminal, complete scan). Both have pid > 0 and a null terminal — scanComplete is the
|
||||
// only thing that tells them apart.
|
||||
FakeHerdr failing = new FakeHerdr().processInfoFailsForPane("w2:p7", "transient");
|
||||
ConnectionIdentity id = new ConnectionIdentity(new PaneLocator(failing), _ -> FakeHerdr.WORKER_PID);
|
||||
|
||||
ConnectionIdentity.Caller c = id.resolve("127.0.0.1", 55555);
|
||||
|
||||
assertTrue(c.resolved(), "the pid itself resolved fine — this is not #317's failure");
|
||||
assertNull(c.terminal(), "the owning pane could not be confirmed");
|
||||
assertFalse(c.scanComplete(), "the scan could not check the pane that owns this pid");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -78,10 +78,15 @@ class FleetMcpAuthzTest {
|
||||
// fleetd #480 correction round: FleetMcp has one constructor now (no defaulting
|
||||
// overloads — see its javadoc), so every feature this test does not exercise is passed
|
||||
// its explicit "off" value here rather than being omitted.
|
||||
//
|
||||
// fleetd #518: callers is now required (never null) either way — the resolver that used
|
||||
// to be omitted to reach "legacy" is now always real, and AuthorizationMode is the
|
||||
// separate, explicit choice that governs enforcement.
|
||||
mcp = new FleetMcp(messages, workers, sessions, identity, sessions.asPresence(),
|
||||
new PrimaryRegistry(null),
|
||||
enforce ? CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
Map::of, new MemberRegistry(null)) : null,
|
||||
CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
Map::of, new MemberRegistry(null)),
|
||||
enforce ? FleetMcp.AuthorizationMode.ENFORCED : FleetMcp.AuthorizationMode.UNENFORCED,
|
||||
metrics, FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), null, FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LeadSeatSource.none(), List.of(), null);
|
||||
@@ -187,7 +192,31 @@ class FleetMcpAuthzTest {
|
||||
// The 22 pre-existing FleetMcpTest cases rely on no authorization being enforced.
|
||||
FleetMcp m = mcp(false);
|
||||
assertNull(m.denyFor(ANON, Authz.Action.SPAWN, null),
|
||||
"no CallerResolver supplied ⇒ authorization not enforced (legacy behaviour)");
|
||||
"AuthorizationMode.UNENFORCED chosen explicitly ⇒ authorization not enforced "
|
||||
+ "(legacy behaviour) — fleetd #518 replaced the old callers == null idiom");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #509 was originally proven against {@code FleetMcp.legacyPrincipal} — a second,
|
||||
* separately-maintained principal-resolution heuristic that only ran when {@code callers} was
|
||||
* omitted (null). fleetd #518 deleted that whole heuristic: {@code callers} is now required
|
||||
* and non-null under every {@link FleetMcp.AuthorizationMode}, so the ONE real
|
||||
* {@link CallerResolver} resolves every caller, enforced or not, and #509's property (a
|
||||
* non-loopback / unresolved caller must never earn the primary's authority) is exactly what
|
||||
* {@code CallerResolverTest.aNonLoopbackCallerIsNeverThePrimaryUnderLoopbackTrust} already
|
||||
* proves on that one real path. There is no longer a second heuristic here to test.
|
||||
*/
|
||||
@Test
|
||||
void anUnresolvedNonLoopbackCallerIsAnonymousUnderTheOneRealResolver() {
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> 999_999);
|
||||
CallerResolver resolver = CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
Map::of, new MemberRegistry(null));
|
||||
// A non-loopback address never even reaches the pane scan — resolve() short-circuits it
|
||||
// to Caller(null, -1, true), the same "no terminal" shape a genuine primary's connection
|
||||
// produces on loopback. The real resolver must not conflate the two.
|
||||
Principal p = resolver.resolve("8.8.8.8", 1234, null);
|
||||
assertEquals(Principal.anonymous(), p,
|
||||
"an unresolved, non-loopback caller must earn no authority, not the primary's");
|
||||
}
|
||||
|
||||
// --- fleetd #439: who may see fleet_list's coordinator row ----------------------------------
|
||||
|
||||
@@ -0,0 +1,147 @@
|
||||
package dev.ltms.fleet.mcp;
|
||||
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.session.FakeWorktrees;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import io.modelcontextprotocol.client.McpClient;
|
||||
import io.modelcontextprotocol.client.McpSyncClient;
|
||||
import io.modelcontextprotocol.client.transport.HttpClientStreamableHttpTransport;
|
||||
import io.modelcontextprotocol.spec.McpClientTransport;
|
||||
import io.modelcontextprotocol.spec.McpSchema;
|
||||
import org.eclipse.jetty.server.Server;
|
||||
import org.eclipse.jetty.server.ServerConnector;
|
||||
import org.eclipse.jetty.servlet.ServletContextHandler;
|
||||
import org.eclipse.jetty.servlet.ServletHolder;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.net.http.HttpRequest;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #518 — Part 2: drive the {@code contextExtractor} closure for real.
|
||||
*
|
||||
* <p>{@code FleetMcp.deny()}/{@code denyFor()} has a full policy table of tests
|
||||
* ({@code FleetMcpAuthzTest}), and {@code CallerResolver.resolve()} has its own full suite
|
||||
* ({@code CallerResolverTest}). Neither one ever exercises the closure that WIRES them together
|
||||
* inside {@code FleetMcp}'s constructor: it is built once, handed to the MCP SDK's transport, and
|
||||
* only ever runs when a real MCP client makes a real HTTP request. Every existing test either
|
||||
* calls {@code denyFor(Principal, ...)} with a hand-built {@link dev.ltms.fleet.auth.Principal}
|
||||
* (never asking who the transport would actually have resolved) or drives a static handler method
|
||||
* directly. A mutation that swapped the whole resolution decision for an unconditional fallback —
|
||||
* bypassing {@link CallerResolver} entirely — passed the full suite, including every
|
||||
* {@code FleetMcpAuthzTest} case, because none of them go through the transport at all.
|
||||
*
|
||||
* <p>This test boots the real {@code HttpServletStreamableServerTransportProvider} on a real
|
||||
* Jetty server, drives it with a real MCP client over HTTP, and checks a result that only the
|
||||
* real {@link CallerResolver} can produce: token-mode inspects the {@code Authorization} header
|
||||
* and grants {@code PRIMARY} only for the right bearer token. The connection never resolves to a
|
||||
* worker pane (the fake peer-pid lookup always misses), so the ONLY way {@code fleet_whoami} can
|
||||
* come back as {@code primary} is if the closure actually called {@code callers.resolve(...)} and
|
||||
* read that header — a behaviour the deleted {@code legacyPrincipal} heuristic never had at all.
|
||||
*/
|
||||
class FleetMcpContextExtractorTest {
|
||||
|
||||
private static final String TOKEN = "s3cret-mcp-token";
|
||||
|
||||
private final FakeHerdr herdr = new FakeHerdr();
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
private FleetMcp mcp;
|
||||
private Server server;
|
||||
|
||||
@AfterEach
|
||||
void tearDown() throws Exception {
|
||||
if (server != null) {
|
||||
server.stop();
|
||||
}
|
||||
if (mcp != null) {
|
||||
mcp.close();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aRealMcpRequestIsResolvedByTheRealCallerResolverNotAFallback() throws Exception {
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN", null,
|
||||
"tab", "fleetd-workers", "worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(agents, new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
_ -> "tok");
|
||||
SessionManager sessions = new SessionManager(workers, new FakeWorktrees());
|
||||
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(),
|
||||
new InMemoryReplyInbox());
|
||||
// The peer-pid lookup always misses (-1), so no connection here is ever resolved to a
|
||||
// worker pane — every call falls through to CallerResolver's token check, the one branch
|
||||
// that is unreachable through the deleted legacy heuristic.
|
||||
ConnectionIdentity identity = new ConnectionIdentity(new PaneLocator(herdr), _ -> -1);
|
||||
CallerResolver callers = CallerResolver.withLeadsAndMembers(identity, true, TOKEN,
|
||||
Map::of, new MemberRegistry(null));
|
||||
|
||||
mcp = new FleetMcp(messages, workers, sessions, identity, sessions.asPresence(),
|
||||
new PrimaryRegistry(null), callers, FleetMcp.AuthorizationMode.ENFORCED,
|
||||
null, FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), null, FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LeadSeatSource.none(), List.of(), null);
|
||||
|
||||
ServletContextHandler handler = new ServletContextHandler();
|
||||
handler.setContextPath("/");
|
||||
handler.addServlet(new ServletHolder(mcp.servlet()), "/mcp");
|
||||
server = new Server(0);
|
||||
server.setHandler(handler);
|
||||
server.start();
|
||||
String baseUrl = "http://127.0.0.1:"
|
||||
+ ((ServerConnector) server.getConnectors()[0]).getLocalPort();
|
||||
|
||||
// The right bearer token: the real CallerResolver grants PRIMARY, which fleet_whoami's
|
||||
// READ gate lets through.
|
||||
McpSchema.CallToolResult authorized = callWhoami(baseUrl, "Bearer " + TOKEN);
|
||||
assertFalse(authorized.isError(), "a valid bearer token must resolve as PRIMARY and pass "
|
||||
+ "fleet_whoami's READ gate: " + textOf(authorized));
|
||||
assertTrue(textOf(authorized).contains("\"role\":\"primary\""),
|
||||
"fleet_whoami must report the role the real CallerResolver resolved over this "
|
||||
+ "connection, not a fallback: " + textOf(authorized));
|
||||
|
||||
// No credential at all, over the SAME wiring: the real resolver refuses it as ANONYMOUS.
|
||||
// legacyPrincipal never looked at the Authorization header, so it could not have told
|
||||
// these two calls apart at all -- this is the assertion the deleted mutation would fail.
|
||||
McpSchema.CallToolResult unauthorized = callWhoami(baseUrl, null);
|
||||
assertTrue(unauthorized.isError(), "no credential must be refused, not silently let "
|
||||
+ "through: " + textOf(unauthorized));
|
||||
}
|
||||
|
||||
private static McpSchema.CallToolResult callWhoami(String baseUrl, String authorizationHeader) {
|
||||
HttpRequest.Builder requestTemplate = HttpRequest.newBuilder();
|
||||
if (authorizationHeader != null) {
|
||||
requestTemplate.header("Authorization", authorizationHeader);
|
||||
}
|
||||
McpClientTransport transport = HttpClientStreamableHttpTransport.builder(baseUrl)
|
||||
.endpoint("/mcp")
|
||||
.requestBuilder(requestTemplate)
|
||||
.build();
|
||||
try (McpSyncClient client = McpClient.sync(transport).build()) {
|
||||
client.initialize();
|
||||
return client.callTool(McpSchema.CallToolRequest.builder("fleet_whoami").arguments(Map.of()).build());
|
||||
}
|
||||
}
|
||||
|
||||
private static String textOf(McpSchema.CallToolResult r) {
|
||||
return ((McpSchema.TextContent) r.content().getFirst()).text();
|
||||
}
|
||||
}
|
||||
@@ -89,6 +89,7 @@ class FleetMcpHandoverTest {
|
||||
mcp = new FleetMcp(messages, workers, sessions, identity, sessions.asPresence(),
|
||||
new PrimaryRegistry(null),
|
||||
CallerResolver.withLeadsAndMembers(identity, false, null, Map::of, new MemberRegistry(null)),
|
||||
FleetMcp.AuthorizationMode.ENFORCED,
|
||||
null, FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), null, FleetMcp.OutageSource.none(),
|
||||
FleetMcp.LeadSeatSource.none(), List.of(), leadRollover);
|
||||
|
||||
@@ -1,17 +1,17 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
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 com.rabbitmq.client.Channel;
|
||||
import com.rabbitmq.client.ConnectionFactory;
|
||||
import com.rabbitmq.client.impl.DefaultExceptionHandler;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
@@ -32,21 +32,17 @@ class AmqpConnectionFailureLoggerTest {
|
||||
assertEquals(AmqpConnectionFailureLogger.REPLY_INBOX, inboxHandler.connectionName());
|
||||
assertEquals(AmqpConnectionFailureLogger.LEAD_MAILBOX, mailboxHandler.connectionName());
|
||||
|
||||
ListAppender<ILoggingEvent> inboxEvents = attach(AmqpReplyInbox.class);
|
||||
ListAppender<ILoggingEvent> mailboxEvents = attach(LeadMailbox.class);
|
||||
IllegalStateException inboxFailure = new IllegalStateException("inbox failure");
|
||||
IllegalStateException mailboxFailure = new IllegalStateException("mailbox failure");
|
||||
try {
|
||||
try (CapturedLog inboxLog = attach(AmqpReplyInbox.class);
|
||||
CapturedLog mailboxLog = attach(LeadMailbox.class)) {
|
||||
inboxHandler.handleUnexpectedConnectionDriverException(null, inboxFailure);
|
||||
mailboxHandler.handleConnectionRecoveryException(null, mailboxFailure);
|
||||
|
||||
assertError(inboxEvents, "AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred",
|
||||
assertError(inboxLog.events(), "AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred",
|
||||
inboxFailure, "inbox failure line");
|
||||
assertError(mailboxEvents, "AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!",
|
||||
assertError(mailboxLog.events(), "AMQP connection fleetd-lead-mailbox: Caught an exception during connection recovery!",
|
||||
mailboxFailure, "mailbox recovery line");
|
||||
} finally {
|
||||
detach(AmqpReplyInbox.class, inboxEvents);
|
||||
detach(LeadMailbox.class, mailboxEvents);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -54,17 +50,14 @@ class AmqpConnectionFailureLoggerTest {
|
||||
void connectionResetKeepsForgivingHandlerWarningSemantics() {
|
||||
AmqpConnectionFailureLogger handler = new AmqpConnectionFailureLogger(
|
||||
AmqpConnectionFailureLogger.REPLY_INBOX, LoggerFactory.getLogger(AmqpReplyInbox.class));
|
||||
ListAppender<ILoggingEvent> events = attach(AmqpReplyInbox.class);
|
||||
try {
|
||||
try (CapturedLog log = attach(AmqpReplyInbox.class)) {
|
||||
handler.handleUnexpectedConnectionDriverException(null, new IOException("Connection reset"));
|
||||
assertEquals(1, events.list.size(), "the handler must still log a reset");
|
||||
ILoggingEvent event = events.list.getFirst();
|
||||
assertEquals(1, log.events().size(), "the handler must still log a reset");
|
||||
ILoggingEvent event = log.events().getFirst();
|
||||
assertEquals(Level.WARN, event.getLevel(), "ForgivingExceptionHandler logs connection resets at WARN");
|
||||
assertEquals("AMQP connection fleetd-reply-inbox: An unexpected connection driver error occurred "
|
||||
+ "(Exception message: Connection reset)", event.getFormattedMessage());
|
||||
assertTrue(event.getThrowableProxy() == null, "ForgivingExceptionHandler does not attach a reset stack trace");
|
||||
} finally {
|
||||
detach(AmqpReplyInbox.class, events);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -103,17 +96,8 @@ class AmqpConnectionFailureLoggerTest {
|
||||
"all exception-handling methods must remain inherited from DefaultExceptionHandler");
|
||||
}
|
||||
|
||||
private static ListAppender<ILoggingEvent> attach(Class<?> owner) {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(owner);
|
||||
logger.setLevel(Level.DEBUG);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detach(Class<?> owner, ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(owner)).detachAppender(appender);
|
||||
private static CapturedLog attach(Class<?> owner) {
|
||||
return CapturedLog.at(owner, Level.DEBUG);
|
||||
}
|
||||
|
||||
private static AmqpConnectionFailureLogger installedStrictHandler(ConnectionFactory factory, String connection) {
|
||||
@@ -123,9 +107,9 @@ class AmqpConnectionFailureLoggerTest {
|
||||
return assertInstanceOf(AmqpConnectionFailureLogger.class, factory.getExceptionHandler());
|
||||
}
|
||||
|
||||
private static void assertError(ListAppender<ILoggingEvent> events, String message, Throwable cause, String name) {
|
||||
assertEquals(1, events.list.size(), name);
|
||||
ILoggingEvent event = events.list.getFirst();
|
||||
private static void assertError(List<ILoggingEvent> events, String message, Throwable cause, String name) {
|
||||
assertEquals(1, events.size(), name);
|
||||
ILoggingEvent event = events.getFirst();
|
||||
assertEquals(Level.ERROR, event.getLevel(), name);
|
||||
assertEquals(message, event.getFormattedMessage(), name);
|
||||
assertEquals(cause.toString(), event.getThrowableProxy().getClassName() + ": "
|
||||
|
||||
@@ -1,16 +1,13 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.IThrowableProxy;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
@@ -489,36 +486,33 @@ class GitWorktreesTest {
|
||||
// ---- CB-189: broader remote-URL coverage — every remote, both fetch and push URLs, any
|
||||
// non-SSH scheme. Reporting only, additive to the origin/https strip-and-refuse tests above. ----
|
||||
|
||||
private Logger reportingLogger;
|
||||
private ListAppender<ILoggingEvent> reportingAppender;
|
||||
private CapturedLog reportingLog;
|
||||
|
||||
/** {@link GitWorktrees}'s own logger, captured fresh for each test so assertions never see a
|
||||
* message left over from a previous test. */
|
||||
* message left over from a previous test. fleetd #529: {@link CapturedLog#close} restores the
|
||||
* level it captured here (whatever it truly was before this test, not just WARN), so a test
|
||||
* below that further lowers the level to INFO for its own assertion (via {@link
|
||||
* #reportingLog}'s {@link CapturedLog#setLevel}) can never leak that INFO pin past its own
|
||||
* {@code @AfterEach} — every test's window is self-contained. */
|
||||
@BeforeEach
|
||||
void attachReportingLogCapture() {
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
reportingLogger = ctx.getLogger(GitWorktrees.class);
|
||||
reportingLogger.setLevel(Level.WARN);
|
||||
reportingAppender = new ListAppender<>();
|
||||
reportingAppender.setContext(ctx);
|
||||
reportingAppender.start();
|
||||
reportingLogger.addAppender(reportingAppender);
|
||||
reportingLog = CapturedLog.at(GitWorktrees.class, Level.WARN);
|
||||
}
|
||||
|
||||
@AfterEach
|
||||
void detachReportingLogCapture() {
|
||||
reportingLogger.detachAppender(reportingAppender);
|
||||
reportingLog.close();
|
||||
}
|
||||
|
||||
private List<String> capturedMessages() {
|
||||
return reportingAppender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
return reportingLog.events().stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
}
|
||||
|
||||
/** Asserts {@code secret} appears in no captured message, and in no attached exception's
|
||||
* message either — the constraint is that a credential must never reach a log, however it
|
||||
* would have gotten there. */
|
||||
private void assertNoLeak(String secret) {
|
||||
for (ILoggingEvent event : reportingAppender.list) {
|
||||
for (ILoggingEvent event : reportingLog.events()) {
|
||||
assertFalse(event.getFormattedMessage().contains(secret),
|
||||
"log message leaked a credential (" + secret + "): " + event.getFormattedMessage());
|
||||
IThrowableProxy thrown = event.getThrowableProxy();
|
||||
@@ -596,7 +590,7 @@ class GitWorktreesTest {
|
||||
|
||||
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-189-d", "HEAD");
|
||||
|
||||
assertTrue(reportingAppender.list.isEmpty(),
|
||||
assertTrue(reportingLog.events().isEmpty(),
|
||||
"expected no report for an ssh remote and a credential-free https remote, got:\n"
|
||||
+ capturedMessages());
|
||||
}
|
||||
@@ -812,7 +806,7 @@ class GitWorktreesTest {
|
||||
/** Criterion 1, all three present: the summary names the denominator and every neutralized file. */
|
||||
@Test
|
||||
void isolateToolSurfaceLogsAllThreeConfigsNeutralized(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepoWithAllThreeConfigs(tmp.resolve("repo"));
|
||||
|
||||
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-134-log-all", "HEAD");
|
||||
@@ -827,7 +821,7 @@ class GitWorktreesTest {
|
||||
/** Criterion 1, two absent: the summary must still name the denominator and say why. */
|
||||
@Test
|
||||
void isolateToolSurfaceLogsAbsentConfigsWithReason(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepo(tmp.resolve("repo")); // only .mcp.json + README committed
|
||||
|
||||
new GitWorktrees(tmp.resolve("wts").toString()).add(repo.toString(), "cb-134-log-partial", "HEAD");
|
||||
@@ -1435,7 +1429,7 @@ class GitWorktreesTest {
|
||||
/** Criterion 2: both candidates present — the summary line names both and the denominator. */
|
||||
@Test
|
||||
void overlayParityLogsBothCopiedWhenBothCandidatesArePresent(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Files.writeString(repo.resolve(".env"), "A=1\n");
|
||||
Files.writeString(repo.resolve(".envrc"), "export A=1\n");
|
||||
@@ -1452,7 +1446,7 @@ class GitWorktreesTest {
|
||||
* denominator, and why the other candidate was not copied. */
|
||||
@Test
|
||||
void overlayParityLogsOneCopiedOneAbsent(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Files.writeString(repo.resolve(".env"), "A=1\n");
|
||||
// .envrc deliberately not created — the absent candidate.
|
||||
@@ -1472,7 +1466,7 @@ class GitWorktreesTest {
|
||||
* change, silently, unless this line told it so beforehand. */
|
||||
@Test
|
||||
void overlayParityLogsSkipWorktreeConsequenceForATrackedFile(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Files.writeString(repo.resolve(".env"), "A=1\n");
|
||||
git(repo, "add", ".env");
|
||||
@@ -1495,7 +1489,7 @@ class GitWorktreesTest {
|
||||
/** Criterion 5: null and empty overlay lists return quietly — no exception, no log noise. */
|
||||
@Test
|
||||
void overlayParityWithNoCandidatesLogsNothing(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
Path wt = bareWorktree(repo, tmp.resolve("wt"), "cb134-empty");
|
||||
GitWorktrees worktrees = new GitWorktrees(tmp.resolve("wts").toString());
|
||||
@@ -1503,7 +1497,7 @@ class GitWorktreesTest {
|
||||
worktrees.overlayParity(repo.toString(), wt.toString(), null);
|
||||
worktrees.overlayParity(repo.toString(), wt.toString(), List.of());
|
||||
|
||||
assertTrue(reportingAppender.list.isEmpty(),
|
||||
assertTrue(reportingLog.events().isEmpty(),
|
||||
"a null/empty overlay must log nothing, got:\n" + capturedMessages());
|
||||
}
|
||||
|
||||
@@ -1683,7 +1677,7 @@ class GitWorktreesTest {
|
||||
* denominator, what was seeded, and what was kept because the repo already had it. */
|
||||
@Test
|
||||
void seedSkillsLogsSeededAndKept(@TempDir Path tmp) throws Exception {
|
||||
reportingLogger.setLevel(Level.INFO);
|
||||
reportingLog.setLevel(Level.INFO);
|
||||
Path repo = tmp.resolve("repo");
|
||||
Files.createDirectories(repo);
|
||||
git(repo, "init", "-q", "-b", "main");
|
||||
|
||||
@@ -1,9 +1,7 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.auth.MemberLifecycle;
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
@@ -25,7 +23,13 @@ import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.AfterAll;
|
||||
import org.junit.jupiter.api.BeforeAll;
|
||||
import org.junit.jupiter.api.MethodOrderer;
|
||||
import org.junit.jupiter.api.Order;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.TestMethodOrder;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -50,9 +54,46 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
* CB-301 / CB-303 acceptance tests for the authoritative session registry, one-shot lifecycle FSM,
|
||||
* and configurable lifecycle limits (idle TTL, context cap, drain).
|
||||
* No live herdr — everything runs against the same {@link FakeHerdr} the rest of the project uses.
|
||||
*
|
||||
* <p>fleetd #525: only {@link #onTurnFailedIsLoggedAtWarnWithThePriorState} (explicitly
|
||||
* {@link Order#value() @Order(1)}) and the proving test right after it
|
||||
* ({@link #sharedSessionManagerLoggerLevelIsRestoredAfterOnTurnFailedPinsWarn}, {@code @Order(2)})
|
||||
* care about method order — every other test here has no {@code @Order} and so runs after both of
|
||||
* these (JUnit 5's {@link MethodOrderer.OrderAnnotation} gives an unannotated method the lowest
|
||||
* priority), in whatever relative order it already ran in.
|
||||
*/
|
||||
@TestMethodOrder(MethodOrderer.OrderAnnotation.class)
|
||||
class SessionManagerTest {
|
||||
|
||||
/**
|
||||
* fleetd #525: the level {@link SessionManager}'s logger had when this class started, captured
|
||||
* before any test here — including the leak this ticket fixes — can touch it. {@code
|
||||
* pinSessionManagerLoggerToAKnownBaseline} then forces a distinctive, known value (DEBUG) so
|
||||
* {@link #sharedSessionManagerLoggerLevelIsRestoredAfterOnTurnFailedPinsWarn} can tell "the
|
||||
* level came back to what it was" apart from "the level happens to already be WARN because
|
||||
* some earlier test class in this JVM fork (surefire reuses forks by default) left it there" —
|
||||
* a real risk, since {@code ch.qos.logback.classic.Logger} instances are cached per class and
|
||||
* shared across the whole JVM, and this exact logger is also touched by
|
||||
* {@code WorktreeSessionManagerTest#releasePreservesDirtyWorktreeAndLogsWarn}, which has the
|
||||
* same unfixed leak (reported, not fixed — out of this ticket's scope).
|
||||
*/
|
||||
private static Level sessionManagerLevelBeforeThisClass;
|
||||
|
||||
@BeforeAll
|
||||
static void pinSessionManagerLoggerToAKnownBaseline() {
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
sessionManagerLevelBeforeThisClass = sessionLog.getLevel();
|
||||
sessionLog.setLevel(Level.DEBUG);
|
||||
}
|
||||
|
||||
@AfterAll
|
||||
static void restoreSessionManagerLoggerLevel() {
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
sessionLog.setLevel(sessionManagerLevelBeforeThisClass);
|
||||
}
|
||||
|
||||
private SessionManager sessionManager(FakeHerdr herdr) {
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
@@ -81,6 +122,9 @@ class SessionManagerTest {
|
||||
return new SessionManager(workers, worktrees, clock);
|
||||
}
|
||||
|
||||
// fleetd #529: CapturedLog moved to dev.ltms.fleet.testing.CapturedLog (imported above) so
|
||||
// every test file shares one implementation instead of each hand-rolling its own capture.
|
||||
|
||||
/**
|
||||
* CB-581: a {@link Worktrees} test double whose {@code hasUncommitted} and {@code remove} can
|
||||
* be told to throw, so {@link SessionManager#release} can be exercised against exactly the
|
||||
@@ -466,24 +510,15 @@ class SessionManagerTest {
|
||||
|
||||
@Test
|
||||
void backendErrorForUnknownTargetIsWarnedAndDoesNotCreateASession() {
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.of(SessionManager.class)) {
|
||||
SessionManager sessions = sessionManager(new FakeHerdr());
|
||||
|
||||
assertFalse(sessions.onBackendError("term_missing", "backend exited"));
|
||||
|
||||
assertTrue(sessions.roster().isEmpty(), "unknown target must not create a session");
|
||||
assertTrue(appender.list.stream().anyMatch(e -> e.getLevel().equals(Level.WARN)
|
||||
assertTrue(log.events().stream().anyMatch(e -> e.getLevel().equals(Level.WARN)
|
||||
&& e.getFormattedMessage().contains("term_missing")),
|
||||
"unknown target is logged at WARN");
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -516,19 +551,12 @@ class SessionManagerTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
@Order(1)
|
||||
void onTurnFailedIsLoggedAtWarnWithThePriorState() {
|
||||
// CB-564: this transition used to be a bare DEBUG "session marked failed" — a symptom with no
|
||||
// cause. A member that can no longer be delegated to must be at least WARN, and should name
|
||||
// what stage it failed at (here: BUSY, i.e. a turn was in flight and never resolved).
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
sessionLog.setLevel(Level.WARN);
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.at(SessionManager.class, Level.WARN)) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
@@ -538,18 +566,37 @@ class SessionManagerTest {
|
||||
|
||||
sessions.onTurnFailed(terminal);
|
||||
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
.orElse("no turn-failed WARN logged");
|
||||
assertTrue(warn.contains(terminal), "the log names the member: " + warn);
|
||||
assertTrue(warn.contains("BUSY"), "the log names the stage it failed at: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #525: proves the leak in {@link #onTurnFailedIsLoggedAtWarnWithThePriorState} above
|
||||
* (which runs immediately before this, via {@code @Order}) is closed. That test pins the
|
||||
* shared {@link SessionManager} logger to WARN through a {@link CapturedLog}; if {@link
|
||||
* CapturedLog#close} only detached the appender — the original bug, before this ticket's fix —
|
||||
* the level would still read WARN here instead of the {@code DEBUG} baseline this class's
|
||||
* {@code @BeforeAll} set. Runs at {@code @Order(2)}, guaranteed after {@code @Order(1)} and
|
||||
* before every other (unannotated) test in this class.
|
||||
*/
|
||||
@Test
|
||||
@Order(2)
|
||||
void sharedSessionManagerLoggerLevelIsRestoredAfterOnTurnFailedPinsWarn() {
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
assertEquals(Level.DEBUG, sessionLog.getLevel(),
|
||||
"onTurnFailedIsLoggedAtWarnWithThePriorState pins the shared SessionManager logger "
|
||||
+ "to WARN; its cleanup must restore the level it captured (DEBUG, set by "
|
||||
+ "this class's @BeforeAll) rather than leaving WARN pinned for every test "
|
||||
+ "that runs after it");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #226: a contended slot is refused through the real {@link SessionManager#acquire}
|
||||
* path before the real launcher can hand an architect charter to a process.
|
||||
@@ -576,28 +623,18 @@ class SessionManagerTest {
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberRegistry members = architectRegistry();
|
||||
sessions.setMemberLifecycle(bindFailureAfterReservation(members));
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger registryLog = (ch.qos.logback.classic.Logger)
|
||||
LoggerFactory.getLogger(MemberRegistry.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
registryLog.addAppender(appender);
|
||||
registryLog.setLevel(Level.WARN);
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.at(MemberRegistry.class, Level.WARN)) {
|
||||
MemberSession session = sessions.acquire("ltms-local", MemberRole.ARCHITECT, null,
|
||||
"/caller", "term_primary", null);
|
||||
|
||||
assertEquals(MemberRole.DEV, session.role(), "a failed reservation bind must use the fallback");
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.findFirst()
|
||||
.orElse("no slot-exhaustion WARN logged");
|
||||
assertTrue(warn.contains("ltms-local"), "the WARN names the profile: " + warn);
|
||||
assertTrue(warn.contains(session.terminalId()), "the WARN names the terminal: " + warn);
|
||||
} finally {
|
||||
registryLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -902,6 +939,82 @@ class SessionManagerTest {
|
||||
.count();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #512: a drain that releases every session cleanly used to log nothing at all — the
|
||||
* two log calls in {@code drainAll}/{@code drainSnapshot} both sit on abnormal paths, so
|
||||
* "clean drain" and "died on the first session" were indistinguishable. This asserts the new
|
||||
* {@code log.info} line fires on the ordinary, nothing-went-wrong path, and that its numbers
|
||||
* are the real counts (two released, zero abandoned) rather than just a non-empty string.
|
||||
*/
|
||||
@Test
|
||||
void drainAllLogsACompletionLineWithTheRealCountsOnACleanDrain() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession first = sessions.acquire("ltms-local", "/one", "/caller", "ownerOne");
|
||||
MemberSession second = sessions.acquire("ltms-local", "/two", "/caller", "ownerTwo");
|
||||
sessions.asPresence().markPresent(first.terminalId());
|
||||
sessions.asPresence().markPresent(second.terminalId());
|
||||
// Both stay READY — neither is delivered a turn, so neither is BUSY and the drain below
|
||||
// has nothing abnormal to hit.
|
||||
|
||||
// Pin INFO explicitly: fleetd #525 made CapturedLog itself restore the level it pins, but
|
||||
// this pin stays anyway as belt-and-braces — a later change to the sweep must not be able
|
||||
// to make this INFO assertion vacuous again by leaving some other test's WARN pin in place.
|
||||
try (CapturedLog log = CapturedLog.at(SessionManager.class, Level.INFO)) {
|
||||
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
|
||||
|
||||
assertTrue(sessions.roster().isEmpty(), "precondition: the drain actually ran");
|
||||
String info = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.INFO))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains("drain complete"))
|
||||
.findFirst()
|
||||
.orElse("no drain-complete INFO logged");
|
||||
assertTrue(info.contains("released=2"),
|
||||
"both released sessions must be counted: " + info);
|
||||
assertTrue(info.contains("abandoned=0"),
|
||||
"neither session was BUSY, so nothing was abandoned mid-turn: " + info);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #512: the same completion line must also report a non-zero abandoned count when a
|
||||
* session is still {@code BUSY} once the whole-drain deadline passes — the case the ticket
|
||||
* calls out as the one a script needs to be able to see. Reuses the same BUSY/READY mix as
|
||||
* {@link #drainAllReleasesBusyAndReadySessionsAndWaitsForBusy}, which already forces the busy
|
||||
* session to spin until the real-time deadline expires (its state never leaves BUSY on its
|
||||
* own), and adds the log assertion that test does not make.
|
||||
*/
|
||||
@Test
|
||||
void drainAllLogsANonZeroAbandonedCountForASessionStillBusyAtTheDeadline() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
SessionManager sessions = sessionManager(herdr);
|
||||
MemberSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR");
|
||||
MemberSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "ownerB");
|
||||
sessions.asPresence().markPresent(ready.terminalId());
|
||||
sessions.asPresence().markPresent(busy.terminalId());
|
||||
sessions.onDelivered(busy.terminalId(), TestTurnTokens.inert(busy.terminalId()));
|
||||
// busy never leaves BUSY — no completion is delivered — so the drain below must spin the
|
||||
// full timeout and then release it anyway, counting it abandoned.
|
||||
|
||||
// Pin INFO explicitly — see the comment in drainAllLogsACompletionLineWithTheRealCountsOnACleanDrain.
|
||||
try (CapturedLog log = CapturedLog.at(SessionManager.class, Level.INFO)) {
|
||||
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
|
||||
|
||||
assertTrue(sessions.roster().isEmpty(), "precondition: the drain actually ran");
|
||||
String info = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.INFO))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains("drain complete"))
|
||||
.findFirst()
|
||||
.orElse("no drain-complete INFO logged");
|
||||
assertTrue(info.contains("released=2"),
|
||||
"both the ready and the busy session are released: " + info);
|
||||
assertTrue(info.contains("abandoned=1"),
|
||||
"the busy session hit the deadline still BUSY and must be counted: " + info);
|
||||
}
|
||||
}
|
||||
|
||||
// --- fleetd #308: a spawn accepted while the shutdown drain is running must not orphan ---
|
||||
|
||||
@Test
|
||||
@@ -1202,21 +1315,13 @@ class SessionManagerTest {
|
||||
new WorktreeRequest("cb-581a", null));
|
||||
worktrees.failHasUncommittedWith(new WorktreeException("git status exited 128"));
|
||||
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
sessionLog.setLevel(Level.WARN);
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.at(SessionManager.class, Level.WARN)) {
|
||||
assertDoesNotThrow(() -> sessions.release(s.paneId()),
|
||||
"a throwing dirty check must not abort the release");
|
||||
|
||||
assertTrue(worktrees.removeCalls().isEmpty(),
|
||||
"the worktree is preserved when its dirty state cannot be determined");
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains(s.worktree()))
|
||||
@@ -1224,8 +1329,6 @@ class SessionManagerTest {
|
||||
.orElse("no warn logged naming the worktree");
|
||||
assertTrue(warn.contains(s.paneId()), "the WARN names the pane: " + warn);
|
||||
assertTrue(warn.contains(s.terminalId()), "the WARN names the terminal: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1392,20 +1495,12 @@ class SessionManagerTest {
|
||||
// this itself, so it no longer propagates out of release() at all.
|
||||
worktrees.failRemoveFor(b.worktree());
|
||||
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
sessionLog.setLevel(Level.WARN);
|
||||
int reaped;
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.at(SessionManager.class, Level.WARN)) {
|
||||
clock[0] = 100;
|
||||
reaped = sessions.reapIdle(10);
|
||||
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains(b.paneId()))
|
||||
@@ -1413,8 +1508,6 @@ class SessionManagerTest {
|
||||
.orElse("no worktree-removal-failure WARN logged");
|
||||
assertTrue(warn.contains(b.terminalId()), "the WARN names the failed session's terminal: " + warn);
|
||||
assertTrue(warn.contains(b.worktree()), "the WARN names the failed session's worktree: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertEquals(3, reaped,
|
||||
@@ -1460,28 +1553,18 @@ class SessionManagerTest {
|
||||
// trigger reapIdle's own guard is for, now that #283 closed the worktree-removal trigger.
|
||||
herdr.paneCloseFailsForPane("w9:pRoot_2", "internal_error");
|
||||
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
sessionLog.setLevel(Level.WARN);
|
||||
int reaped;
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.at(SessionManager.class, Level.WARN)) {
|
||||
clock[0] = 100;
|
||||
reaped = sessions.reapIdle(10);
|
||||
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains("reap failed") && m.contains(b.paneId()))
|
||||
.findFirst()
|
||||
.orElse("no reap-failed WARN logged for the failing session");
|
||||
assertTrue(warn.contains(b.terminalId()), "the WARN names the failed session's terminal: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
|
||||
assertEquals(2, reaped,
|
||||
|
||||
@@ -0,0 +1,91 @@
|
||||
package dev.ltms.fleet.session;
|
||||
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotSame;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class SessionReaperResilienceTest {
|
||||
|
||||
@Test
|
||||
void anErrorInOneIterationDoesNotStopTheNextIteration() throws Exception {
|
||||
CountDownLatch errorThrown = new CountDownLatch(1);
|
||||
CountDownLatch nextIteration = new CountDownLatch(1);
|
||||
AtomicBoolean first = new AtomicBoolean(true);
|
||||
LongSupplier clock = () -> {
|
||||
if (first.compareAndSet(true, false)) {
|
||||
errorThrown.countDown();
|
||||
throw new AssertionError("test Error from the first reap iteration");
|
||||
}
|
||||
nextIteration.countDown();
|
||||
return System.nanoTime();
|
||||
};
|
||||
SessionReaper reaper = new SessionReaper(sessionManager(clock), 60, 1);
|
||||
reaper.start();
|
||||
try {
|
||||
assertTrue(errorThrown.await(2, TimeUnit.SECONDS),
|
||||
"the first reap iteration must throw its test Error");
|
||||
assertTrue(nextIteration.await(2, TimeUnit.SECONDS),
|
||||
"the reaper must continue to the next iteration after an Error");
|
||||
} finally {
|
||||
reaper.stop();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anAbnormalExitClearsRunningSoStartCreatesANewLoop() throws Exception {
|
||||
SessionReaper reaper = new SessionReaper(sessionManager(System::nanoTime), 60, -1);
|
||||
reaper.start();
|
||||
Thread first = threadOf(reaper);
|
||||
first.join(2000);
|
||||
assertFalse(runningOf(reaper), "an abnormal loop exit must clear running");
|
||||
|
||||
reaper.start();
|
||||
Thread restarted = threadOf(reaper);
|
||||
try {
|
||||
assertNotSame(first, restarted, "start() must create a new loop after an abnormal exit");
|
||||
restarted.join(2000);
|
||||
} finally {
|
||||
reaper.stop();
|
||||
}
|
||||
}
|
||||
|
||||
private static SessionManager sessionManager(LongSupplier clock) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
|
||||
"worker: {profile} #{n}", null, null, null);
|
||||
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
return new SessionManager(launcher, new FakeWorktrees(), clock);
|
||||
}
|
||||
|
||||
private static Thread threadOf(SessionReaper reaper) throws ReflectiveOperationException {
|
||||
Field field = SessionReaper.class.getDeclaredField("thread");
|
||||
field.setAccessible(true);
|
||||
return (Thread) field.get(reaper);
|
||||
}
|
||||
|
||||
private static boolean runningOf(SessionReaper reaper) throws ReflectiveOperationException {
|
||||
Field field = SessionReaper.class.getDeclaredField("running");
|
||||
field.setAccessible(true);
|
||||
return field.getBoolean(reaper);
|
||||
}
|
||||
}
|
||||
@@ -7,13 +7,17 @@ import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.WorkspaceControl;
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.fleet.msg.TestTurnTokens;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.testing.CapturedLog;
|
||||
import org.junit.jupiter.api.AfterAll;
|
||||
import org.junit.jupiter.api.BeforeAll;
|
||||
import org.junit.jupiter.api.MethodOrderer;
|
||||
import org.junit.jupiter.api.Order;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.TestMethodOrder;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.List;
|
||||
@@ -27,9 +31,41 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
* CB-301-ext acceptance tests for worktree provisioning and config-parity overlay.
|
||||
* No live git — every Worktrees call is handled by {@link FakeWorktrees} and every herdr
|
||||
* call by {@link FakeHerdr}, matching the project's fake-based test style.
|
||||
*
|
||||
* <p>fleetd #529: only {@link #releasePreservesDirtyWorktreeAndLogsWarn} (explicitly
|
||||
* {@link Order#value() @Order(1)}) and the proving test right after it
|
||||
* ({@link #sharedSessionManagerLoggerLevelIsRestoredAfterDirtyWorktreeReleasePinsWarn},
|
||||
* {@code @Order(2)}) care about method order — every other test here has no {@code @Order} and so
|
||||
* runs after both (JUnit 5's {@link MethodOrderer.OrderAnnotation} gives an unannotated method the
|
||||
* lowest priority), in whatever relative order it already ran in.
|
||||
*/
|
||||
@TestMethodOrder(MethodOrderer.OrderAnnotation.class)
|
||||
class WorktreeSessionManagerTest {
|
||||
|
||||
/**
|
||||
* fleetd #529: the level {@link SessionManager}'s logger had when this class started, captured
|
||||
* before any test here touches it, then forced to a distinctive, known value (TRACE) so {@link
|
||||
* #sharedSessionManagerLoggerLevelIsRestoredAfterDirtyWorktreeReleasePinsWarn} can tell "the
|
||||
* level came back to what it was" apart from "the level happens to already be WARN because
|
||||
* some other test class in this JVM fork (surefire reuses forks by default) left it there".
|
||||
*/
|
||||
private static ch.qos.logback.classic.Level sessionManagerLevelBeforeThisClass;
|
||||
|
||||
@BeforeAll
|
||||
static void pinSessionManagerLoggerToAKnownBaseline() {
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
sessionManagerLevelBeforeThisClass = sessionLog.getLevel();
|
||||
sessionLog.setLevel(ch.qos.logback.classic.Level.TRACE);
|
||||
}
|
||||
|
||||
@AfterAll
|
||||
static void restoreSessionManagerLoggerLevel() {
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
sessionLog.setLevel(sessionManagerLevelBeforeThisClass);
|
||||
}
|
||||
|
||||
private static MemberRegistry members() {
|
||||
return new MemberRegistry(new FleetConfig.Fleet(Map.of(),
|
||||
Map.of("architect", new FleetConfig.Slot("ltms-local")),
|
||||
@@ -254,6 +290,7 @@ class WorktreeSessionManagerTest {
|
||||
* path, the session, and the cause an operator needs to find the work.
|
||||
*/
|
||||
@Test
|
||||
@Order(1)
|
||||
void releasePreservesDirtyWorktreeAndLogsWarn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
|
||||
@@ -262,21 +299,13 @@ class WorktreeSessionManagerTest {
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-576", null));
|
||||
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
sessionLog.setLevel(Level.WARN);
|
||||
try {
|
||||
try (CapturedLog log = CapturedLog.at(SessionManager.class, Level.WARN)) {
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertTrue(herdr.called("pane.close"), "release still tears the worker pane down");
|
||||
assertTrue(worktrees.removeCalls().isEmpty(),
|
||||
"a dirty worktree is never removed — it holds the only copy of the work");
|
||||
String warn = appender.list.stream()
|
||||
String warn = log.events().stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains("dirty worktree"))
|
||||
@@ -285,11 +314,38 @@ class WorktreeSessionManagerTest {
|
||||
assertTrue(warn.contains(s.worktree()), "the WARN names the worktree path: " + warn);
|
||||
assertTrue(warn.contains(s.terminalId()), "the WARN names the session: " + warn);
|
||||
assertTrue(warn.contains("COMPLETED"), "the WARN names the release cause: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #529 proving test: pins that the leak this ticket fixes stays fixed. {@link
|
||||
* #releasePreservesDirtyWorktreeAndLogsWarn} above (which runs immediately before this, via
|
||||
* {@code @Order}) pins the shared {@link SessionManager} logger to WARN through a {@link
|
||||
* CapturedLog}; if {@link CapturedLog#close} only detached the appender — the original bug —
|
||||
* the level would still read WARN here instead of the {@code TRACE} baseline this class's
|
||||
* {@code @BeforeAll} set. Runs at {@code @Order(2)}, guaranteed after {@code @Order(1)} and
|
||||
* before every other (unannotated) test in this class.
|
||||
*
|
||||
* <p>This proves only the WITHIN-CLASS case: JUnit 5's {@code @TestMethodOrder} orders methods
|
||||
* inside one class, not test classes relative to each other, and surefire's default class order
|
||||
* is not something a single test can force. The cross-class leak fleetd #525 measured — this
|
||||
* class's {@code releasePreservesDirtyWorktreeAndLogsWarn} pinning WARN and bleeding into a
|
||||
* later-running {@code SessionManagerTest} in the same fork — is fixed by the same {@link
|
||||
* CapturedLog} mechanism proven here, but that cross-class ordering itself is NOT asserted by
|
||||
* any test and remains unproven by construction.
|
||||
*/
|
||||
@Test
|
||||
@Order(2)
|
||||
void sharedSessionManagerLoggerLevelIsRestoredAfterDirtyWorktreeReleasePinsWarn() {
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
assertEquals(ch.qos.logback.classic.Level.TRACE, sessionLog.getLevel(),
|
||||
"releasePreservesDirtyWorktreeAndLogsWarn pins the shared SessionManager logger to "
|
||||
+ "WARN; its cleanup must restore the level it captured (TRACE, set by this "
|
||||
+ "class's @BeforeAll) rather than leaving WARN pinned for every test that "
|
||||
+ "runs after it");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-576 review (fleetd #116). A worktree that is already gone (operator cleanup,
|
||||
* {@code git worktree prune}, an earlier half-completed release) must not break teardown.
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
package dev.ltms.fleet.testing;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* fleetd #529 (promoted from {@code SessionManagerTest}, merged in #527 for fleetd #525): captures
|
||||
* a logger's output and, on {@link #close}, restores <em>both</em> the appender and the level to
|
||||
* what they were before.
|
||||
*
|
||||
* <p>{@code ch.qos.logback.classic.Logger} instances are cached per class and shared across the
|
||||
* whole JVM — and surefire reuses forks by default — so a bare {@code addAppender}/{@code
|
||||
* setLevel} pair whose {@code finally} only detaches the appender leaves the level pinned for
|
||||
* every test that runs after it, in the same class or in a completely unrelated one sharing the
|
||||
* fork. try-with-resources makes "restored the appender but not the level" impossible to write,
|
||||
* because there is only one thing to close.
|
||||
*
|
||||
* <p>New code must use this rather than hand-rolling the {@code ListAppender} + {@code setLevel} +
|
||||
* {@code finally detachAppender} pattern: use {@link #at} or {@link #of}. It is <em>not</em> yet
|
||||
* the only instance of the pattern in this test tree, and the earlier wording here said it was —
|
||||
* which would leave a reader who greps unable to tell a leftover from a violation.
|
||||
*
|
||||
* <p>Measured on main at af95897 (2026-09-12): nine test files still hand-roll it, with 42
|
||||
* {@code setLevel} calls on a raw logback {@code Logger} between them — {@code
|
||||
* FleetdStartupReportTest}, {@code GitHostShapeReportTest}, {@code MemberCredentialsGapReportTest},
|
||||
* {@code MemberTrustModelReportTest}, {@code FleetHealthMonitorTest}, {@code
|
||||
* ClaudeCodeLauncherTest}, {@code HerdrPeerLauncherAllowListWiringTest}, {@code
|
||||
* HerdrPeerLauncherCharterTest} and {@code OpenCodeLauncherTest}. Every one of them pairs its pin
|
||||
* with a restore, so none is the fleetd #525 leak and none was in fleetd #529's scope, which was
|
||||
* the 19 <em>unrestored</em> pins only. They are unmigrated, not broken.
|
||||
*
|
||||
* <p>Re-measure with the two commands below, from the repo root. A file that appears in the first
|
||||
* list and not the second still hand-rolls the pattern. When the first list comes back empty, this
|
||||
* paragraph is spent and the sentence above can go back to saying "the one way" — delete the
|
||||
* paragraph then rather than updating the count.
|
||||
*
|
||||
* <pre>{@code
|
||||
* grep -rlE '\.setLevel\(' fleetd/src/test/java --include='*.java' | grep -v CapturedLog.java
|
||||
* grep -rl 'CapturedLog' fleetd/src/test/java --include='*.java'
|
||||
* }</pre>
|
||||
*/
|
||||
public final class CapturedLog implements AutoCloseable {
|
||||
private final Logger logger;
|
||||
private final Level originalLevel;
|
||||
private final ListAppender<ILoggingEvent> appender;
|
||||
|
||||
private CapturedLog(Logger logger, Level pinnedLevel) {
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
this.logger = logger;
|
||||
this.originalLevel = logger.getLevel();
|
||||
this.appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
if (pinnedLevel != null) {
|
||||
logger.setLevel(pinnedLevel);
|
||||
}
|
||||
}
|
||||
|
||||
/** Capture {@code loggerClass}'s output, pinning its level to {@code pinnedLevel} for the
|
||||
* duration of the try-with-resources block. */
|
||||
public static CapturedLog at(Class<?> loggerClass, Level pinnedLevel) {
|
||||
return new CapturedLog((Logger) LoggerFactory.getLogger(loggerClass), pinnedLevel);
|
||||
}
|
||||
|
||||
/** Capture {@code loggerClass}'s output without changing its level. */
|
||||
public static CapturedLog of(Class<?> loggerClass) {
|
||||
return new CapturedLog((Logger) LoggerFactory.getLogger(loggerClass), null);
|
||||
}
|
||||
|
||||
/**
|
||||
* Capture the named logger's output, pinning its level to {@code pinnedLevel}. For a logger
|
||||
* obtained in production code via {@code LoggerFactory.getLogger("some-name")} rather than a
|
||||
* class — e.g. {@code AuditLog}'s {@code "audit"} logger — where {@link #at(Class, Level)}
|
||||
* would capture the wrong {@code Logger} instance (class name and logger name are different
|
||||
* strings and resolve to different cached loggers).
|
||||
*/
|
||||
public static CapturedLog at(String loggerName, Level pinnedLevel) {
|
||||
return new CapturedLog((Logger) LoggerFactory.getLogger(loggerName), pinnedLevel);
|
||||
}
|
||||
|
||||
/** Capture the named logger's output without changing its level. See {@link #at(String, Level)}. */
|
||||
public static CapturedLog of(String loggerName) {
|
||||
return new CapturedLog((Logger) LoggerFactory.getLogger(loggerName), null);
|
||||
}
|
||||
|
||||
public List<ILoggingEvent> events() {
|
||||
return appender.list;
|
||||
}
|
||||
|
||||
/**
|
||||
* Re-pin the level while this capture is still open — for example to lower it further for one
|
||||
* assertion inside a test whose fixture already pinned a coarser baseline in {@code @BeforeEach}.
|
||||
* This does not change what {@link #close} restores: that is always the level captured when
|
||||
* this instance was created, never a value set through this method.
|
||||
*/
|
||||
public void setLevel(Level level) {
|
||||
logger.setLevel(level);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(originalLevel);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,96 @@
|
||||
package dev.ltms.fleet.testing;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #537: pins {@link CapturedLog#close}'s own contract — the appender detach half, the level
|
||||
* restore half, and that {@link CapturedLog#setLevel} does not change what {@code close} restores.
|
||||
* Before this, only the level-restore half was pinned (by {@code
|
||||
* WorktreeSessionManagerTest.sharedSessionManagerLoggerLevelIsRestoredAfterDirtyWorktreeReleasePinsWarn}).
|
||||
* Measured: deleting {@code logger.detachAppender(appender);} from {@code close()} still left
|
||||
* {@code mvn clean install} green — 1701 tests, 0 failures — before this file existed.
|
||||
*
|
||||
* <p>Every test below uses a logger name no production class uses, and unique per test, so this
|
||||
* file cannot become the next entry in fleetd #525's leak family: {@link CapturedLog}'s own
|
||||
* javadoc re-measure commands (top of that file) would otherwise need to start naming this class.
|
||||
*/
|
||||
class CapturedLogTest {
|
||||
|
||||
/**
|
||||
* The appender must be detached on close: an event logged through the raw logger after close
|
||||
* must not land in {@link CapturedLog#events()}. Asserting on the observable list (rather than
|
||||
* {@code logger.iteratorForAppenders()}) is what the ticket asked for, and it is also what a
|
||||
* real leak would actually break — a later test's own {@code ListAppender} silently gaining
|
||||
* events emitted by code under test that has nothing to do with it.
|
||||
*/
|
||||
@Test
|
||||
void closeDetachesTheAppenderSoALaterLogIsNotCaptured() {
|
||||
String loggerName = "capturedlog-test-only.appender-detach";
|
||||
Logger rawLogger = (Logger) LoggerFactory.getLogger(loggerName);
|
||||
|
||||
CapturedLog log = CapturedLog.of(loggerName);
|
||||
rawLogger.info("while open");
|
||||
int eventsWhileOpen = log.events().size();
|
||||
assertEquals(1, eventsWhileOpen, "the event logged while open must be captured");
|
||||
|
||||
log.close();
|
||||
rawLogger.info("after close");
|
||||
|
||||
assertEquals(eventsWhileOpen, log.events().size(),
|
||||
"close() must detach the appender: an event logged after close must not be "
|
||||
+ "captured, but the captured list grew from " + eventsWhileOpen + " to "
|
||||
+ log.events().size());
|
||||
}
|
||||
|
||||
/**
|
||||
* The helper's own headline contract, pinned in one place independent of any production
|
||||
* class's behaviour: {@code close()} restores the level the logger had before {@link
|
||||
* CapturedLog#at} pinned it.
|
||||
*/
|
||||
@Test
|
||||
void closeRestoresTheLevelCapturedAtOpen() {
|
||||
String loggerName = "capturedlog-test-only.level-restore";
|
||||
Logger rawLogger = (Logger) LoggerFactory.getLogger(loggerName);
|
||||
rawLogger.setLevel(Level.DEBUG);
|
||||
|
||||
CapturedLog log = CapturedLog.at(loggerName, Level.ERROR);
|
||||
assertEquals(Level.ERROR, rawLogger.getLevel(), "the pinned level took effect while open");
|
||||
|
||||
log.close();
|
||||
|
||||
assertEquals(Level.DEBUG, rawLogger.getLevel(),
|
||||
"close() must restore the level captured when at() was called (DEBUG), not leave "
|
||||
+ "the pinned level (ERROR) in place");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link CapturedLog#setLevel}'s javadoc claims that re-pinning the level mid-capture does not
|
||||
* change what {@code close()} restores — that restore always uses the level captured when the
|
||||
* instance was created, never a value set through {@code setLevel}. Nothing checked this
|
||||
* before: open with a pinned WARN, call {@code setLevel(TRACE)}, close, and the result must be
|
||||
* the level from BEFORE {@code at} — neither WARN nor TRACE.
|
||||
*/
|
||||
@Test
|
||||
void setLevelDuringCaptureDoesNotChangeWhatCloseRestores() {
|
||||
String loggerName = "capturedlog-test-only.setlevel-no-effect";
|
||||
Logger rawLogger = (Logger) LoggerFactory.getLogger(loggerName);
|
||||
rawLogger.setLevel(Level.DEBUG);
|
||||
|
||||
CapturedLog log = CapturedLog.at(loggerName, Level.WARN);
|
||||
log.setLevel(Level.TRACE);
|
||||
assertEquals(Level.TRACE, rawLogger.getLevel(), "setLevel took effect immediately");
|
||||
|
||||
log.close();
|
||||
|
||||
assertEquals(Level.DEBUG, rawLogger.getLevel(),
|
||||
"close() must restore the level captured at open (DEBUG) regardless of any "
|
||||
+ "later setLevel() call: it must be neither WARN (the level pinned by "
|
||||
+ "at()) nor TRACE (the level set via setLevel() mid-capture), but got "
|
||||
+ rawLogger.getLevel());
|
||||
}
|
||||
}
|
||||
@@ -64,7 +64,106 @@
|
||||
# The two outputs side by side are the finding: any name whose hash matches between them is a
|
||||
# credential the member holds in full.
|
||||
#
|
||||
# One parse pass: line 1 = present (true/false/null), line 2 = policy mode (possibly blank),
|
||||
# lines 3-5 = knownCount/allowedCount/blockedCount, remaining lines = the known[] names. A single
|
||||
# pass avoids re-parsing (and re-risking a truthiness bug) five separate times.
|
||||
#
|
||||
# This used to feed the parser straight into `mapfile -t _FIELDS < <(producer)`. That form cannot
|
||||
# see the producer fail: `<` `<(...)` is a process substitution, not a pipeline, so `set -o
|
||||
# pipefail` does not reach inside it, and mapfile's own exit status reports whether the BUILTIN
|
||||
# ran, not whether the command substituted into it succeeded — a failing jq or python3 there still
|
||||
# leaves mapfile at rc=0 with an empty array, read as a parse that genuinely found nothing (fleetd
|
||||
# #500). Capturing the parser's output with command substitution first, and checking ITS exit
|
||||
# status, reports the producer's real failure while the fact still exists — before it is handed to
|
||||
# mapfile at all.
|
||||
#
|
||||
# mapfile then reads from that captured string with `<<<` (a herestring), not `< <(...)`: `<<<`
|
||||
# materialises the whole string in memory first, where `< <(...)` would stream it. That only
|
||||
# matters for a large producer; this one is a short credential-name policy response, so the
|
||||
# tradeoff is irrelevant here — noted because it would not be for every producer.
|
||||
parse_policy_fields() {
|
||||
if command -v jq >/dev/null 2>&1; then
|
||||
_FIELDS_RAW="$(printf '%s' "$POLICY_JSON" | jq -r '
|
||||
(.present | tostring),
|
||||
(.policy // ""),
|
||||
(.knownCount // 0 | tostring),
|
||||
(.allowedCount // 0 | tostring),
|
||||
(.blockedCount // 0 | tostring),
|
||||
(.known[]? // empty)')"
|
||||
_PARSE_STATUS=$?
|
||||
_PARSER_NAME="jq"
|
||||
else
|
||||
_FIELDS_RAW="$(printf '%s' "$POLICY_JSON" | python3 - <<'PY'
|
||||
import json, sys
|
||||
data = json.load(sys.stdin)
|
||||
print(str(data.get("present")))
|
||||
print(data.get("policy") or "")
|
||||
print(data.get("knownCount") if data.get("knownCount") is not None else 0)
|
||||
print(data.get("allowedCount") if data.get("allowedCount") is not None else 0)
|
||||
print(data.get("blockedCount") if data.get("blockedCount") is not None else 0)
|
||||
for n in (data.get("known") or []):
|
||||
print(n)
|
||||
PY
|
||||
)"
|
||||
_PARSE_STATUS=$?
|
||||
_PARSER_NAME="python3"
|
||||
fi
|
||||
|
||||
if [ "$_PARSE_STATUS" -ne 0 ]; then
|
||||
echo "refusing to run: could not parse the policy fetched from $POLICY_URL — $_PARSER_NAME exited" \
|
||||
"non-zero (status $_PARSE_STATUS). That is a parser failure, not a claim about the policy" \
|
||||
"itself; the policy response has not been read." >&2
|
||||
return 4
|
||||
fi
|
||||
|
||||
# A herestring adds a newline, so mapfile would turn an empty parser result into one empty field.
|
||||
# Keep that case separate so the refusal reports what the parser actually returned: zero fields.
|
||||
if [ -z "$_FIELDS_RAW" ]; then
|
||||
_FIELDS=()
|
||||
else
|
||||
mapfile -t _FIELDS <<< "$_FIELDS_RAW"
|
||||
fi
|
||||
|
||||
# Arity check — the CORRECTNESS fix (fleetd #500). A parser that exits 0 can still return fewer
|
||||
# than the 5 fixed fields (present, policy mode, 3 counts) that every fixed-field read in main() expects,
|
||||
# whatever the reason: a producer that printed nothing, malformed JSON that jq/python3 still
|
||||
# accepted, or a schema change upstream. main()'s slice (`_FIELDS[@]:5`) does not fire
|
||||
# `set -u` on an unset OR a short array, and every fixed-field read there used a `:-` default, so
|
||||
# without this check a short `_FIELDS` reaches the "0 known names" guard further down with the
|
||||
# same look as a policy that genuinely has 0 names. Check the count here, at the one point the
|
||||
# fact is still present, before the slice consumes it.
|
||||
if (( ${#_FIELDS[@]} < 5 )); then
|
||||
echo "refusing to run: the policy parser ($_PARSER_NAME) returned ${#_FIELDS[@]} field(s); at" \
|
||||
"least 5 are required (present, policy mode, knownCount, allowedCount, blockedCount). The" \
|
||||
"parse ran but its shape is wrong — this is not a claim about how many names the policy" \
|
||||
"knows." >&2
|
||||
return 5
|
||||
fi
|
||||
}
|
||||
|
||||
main() {
|
||||
set -uo pipefail
|
||||
# `pipefail` is not what catches the parser failure handled in parse_policy_fields() above (fleetd #500): in
|
||||
# `printf '%s' "$POLICY_JSON" | jq -r '...'`, jq is the LAST element of the pipe, so the pipeline's
|
||||
# own exit status is already jq's status, with or without pipefail. It is kept as insurance for if
|
||||
# a post-processing stage is ever appended after the parser (e.g. `| tail -n +2`) — at that point
|
||||
# the parser would sit upstream and pipefail becomes the only thing that still reports its status.
|
||||
|
||||
# --- refuse on an interpreter that cannot run this script (fleetd #500) -------------------------
|
||||
#
|
||||
# mapfile, used below to parse the policy response, was added in bash 4.0. macOS ships bash 3.2.57
|
||||
# at /bin/bash, which predates it. This script's own `set -uo pipefail` does not catch a missing
|
||||
# mapfile: the builtin just fails with "command not found" on stderr, and every line below that
|
||||
# reads the array it would have filled uses a `:-` default or a slice, neither of which `set -u`
|
||||
# catches on an unset array. Left unguarded, that chain ends in the "0 known names" refusal further
|
||||
# down — a claim about the POLICY, for a failure that is actually about the INTERPRETER. So the
|
||||
# interpreter is checked once, explicitly, before it is asked to do anything mapfile depends on.
|
||||
if (( ${BASH_VERSINFO[0]} < 4 )); then
|
||||
echo "refusing to run: this script uses mapfile, which needs bash 4 or newer. This shell is bash" \
|
||||
"${BASH_VERSION:-<unknown, no \$BASH_VERSION>}. Re-run it under a newer bash, for example:" \
|
||||
"\"\$(command -v bash)\" \"$0\"" "$@" >&2
|
||||
exit 3
|
||||
fi
|
||||
|
||||
FLEETD_HOST="${FLEETD_HOST:-http://127.0.0.1:8765}"
|
||||
POLICY_URL="${FLEETD_HOST%/}/member-credentials"
|
||||
@@ -121,31 +220,10 @@ EOF
|
||||
exit 1
|
||||
fi
|
||||
|
||||
# One parse pass: line 1 = present (true/false/null), line 2 = policy mode (possibly blank),
|
||||
# lines 3-5 = knownCount/allowedCount/blockedCount, remaining lines = the known[] names. A single
|
||||
# pass avoids re-parsing (and re-risking a truthiness bug) five separate times.
|
||||
if command -v jq >/dev/null 2>&1; then
|
||||
mapfile -t _FIELDS < <(printf '%s' "$POLICY_JSON" | jq -r '
|
||||
(.present | tostring),
|
||||
(.policy // ""),
|
||||
(.knownCount // 0 | tostring),
|
||||
(.allowedCount // 0 | tostring),
|
||||
(.blockedCount // 0 | tostring),
|
||||
(.known[]? // empty)')
|
||||
else
|
||||
mapfile -t _FIELDS < <(printf '%s' "$POLICY_JSON" | python3 - <<'PY'
|
||||
import json, sys
|
||||
data = json.load(sys.stdin)
|
||||
print(str(data.get("present")))
|
||||
print(data.get("policy") or "")
|
||||
print(data.get("knownCount") if data.get("knownCount") is not None else 0)
|
||||
print(data.get("allowedCount") if data.get("allowedCount") is not None else 0)
|
||||
print(data.get("blockedCount") if data.get("blockedCount") is not None else 0)
|
||||
for n in (data.get("known") or []):
|
||||
print(n)
|
||||
PY
|
||||
)
|
||||
fi
|
||||
# Parse the policy in one pass and refuse on any of the three failure causes. The decision, the
|
||||
# three refusals and the reasoning behind each live in parse_policy_fields() above — kept there
|
||||
# with the code rather than here, so the explanation cannot drift away from what it explains.
|
||||
parse_policy_fields || exit $?
|
||||
|
||||
PRESENT="${_FIELDS[0]:-null}"
|
||||
POLICY_MODE="${_FIELDS[1]:-}"
|
||||
@@ -164,19 +242,21 @@ case "$KNOWN_COUNT_REPORTED" in
|
||||
;;
|
||||
esac
|
||||
|
||||
# --- guard the denominator explicitly — never proceed on a zero/short count ---------------------
|
||||
# --- guard the denominator explicitly — never proceed on a zero count ---------------------------
|
||||
#
|
||||
# This is the exact trap named in the ticket: an empty (or truncated) NAMES array passes every
|
||||
# subsequent "is it set" check vacuously and prints a table that LOOKS complete. So this is checked
|
||||
# before anything else runs, with a message that says why, not just that it failed.
|
||||
# This is the exact trap named in the ticket: an empty NAMES array passes every subsequent "is it
|
||||
# set" check vacuously and prints a table that LOOKS complete. By this point the interpreter gate,
|
||||
# the parser-exit-status check, and the arity check above have already ruled out "the interpreter
|
||||
# couldn't run mapfile", "the parser failed", and "the parser returned the wrong shape" — so a zero
|
||||
# count reaching here really does mean the policy itself reports 0 known names, not a swallowed
|
||||
# failure upstream. That is still checked before anything else runs, with a message that says so.
|
||||
if [ "${#NAMES[@]}" -eq 0 ] || [ "$KNOWN_COUNT_REPORTED" -eq 0 ]; then
|
||||
cat >&2 <<EOF
|
||||
refusing to run: the policy fetched from $POLICY_URL contains 0 known names (present=${PRESENT:-unknown}).
|
||||
|
||||
Either memberCredentials: is absent/empty on the running daemon (nothing is protected — see fleetd's
|
||||
own startup warning), or the response could not be parsed. Either way, checking zero names would
|
||||
print a clean-looking table for a policy that protects nothing, or for a probe that read nothing.
|
||||
This is refused rather than reported as a pass.
|
||||
memberCredentials: is absent or empty on the running daemon — nothing is protected (see fleetd's own
|
||||
startup warning). Checking zero names would print a clean-looking table for a policy that protects
|
||||
nothing. This is refused rather than reported as a pass.
|
||||
EOF
|
||||
exit 1
|
||||
fi
|
||||
@@ -251,3 +331,8 @@ How to read this:
|
||||
hardcoded list did. If the daemon's policy changes, the next run of this script reflects it
|
||||
with no edit to this file.
|
||||
EOF
|
||||
}
|
||||
|
||||
if [[ "${BASH_SOURCE[0]}" == "$0" ]]; then
|
||||
main "$@"
|
||||
fi
|
||||
|
||||
+867
-70
File diff suppressed because it is too large
Load Diff
Executable
+109
@@ -0,0 +1,109 @@
|
||||
#!/usr/bin/env bash
|
||||
# Self-contained checks for the policy parsing guards in probe-member-credentials.sh.
|
||||
|
||||
set -euo pipefail
|
||||
|
||||
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
|
||||
PROBE="$ROOT/scripts/probe-member-credentials.sh"
|
||||
TMP="$(mktemp -d "$ROOT/.probe-member-credentials-test.XXXXXX")"
|
||||
trap 'rm -rf "$TMP"' EXIT
|
||||
|
||||
# The SOURCED guard exposes this pure parser without contacting POLICY_URL.
|
||||
source "$PROBE"
|
||||
|
||||
fail() {
|
||||
printf 'FAIL: %s\n' "$*" >&2
|
||||
return 1
|
||||
}
|
||||
|
||||
assert_equals() {
|
||||
local expected="$1" actual="$2" description="$3"
|
||||
[ "$expected" = "$actual" ] || fail "$description: expected $expected, got $actual"
|
||||
}
|
||||
|
||||
assert_contains() {
|
||||
local needle="$1" text="$2" description="$3"
|
||||
printf '%s' "$text" | grep -qF "$needle" || fail "$description: missing $needle"
|
||||
}
|
||||
|
||||
make_jq() {
|
||||
local body="$1"
|
||||
mkdir -p "$TMP/bin"
|
||||
printf '%s\n' '#!/usr/bin/env bash' "$body" > "$TMP/bin/jq"
|
||||
chmod +x "$TMP/bin/jq"
|
||||
}
|
||||
|
||||
run_parser() {
|
||||
local output rc=0
|
||||
POLICY_JSON="$(< "$TMP/policy.json")"
|
||||
POLICY_URL="fixture://member-credentials"
|
||||
output="$(PATH="$TMP/bin:$PATH" parse_policy_fields 2>&1)" || rc=$?
|
||||
PARSER_OUTPUT="$output"
|
||||
PARSER_RC="$rc"
|
||||
}
|
||||
|
||||
test_bash_older_than_four_refuses() {
|
||||
local output rc=0 version
|
||||
version="$(/bin/bash -c 'printf %s "$BASH_VERSION"')"
|
||||
output="$(/bin/bash "$PROBE" 2>&1)" || rc=$?
|
||||
assert_equals 3 "$rc" "bash 3 refusal status"
|
||||
assert_contains 'This shell is bash' "$output" "bash 3 refusal"
|
||||
assert_contains "$version" "$output" "bash 3 refusal version"
|
||||
}
|
||||
|
||||
test_parser_non_zero_refuses() {
|
||||
make_jq 'exit 17'
|
||||
run_parser
|
||||
assert_equals 4 "$PARSER_RC" "parser failure status"
|
||||
assert_contains 'jq exited non-zero (status 17)' "$PARSER_OUTPUT" "parser failure message"
|
||||
}
|
||||
|
||||
test_short_parser_output_refuses() {
|
||||
make_jq "printf '%s\\n' true enforce 3 2"
|
||||
run_parser
|
||||
assert_equals 5 "$PARSER_RC" "short parser output status"
|
||||
assert_contains 'policy parser (jq) returned 4 field(s)' "$PARSER_OUTPUT" "short parser output count"
|
||||
}
|
||||
|
||||
test_empty_parser_output_reports_zero_fields() {
|
||||
make_jq ':'
|
||||
run_parser
|
||||
assert_equals 5 "$PARSER_RC" "empty parser output status"
|
||||
assert_contains 'policy parser (jq) returned 0 field(s)' "$PARSER_OUTPUT" "empty parser output count"
|
||||
}
|
||||
|
||||
test_well_formed_policy_prints_name_table() {
|
||||
local output rc=0
|
||||
make_jq "cat '$TMP/policy.fields'"
|
||||
# Shell functions cannot be passed in an environment assignment. Run the executable through bash.
|
||||
output="$(BRIDGED_MEMBER=1 FIXTURE="$TMP/policy.json" PROBE="$PROBE" PATH="$TMP/bin:$PATH" bash -c '
|
||||
curl() { cat "$FIXTURE"; }
|
||||
export -f curl
|
||||
exec "$PROBE"
|
||||
' 2>&1)" || rc=$?
|
||||
assert_equals 0 "$rc" "well-formed policy status"
|
||||
assert_contains 'ALPHA_TOKEN' "$output" "name table"
|
||||
assert_contains 'BETA_TOKEN' "$output" "name table"
|
||||
assert_contains 'GAMMA_TOKEN' "$output" "name table"
|
||||
}
|
||||
|
||||
cat > "$TMP/policy.json" <<'JSON'
|
||||
{"present":true,"policy":"enforce","knownCount":3,"allowedCount":2,"blockedCount":1,"known":["ALPHA_TOKEN","BETA_TOKEN","GAMMA_TOKEN"]}
|
||||
JSON
|
||||
cat > "$TMP/policy.fields" <<'FIELDS'
|
||||
true
|
||||
enforce
|
||||
3
|
||||
2
|
||||
1
|
||||
ALPHA_TOKEN
|
||||
BETA_TOKEN
|
||||
GAMMA_TOKEN
|
||||
FIELDS
|
||||
|
||||
test_bash_older_than_four_refuses
|
||||
test_parser_non_zero_refuses
|
||||
test_short_parser_output_refuses
|
||||
test_empty_parser_output_reports_zero_fields
|
||||
test_well_formed_policy_prints_name_table
|
||||
printf 'PASS: probe member credentials guards\n'
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user