Compare commits
37 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| faefea14c4 | |||
| ea6896f2ef | |||
| d7f94cafa2 | |||
| 096f08c866 | |||
| 4eb720029c | |||
| 0db6d31dc2 | |||
| 158a2a84b5 | |||
| 41ebc9cf69 | |||
| 619769bb52 | |||
| 180c953c42 | |||
| 4af919af0c | |||
| 09c37061c1 | |||
| 2be287ea03 | |||
| c3b0406826 | |||
| e76fa1660b | |||
| 9dca376604 | |||
| 91be5079d6 | |||
| 26f198675c | |||
| 72f46d7c0c | |||
| 640f4d5f23 | |||
| 6edeb70bc4 | |||
| bc49d87cb8 | |||
| 14b169c410 | |||
| 2b52324d9a | |||
| b6147a39f6 | |||
| dbfc34cb6d | |||
| a20cb96730 | |||
| cda1a6a917 | |||
| 608e4496be | |||
| fa61dc587c | |||
| b6b006c651 | |||
| 63eec8a0da | |||
| cbe872b538 | |||
| 7d9a807243 | |||
| 8915e40c7d | |||
| 6cb31a10e4 | |||
| 8368a274a0 |
@@ -167,19 +167,59 @@ fails.
|
||||
- **There is no terminal or session parameter, on purpose.** The pane is always your own, resolved
|
||||
from your connection, so you can only ever roll yourself.
|
||||
- **`operatorConfirmed` is your report of what a human told you.** Do not pass `true` because you
|
||||
are confident. Ask, wait for the answer, then pass what they said. `requireOperatorConfirm`
|
||||
defaults to `true` and this is the only thing standing between a judgement call and a wiped
|
||||
session.
|
||||
are confident. Ask, wait for the answer, then pass what they said.
|
||||
|
||||
- **Whether you must ask at all depends on `leadRollover.requireOperatorConfirm`. Check it; do not
|
||||
assume.** The default is `true` (`FleetConfig.java:1426`), and then `confirm` refuses unless you
|
||||
also pass `operatorConfirmed: true`. **This host set it to `false` on 2026-09-22**, on the
|
||||
operator's explicit grant, because they do not want to approve routine context rolls. Where it is
|
||||
`false`, the three handover-file checks are the whole gate: the file must exist, be fresher than
|
||||
`maxDocAgeSeconds`, and have been modified after the open request.
|
||||
|
||||
Read the live value rather than trusting this line:
|
||||
|
||||
```bash
|
||||
grep -A1 'requireOperatorConfirm' fleetd/fleetd.yaml
|
||||
```
|
||||
|
||||
No match means the key is unset, so the default `true` applies and you must ask. The key is
|
||||
**deferred, not hot** — it is read once at boot, so an edit does nothing until the daemon is
|
||||
redeployed.
|
||||
|
||||
**Until fleetd #621 merges, the nudge text will tell you to ask the operator even where the
|
||||
daemon no longer requires it.** `LeadHeartbeatLoop.contextNotice()` hardcodes "ask the operator"
|
||||
and takes no config, so it cannot know. Trust the config value over the nudge text. Once #621 is
|
||||
merged and deployed, the nudge matches the config and this warning can be deleted.
|
||||
|
||||
- **The roll can still refuse after `confirm` returns**, and by then there is no caller to tell.
|
||||
Those outcomes are logged only, as `lead-rollover:` lines in the daemon log.
|
||||
- **The bootstrap prompt has never yet landed, and the fix is unproven (fleetd #489).** The first
|
||||
real rollover, on 2026-09-12, joined `/clear` and the bootstrap text into one line and Claude Code
|
||||
refused it as `Unknown command: /clearFresh`. The pane was never cleared and no context was lost,
|
||||
so the failure was safe — the roll simply did nothing. PR #490 fixed the cause and is deployed,
|
||||
but no roll has bootstrapped a fresh session end to end yet. **Assume it may still fail, and tell
|
||||
the operator so before you confirm.** The recovery is the same either way: the file is already
|
||||
written, so the operator starts a session and points it at the file. That is why you write the
|
||||
file before you confirm, and never the other way round.
|
||||
- **The bootstrap prompt works end to end. Measured 2026-09-22.** This used to say the fix was
|
||||
unproven (fleetd #489) and told you to expect a failure. That is no longer true. The daemon log
|
||||
now holds four `lead-rollover: rolled` lines, and three of them ran on 2026-09-22 at 10:01:43,
|
||||
10:38:28 and 11:15:47. Each one cleared the old lead and started a fresh session against the
|
||||
handover file, with the configured `bootstrapText` arriving as its first message. No context was
|
||||
lost. The old `Unknown command: /clearFresh` failure from 2026-09-12 does not appear in the log
|
||||
at all. Re-measure both numbers with:
|
||||
|
||||
```bash
|
||||
grep -c "lead-rollover: rolled" fleetd/fleetd.out # successful rolls
|
||||
grep -c "lead-rollover:" fleetd/fleetd.out # positive control: must be larger
|
||||
grep -c "Unknown command" fleetd/fleetd.out # the old failure: expect 0
|
||||
```
|
||||
|
||||
Run the control line too. A broken pattern returns a clean `0` that reads exactly like good news.
|
||||
If the first number stops growing across rolls, or `Unknown command` returns anything above 0,
|
||||
the bootstrap has regressed and this paragraph is stale again.
|
||||
|
||||
**You still write the file before you confirm, and never the other way round.** That order is not
|
||||
about the bootstrap being unreliable. It is what the daemon checks: the handover file must have
|
||||
been modified *after* the open request, or `confirm` refuses it as stale.
|
||||
|
||||
- **One warning in the log is normal and is not a failure.** Every one of the three rolls above also
|
||||
logged `/clear on term_… was never observed as WORKING after 8 consecutive IDLE/DONE polls —
|
||||
releasing rather than wedging the roll`. The daemon could not see the pane go WORKING after
|
||||
`/clear`, so it released instead of hanging. The roll then succeeded anyway. That is the safe
|
||||
branch behaving correctly. Do not report it as a broken roll.
|
||||
|
||||
## Writing style
|
||||
|
||||
|
||||
@@ -20,6 +20,13 @@ fleetd.out
|
||||
fleetd/fleetd.out
|
||||
logs/
|
||||
|
||||
# fleetd #635 follow-up — scripts/config-edit.sh's backup directory. No leading slash, so this is
|
||||
# ignored at every depth: the real one lives under fleetd/ (also named in fleetd/.gitignore, next
|
||||
# to the config it backs up), and scripts/test-config-edit.sh's own throwaway fixtures build one
|
||||
# under the repo root while the suite runs. --config can point anywhere, so the directory name is
|
||||
# ignored everywhere rather than only where the live daemon happens to use it.
|
||||
.config-backups/
|
||||
|
||||
# fleetd #480: the lead rollover handover file. `leadRollover.handoverPath` points here, and the
|
||||
# outgoing lead rewrites it on every rollover. It is a snapshot of one moment's live state —
|
||||
# unpushed branches, running builds, open questions — so it is stale the moment it is written and
|
||||
|
||||
@@ -7,6 +7,13 @@ dependency-reduced-pom.xml
|
||||
fleetd.yaml
|
||||
bridged.yaml
|
||||
|
||||
# fleetd #635 follow-up — scripts/config-edit.sh's backups of fleetd.yaml. A backup of a file
|
||||
# that must never be committed inherits that requirement. The directory is the real protection
|
||||
# (it keeps working even if the backup naming changes); the glob is a backstop for a stray
|
||||
# backup written the old way, directly beside fleetd.yaml, or by an older copy of the script.
|
||||
.config-backups/
|
||||
fleetd.yaml.bak.*
|
||||
|
||||
# CB-505 audit trail + daemon stdout/stderr — runtime records, never source
|
||||
logs/
|
||||
|
||||
|
||||
@@ -131,6 +131,27 @@ public final class Fleetd {
|
||||
}
|
||||
|
||||
static void main(String[] args) {
|
||||
main(args, ResourcePorts.system());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #625: package-private so a test can drive the literal startup sequence — not a copy
|
||||
* of it — against a non-production {@link ResourcePorts} whose {@link
|
||||
* ResourcePorts#environment()} is tainted, without ever touching the real process environment.
|
||||
* The real process environment is exactly what a test cannot taint from inside the JVM, which
|
||||
* is why nothing could pin {@link SubscriptionGuard#assertPrimaryClean} running at this call
|
||||
* site before this ticket.
|
||||
*
|
||||
* <p>{@link #main(String[])} is the one production caller, passing {@link
|
||||
* ResourcePorts#system()}. Every statement below, in the same order, is otherwise unchanged
|
||||
* from before this ticket — in particular the guard still runs here, before {@code
|
||||
* cfg.validateAll()} and before {@link FleetdAssembly#assembleAndStart} ever touches a socket,
|
||||
* a broker, or HTTP, exactly as it always has. {@code ports} is reused for the guard check and
|
||||
* then handed on to the assembly, rather than a second instance being constructed there, so a
|
||||
* test's fake backs the whole boot path with one consistent view — see {@code
|
||||
* FleetdSubscriptionGuardOrderingTest}.
|
||||
*/
|
||||
static void main(String[] args, ResourcePorts ports) {
|
||||
Path configPath = args.length > 0 ? Path.of(args[0]) : chooseDefaultConfigFile(Path.of(""));
|
||||
// The config file was renamed bridged.yaml -> fleetd.yaml. Name the file we actually
|
||||
// loaded, whichever of the two names it carries.
|
||||
@@ -166,7 +187,7 @@ public final class Fleetd {
|
||||
|
||||
// The primary/host env that launched fleetd must not be tainted.
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
guard.assertPrimaryClean(System.getenv());
|
||||
guard.assertPrimaryClean(ports.environment());
|
||||
|
||||
// Every FleetConfig.validateXxx() the operator's config can fail — CB-501's auth-exposure
|
||||
// check, CB-531's lead-tab-prefix check, CB-542's subscription-profile check, the charter
|
||||
@@ -189,7 +210,9 @@ public final class Fleetd {
|
||||
// after validateAll() (this line) puts both back under test, in the same relative order,
|
||||
// before either one does any I/O — see FleetdAssembly's javadoc for the full boot-order
|
||||
// contract this preserves exactly.
|
||||
FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ResourcePorts.system());
|
||||
// fleetd #625: the same `ports` the guard check above just used, not a second
|
||||
// ResourcePorts.system() instance — see this method's own javadoc.
|
||||
FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -194,7 +194,7 @@ final class FleetdAssembly {
|
||||
// CB-504: under supervision (launchd/systemd) fleetd can start before herdr's socket
|
||||
// exists. Wait, then degrade rather than die: serving with /healthz reporting "degraded" is
|
||||
// strictly more useful than exiting.
|
||||
Fleetd.HerdrAwaitOutcome herdrOutcome = Fleetd.awaitHerdr(herdr, ports.nanoClock(), Fleetd::sleepHerdrPoll);
|
||||
Fleetd.HerdrAwaitOutcome herdrOutcome = Fleetd.awaitHerdr(herdr, ports.nanoClock(), ports.herdrPollWait());
|
||||
boolean herdrUp = Fleetd.logHerdrWaitOutcomeAndShouldReap(herdrOutcome);
|
||||
if (herdrUp) {
|
||||
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died
|
||||
@@ -392,13 +392,21 @@ final class FleetdAssembly {
|
||||
if (cfg.leadHeartbeat() != null) {
|
||||
var hb = cfg.leadHeartbeat();
|
||||
var leadContextGauge = new LeadContextGauge();
|
||||
// fleetd #621: the context-high notice's own wording must track this same effective
|
||||
// value — LeadRollover.confirm(...) already gates the roll on it (LeadRollover.java:480),
|
||||
// and absent `leadRollover:` entirely the roll is unusable regardless (NOT_CONFIGURED),
|
||||
// so `true` (the FleetConfig.LeadRollover default) is the safe, byte-identical fallback.
|
||||
// Carried in from Fleetd.main when #612 Unit A merged main: #622 added this line to the
|
||||
// block Unit A had already moved here, so the merge would otherwise have silently
|
||||
// dropped it — with a fully green suite, because nothing pins it (see the follow-up issue).
|
||||
boolean requireOperatorConfirm = cfg.leadRollover() == null || cfg.leadRollover().requireOperatorConfirm();
|
||||
heartbeat = new LeadHeartbeatLoop(primaryRegistry, router.leadAgents(), replyInbox, sessions::roster,
|
||||
pushLoop, heartbeatScheduler, ports.nanoClock(),
|
||||
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
|
||||
metrics,
|
||||
Fleetd.leadContextSource(leadContextGauge, router.leadAgents(), leads,
|
||||
Fleetd.leadConfigDirLookup(() -> config.get().profiles(), leaders)),
|
||||
Boolean.TRUE.equals(hb.contextHighNudge()));
|
||||
Boolean.TRUE.equals(hb.contextHighNudge()), requireOperatorConfirm);
|
||||
heartbeat.start();
|
||||
} else {
|
||||
heartbeat = null;
|
||||
|
||||
@@ -43,6 +43,19 @@ public interface ResourcePorts {
|
||||
/** A monotonic elapsed-time clock. Production: {@link System#nanoTime()}. */
|
||||
LongSupplier nanoClock();
|
||||
|
||||
/**
|
||||
* fleetd #629: the per-poll wait {@code FleetdAssembly#assembleAndStart} passes to {@code
|
||||
* Fleetd#awaitHerdr} while polling for herdr's socket. Production: {@link
|
||||
* Fleetd#sleepHerdrPoll()} — a real {@code Thread.sleep}. {@link #nanoClock()} alone is not
|
||||
* enough to make {@code awaitHerdr}'s deadline controllable: the old call site passed {@code
|
||||
* Fleetd::sleepHerdrPoll} directly, hardcoded, so a test that injected a fake clock still had
|
||||
* to wait out the real sleep between each poll to ever reach the deadline — the clock looked
|
||||
* injected and was not actually controllable. A test supplies a no-op that advances its own
|
||||
* injected {@link #nanoClock()} instead, so the deadline becomes reachable without any real
|
||||
* wall-clock time passing.
|
||||
*/
|
||||
Runnable herdrPollWait();
|
||||
|
||||
/**
|
||||
* A wall-clock reading, in nanoseconds. Production: {@code System.currentTimeMillis()}
|
||||
* converted to nanoseconds. Kept separate from {@link #nanoClock()} because {@link
|
||||
|
||||
@@ -44,6 +44,11 @@ final class SystemResourcePorts implements ResourcePorts {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
return Fleetd::sleepHerdrPoll;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis());
|
||||
|
||||
@@ -2240,11 +2240,11 @@ public record FleetConfig(
|
||||
* information — which profiles, and now both values, so they can fix it without reading the
|
||||
* source — without ever taking the fleet down.
|
||||
*
|
||||
* <p>Which of the two inputs Claude Code actually follows when they disagree is intentionally
|
||||
* <em>not</em> asserted here. {@code ClaudeCodeArguments}'s javadoc used to state the
|
||||
* environment variable always wins; nobody had measured that, and this host's own
|
||||
* {@code fleetd.yaml} asserts the opposite in a comment. This method only detects and reports
|
||||
* the disagreement — see {@link dev.ltms.fleet.launch.ClaudeCodeArguments}.
|
||||
* <p>fleetd #618 measured which of the two inputs Claude Code actually follows when they
|
||||
* disagree: the environment variable wins, so {@code autoCompactWindow} is inert on a profile
|
||||
* that also sets the env var. This method only detects and reports the disagreement — it does
|
||||
* not correct it — see {@link dev.ltms.fleet.launch.ClaudeCodeArguments} for the full measured
|
||||
* precedence.
|
||||
*
|
||||
* <p>Equal values never warn: either input then produces the same session window, so there is
|
||||
* nothing to reconcile.
|
||||
@@ -2285,9 +2285,11 @@ public record FleetConfig(
|
||||
names.sort(String::compareTo);
|
||||
detail.sort(String::compareTo);
|
||||
log.warn("Claude Code profile(s) {} set disagreeing autoCompactWindow and env."
|
||||
+ "CLAUDE_CODE_AUTO_COMPACT_WINDOW — the daemon starts anyway. Fix by "
|
||||
+ "removing one key or setting equal values on each: {}. Which input Claude "
|
||||
+ "Code actually follows when they disagree is not verified here.",
|
||||
+ "CLAUDE_CODE_AUTO_COMPACT_WINDOW — the daemon starts anyway: {}. fleetd "
|
||||
+ "#618 measured that CLAUDE_CODE_AUTO_COMPACT_WINDOW wins, so "
|
||||
+ "autoCompactWindow is inert on these profiles. Set equal values on each "
|
||||
+ "to resolve this — do not just delete the env var, since that LOWERS the "
|
||||
+ "live window to autoCompactWindow's value rather than fixing anything.",
|
||||
names, String.join(", ", detail));
|
||||
}
|
||||
|
||||
|
||||
@@ -15,13 +15,15 @@ public final class ClaudeCodeArguments {
|
||||
* Append the configured Claude Code auto-compaction window when the profile opts in.
|
||||
*
|
||||
* <p>This flag and the environment variable {@code CLAUDE_CODE_AUTO_COMPACT_WINDOW} can
|
||||
* disagree. Which one Claude Code actually follows when they do is NOT verified here — this
|
||||
* javadoc used to claim the environment variable always wins, but nobody had measured that, and
|
||||
* this host's own {@code fleetd.yaml} asserts the opposite in a comment. So this javadoc no
|
||||
* longer picks a side. {@link FleetConfig#load(java.nio.file.Path)} only WARNS when a Claude
|
||||
* Code profile sets both to different values (see {@code
|
||||
* disagree, and fleetd #618 measured which one Claude Code actually follows: the environment
|
||||
* variable wins, ahead of this {@code --autocompact} flag, ahead of the settings file, ahead of
|
||||
* clientdata, the experiment, and the model default. So when a profile sets both, the flag this
|
||||
* method appends has NO effect — Claude Code reads {@code CLAUDE_CODE_AUTO_COMPACT_WINDOW}
|
||||
* first and never consults the flag. {@link FleetConfig#load(java.nio.file.Path)} only WARNS
|
||||
* when a Claude Code profile sets both to different values (see {@code
|
||||
* FleetConfig.warnConflictingAutoCompactWindows}) — it does not stop the daemon from starting,
|
||||
* and a launched session may end up honouring either window.
|
||||
* and the launched session honours the env var, not this flag. Measured against Claude Code
|
||||
* 2.1.278 (fleetd #618) — a later version could reorder this precedence.
|
||||
*/
|
||||
public static List<String> withAutoCompactWindow(List<String> argv, FleetConfig.Profile profile) {
|
||||
if (profile.autoCompactWindow() == null) {
|
||||
|
||||
@@ -31,6 +31,20 @@ public final class ConnectionIdentity {
|
||||
this.cwds = cwds;
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@link PaneLocator} this identity resolves callers against — fleetd #612 CB-185: lets a
|
||||
* test drive the exact {@link PaneLocator} a real assembly wired up (e.g. {@code
|
||||
* FleetdAssembly}'s {@code new ConnectionIdentity(new PaneLocator(herdr, memberHerdr), ...)})
|
||||
* directly with a chosen pid, bypassing the OS-dependent {@link PeerPidLookup} that {@link
|
||||
* #resolve} otherwise goes through. A full HTTP round trip cannot exercise this: {@code
|
||||
* LsofPeerPidLookup} excludes its own pid, and an in-process test client and server share one
|
||||
* JVM pid, so {@code pidForLocalPort} always returns {@code -1} and {@link PaneLocator} never
|
||||
* gets called at all.
|
||||
*/
|
||||
public PaneLocator panes() {
|
||||
return panes;
|
||||
}
|
||||
|
||||
/**
|
||||
* The caller resolved from the connection: its worker {@code terminal} (or {@code null} for the
|
||||
* primary / an off-host client), its {@code pid} (or {@code -1} if not resolvable), and whether
|
||||
|
||||
@@ -107,6 +107,13 @@ public final class FleetMcp {
|
||||
* untested identity heuristic.
|
||||
*/
|
||||
private final boolean authorizationEnforced;
|
||||
/**
|
||||
* fleetd #612 CB-185: kept as a field (rather than only captured by the {@code
|
||||
* contextExtractor} closure built in the constructor) so a test can reach the exact {@link
|
||||
* ConnectionIdentity} — and, through {@link ConnectionIdentity#panes()}, the exact {@link
|
||||
* dev.ltms.fleet.herdr.PaneLocator} — that a real assembly wired up. See {@link #identity()}.
|
||||
*/
|
||||
private final ConnectionIdentity identity;
|
||||
private final Metrics metrics; // CB-502: null → auth failures not counted
|
||||
private final CapacitySource capacity;
|
||||
private final HealthCoverageSource healthCoverage;
|
||||
@@ -396,6 +403,7 @@ public final class FleetMcp {
|
||||
Objects.requireNonNull(callers, "callers");
|
||||
this.authorizationEnforced = Objects.requireNonNull(authorizationMode, "authorizationMode")
|
||||
== AuthorizationMode.ENFORCED;
|
||||
this.identity = identity;
|
||||
this.leadChannel = leadChannel;
|
||||
this.peers = peers == null ? List.of() : List.copyOf(peers);
|
||||
this.capacity = capacity;
|
||||
@@ -755,6 +763,15 @@ public final class FleetMcp {
|
||||
return transport;
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@link ConnectionIdentity} this server resolves every caller against — fleetd #612
|
||||
* CB-185: lets a test reach the exact {@link dev.ltms.fleet.herdr.PaneLocator} a real assembly
|
||||
* wired up (via {@link ConnectionIdentity#panes()}), rather than a copy built for the test.
|
||||
*/
|
||||
public ConnectionIdentity identity() {
|
||||
return identity;
|
||||
}
|
||||
|
||||
/** Mark a connected spawned member available for the injector readiness gate. */
|
||||
static void markSpawnedMemberPresent(Principal caller, MemberPresence presence) {
|
||||
if (caller.isSpawnedMember()) {
|
||||
@@ -780,6 +797,46 @@ public final class FleetMcp {
|
||||
return server.listTools();
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — same reason as {@link #registeredTools()}: a test that must drive the REAL
|
||||
* {@link QuarantineSource} (and the real {@link BackendQuarantine} it wraps) this daemon was
|
||||
* assembled with, rather than scraping {@code FleetdAssembly.java}'s source text for the
|
||||
* constructor call that built it. Unlike {@link #registeredTools()}'s callers, that test cannot
|
||||
* live in this package: it also builds the {@code ResourcePorts} that drives
|
||||
* {@code FleetdAssembly.assembleAndStart}, and {@code ResourcePorts}' methods return
|
||||
* {@code Fleetd}-nested types that are only visible from package {@code dev.ltms.fleet} — so
|
||||
* this accessor is {@code public}, not package-private, to stay reachable from there. {@code
|
||||
* FleetdBackendQuarantineAssemblyTest} quarantines a credential twice through this exact
|
||||
* instance and checks the second cooldown is longer than the first — the one behavioural
|
||||
* difference {@link BackendQuarantine#withEscalation} and the flat two-argument constructor
|
||||
* actually produce.
|
||||
*/
|
||||
public QuarantineSource quarantineSource() {
|
||||
return quarantine;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — as {@link #quarantineSource()}, {@code public} for the same cross-package
|
||||
* reason, for the real {@link LeadSeatSource} this daemon was assembled with. {@code
|
||||
* FleetdLeadSeatAssemblyTest} calls {@code seatsFor} on this exact instance and checks it
|
||||
* reports a live lead's seat, which {@link LeadSeatSource#none()} can never do (it is a
|
||||
* constant-zero function regardless of input).
|
||||
*/
|
||||
public LeadSeatSource leadSeatSource() {
|
||||
return leadSeats;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — as {@link #quarantineSource()}, {@code public} for the same cross-package
|
||||
* reason, for the real {@link LeadRollover} (or {@code null}) this daemon was assembled with.
|
||||
* {@code FleetdLeadRolloverAssemblyTest} drives {@code open}/{@code confirm} on this exact
|
||||
* instance and waits for the real continuation to send {@code /clear} and {@code bootstrapText}
|
||||
* through the real {@code router.leadAgents()}.
|
||||
*/
|
||||
public LeadRollover leadRollover() {
|
||||
return leadRollover;
|
||||
}
|
||||
|
||||
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -76,6 +76,7 @@ public final class LeadHeartbeatLoop {
|
||||
private final Metrics metrics; // CB-512 pattern: nullable — no registry in unit tests
|
||||
private final LeadContextSource contextSource; // fleetd #609
|
||||
private final boolean contextHighNudge; // fleetd #609: opt-in, like the loop itself
|
||||
private final boolean requireOperatorConfirm; // fleetd #621: mirrors leadRollover.requireOperatorConfirm
|
||||
|
||||
/** When the current idle stretch began (nanos), or {@link #NOT_IDLE}. Single scheduler thread only. */
|
||||
private long idleSinceNanos = NOT_IDLE;
|
||||
@@ -106,12 +107,32 @@ public final class LeadHeartbeatLoop {
|
||||
* fleetd #609: as above, plus the lead's own context source and whether a HIGH reading should
|
||||
* append a hand-over notice to the loop's nudge. Pass {@link LeadContextSource#none()} and
|
||||
* {@code false} to keep the pre-#609 behaviour exactly (both existing public constructors do).
|
||||
*
|
||||
* <p>fleetd #621: delegates to the full constructor with {@code requireOperatorConfirm=true} —
|
||||
* the pre-#621 wording ("ask the operator ... only the operator can approve the roll") assumed
|
||||
* the config default, so every caller of this overload keeps that text byte-identical.
|
||||
*/
|
||||
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock,
|
||||
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
|
||||
LeadContextSource contextSource, boolean contextHighNudge) {
|
||||
this(primaryRegistry, agents, inbox, roster, pushLoop, scheduler, clock,
|
||||
idleAfterNanos, backoffMs, quietNudgeCap, metrics, contextSource, contextHighNudge, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #621: as above, plus the daemon's effective {@code leadRollover.requireOperatorConfirm}
|
||||
* value — threaded into {@link #contextNotice(boolean, LeadContextGauge.Reading, boolean, boolean)}
|
||||
* so the notice's wording tracks the config the daemon actually enforces (see {@code
|
||||
* LeadRollover.confirm}) instead of always asserting the operator gate is on.
|
||||
*/
|
||||
public LeadHeartbeatLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
Supplier<List<MemberSession>> roster, ReplyPushLoop pushLoop,
|
||||
ScheduledExecutorService scheduler, LongSupplier clock,
|
||||
long idleAfterNanos, long backoffMs, int quietNudgeCap, Metrics metrics,
|
||||
LeadContextSource contextSource, boolean contextHighNudge,
|
||||
boolean requireOperatorConfirm) {
|
||||
this.primaryRegistry = primaryRegistry;
|
||||
this.agents = agents;
|
||||
this.inbox = inbox;
|
||||
@@ -125,6 +146,7 @@ public final class LeadHeartbeatLoop {
|
||||
this.metrics = metrics;
|
||||
this.contextSource = contextSource;
|
||||
this.contextHighNudge = contextHighNudge;
|
||||
this.requireOperatorConfirm = requireOperatorConfirm;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -344,7 +366,7 @@ public final class LeadHeartbeatLoop {
|
||||
// d.contextNotified() is the value to persist once delivery is confirmed, not the value the text
|
||||
// itself should be built from. Otherwise a HIGH stretch that is still latched would never see the
|
||||
// notice at all, defeating the very check this fixes.
|
||||
String notice = contextNotice(contextHighNudge, reading, contextNotified);
|
||||
String notice = contextNotice(contextHighNudge, reading, contextNotified, requireOperatorConfirm);
|
||||
var lead = primaryRegistry.primaryTerminal();
|
||||
boolean sent = lead.isPresent() && trySend(lead.get(), fleet.nudgeText() + notice, notice);
|
||||
// The latch becomes true only when all three hold: decide() chose to notify, a notice was
|
||||
@@ -381,8 +403,9 @@ public final class LeadHeartbeatLoop {
|
||||
/**
|
||||
* fleetd #609: the text appended to a nudge when the lead's own context is full — {@code ""}
|
||||
* whenever the notice does not apply, so callers can unconditionally append this without an extra
|
||||
* branch. Wording stays plain (CEFR B1) and honest that only the operator approves a roll — this
|
||||
* loop only ever prints text, it never calls {@code fleet_handover} itself.
|
||||
* branch. Wording stays plain (CEFR B1) and honest about who actually gates the roll — see the
|
||||
* {@code requireOperatorConfirm} overload (fleetd #621) for which check that is. This loop only
|
||||
* ever prints text, it never calls {@code fleet_handover} itself.
|
||||
*
|
||||
* @param enabled the {@code leadHeartbeat.contextHighNudge} config flag
|
||||
* @param reading the lead's current {@link LeadContextGauge} reading
|
||||
@@ -401,9 +424,33 @@ public final class LeadHeartbeatLoop {
|
||||
* closing sentence ("You will not be told again until your context reads ok.") false. {@link
|
||||
* #injectNudge} is the only caller that passes a non-default {@code alreadyNotified}.
|
||||
*
|
||||
* <p>fleetd #621: delegates with {@code requireOperatorConfirm=true} — the pre-#621 default and the
|
||||
* value every existing caller of this overload (including every test written before #621) already
|
||||
* assumed, so the text this overload returns stays byte-identical.
|
||||
*
|
||||
* @param alreadyNotified whether the lead has already been told about the current HIGH stretch
|
||||
*/
|
||||
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified) {
|
||||
return contextNotice(enabled, reading, alreadyNotified, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #621: as {@link #contextNotice(boolean, LeadContextGauge.Reading, boolean)}, but the closing
|
||||
* instructions also track the daemon's effective {@code leadRollover.requireOperatorConfirm} value,
|
||||
* instead of always asserting that only the operator can approve the roll.
|
||||
*
|
||||
* <p>{@code LeadRollover.confirm(...)} already honours this flag: when it is {@code false}, the daemon
|
||||
* itself gates the roll on the three handover-file checks alone (exists, modified after the {@code
|
||||
* open()} request, and no older than {@code maxDocAgeSeconds}) and never consults {@code
|
||||
* operatorConfirmed}. Before this parameter existed, this notice told the lead to ask the operator
|
||||
* regardless — so a lead that followed its own instructions asked anyway, and setting the config knob
|
||||
* to {@code false} stopped the daemon refusing the roll without stopping the operator being
|
||||
* interrupted. This parameter is how the text is kept honest about which gate is actually live.
|
||||
*
|
||||
* @param requireOperatorConfirm the effective {@code leadRollover.requireOperatorConfirm} value
|
||||
*/
|
||||
static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified,
|
||||
boolean requireOperatorConfirm) {
|
||||
if (!enabled || alreadyNotified || reading.state() != LeadContextGauge.State.HIGH) {
|
||||
return "";
|
||||
}
|
||||
@@ -420,10 +467,18 @@ public final class LeadHeartbeatLoop {
|
||||
sb.append(" (").append(reading.compactions()).append(' ').append(compactionWord)
|
||||
.append(" so far).");
|
||||
}
|
||||
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
||||
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
|
||||
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
|
||||
+ "again until your context reads ok.");
|
||||
if (requireOperatorConfirm) {
|
||||
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
||||
+ "write the file it names, ask the operator, then call fleet_handover(action=\"confirm\", "
|
||||
+ "token, operatorConfirmed). Only the operator can approve the roll. You will not be told "
|
||||
+ "again until your context reads ok.");
|
||||
} else {
|
||||
sb.append(" A fresh session would work better. To hand over: call fleet_handover(action=\"open\"), "
|
||||
+ "write the file it names, then call fleet_handover(action=\"confirm\", token). Decide for "
|
||||
+ "yourself when to confirm: the roll goes through if the handover file exists, was "
|
||||
+ "changed after you opened it, and is not older than maxDocAgeSeconds. You will not be "
|
||||
+ "told again until your context reads ok.");
|
||||
}
|
||||
return sb.toString();
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,183 @@
|
||||
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.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.LeadChannel;
|
||||
import dev.ltms.fleet.msg.LeadChannelHandle;
|
||||
import dev.ltms.fleet.msg.LeadMessage;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
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.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertSame;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 step 4 ranks 3 and 8: the assembled daemon must use the AMQP openers from {@link
|
||||
* ResourcePorts}, and its startup report must describe the object the runtime actually owns. These
|
||||
* fakes never open a socket.
|
||||
*/
|
||||
class FleetdAssemblyAmqpOpenersTest {
|
||||
|
||||
private static final String COORD_ID = "assembly-test";
|
||||
|
||||
private static final class DurableReplyInbox implements ReplyInbox {
|
||||
@Override public void own(String target) { }
|
||||
@Override public void release(String target) { }
|
||||
@Override public void publish(String target, String msgId, String content) { }
|
||||
@Override public List<InboxMessage> peek(String target) { return List.of(); }
|
||||
@Override public boolean ack(String target, String msgId) { return false; }
|
||||
}
|
||||
|
||||
private static final class DurableLeadMailbox implements LeadChannelHandle {
|
||||
@Override public void publish(String toCoordId, LeadMessage message) { }
|
||||
@Override public List<LeadMessage> peek() { return List.of(); }
|
||||
@Override public void ack(String msgId) { }
|
||||
@Override public String selfCoordId() { return COORD_ID; }
|
||||
@Override public boolean heldDurable() { return true; }
|
||||
@Override public MailboxState inspect(String coordId) { return MailboxState.unknown(coordId); }
|
||||
@Override public void close() { }
|
||||
}
|
||||
|
||||
private static final class RecordingPorts implements ResourcePorts {
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final DurableReplyInbox replyInbox = new DurableReplyInbox();
|
||||
final DurableLeadMailbox leadMailbox = new DurableLeadMailbox();
|
||||
final AtomicInteger replyOpenCalls = new AtomicInteger();
|
||||
final AtomicInteger mailboxOpenCalls = new AtomicInteger();
|
||||
final boolean openSucceeds;
|
||||
|
||||
RecordingPorts(boolean openSucceeds) {
|
||||
this.openSucceeds = openSucceeds;
|
||||
}
|
||||
|
||||
@Override public Map<String, String> environment() { return Map.of(); }
|
||||
@Override public HerdrClient connectHerdr(Path socketPath) { return herdr; }
|
||||
@Override public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> {
|
||||
replyOpenCalls.incrementAndGet();
|
||||
if (!openSucceeds) throw new IllegalStateException("fake reply broker is down");
|
||||
return replyInbox;
|
||||
};
|
||||
}
|
||||
@Override public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfId, prefetch) -> {
|
||||
mailboxOpenCalls.incrementAndGet();
|
||||
if (!openSucceeds) throw new IllegalStateException("fake coordination broker is down");
|
||||
return leadMailbox;
|
||||
};
|
||||
}
|
||||
@Override public LongSupplier nanoClock() { return System::nanoTime; }
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
@Override public LongSupplier wallClockNanos() { return System::nanoTime; }
|
||||
@Override public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
@Override public void addShutdownHook(Runnable hook) { }
|
||||
@Override public void startHttp(Javalin app, String host, int port) { }
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Files.createDirectories(dir);
|
||||
Path config = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(config, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-reply-broker/vh"
|
||||
coordinator:
|
||||
uri: "amqp://fake-coordination-broker/vh"
|
||||
selfId: "assembly-test"
|
||||
""");
|
||||
return FleetConfig.load(config);
|
||||
}
|
||||
|
||||
private static FleetdRuntime assemble(Path dir, RecordingPorts ports) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
return FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, new ConfigRef(dir.resolve("fleetd.yaml"), cfg),
|
||||
new SubscriptionGuard(cfg.guard().hostSet())), ports);
|
||||
}
|
||||
|
||||
private static boolean reportContains(ListAppender<ILoggingEvent> appender, String text) {
|
||||
return appender.list.stream().map(ILoggingEvent::getFormattedMessage).anyMatch(message -> message.contains(text));
|
||||
}
|
||||
|
||||
@Test
|
||||
void assembledAmqpOpenersAndTheirReportsAgreeOnDurableAndFallbackStates(@TempDir Path dir) throws Exception {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
Level oldLevel = logger.getLevel();
|
||||
ListAppender<ILoggingEvent> reports = new ListAppender<>();
|
||||
reports.start();
|
||||
logger.setLevel(Level.INFO);
|
||||
logger.addAppender(reports);
|
||||
try {
|
||||
RecordingPorts durablePorts = new RecordingPorts(true);
|
||||
FleetdRuntime durable = assemble(dir.resolve("durable"), durablePorts);
|
||||
try {
|
||||
// Control: this fails loudly if the assembly did not run or used an inert opener.
|
||||
assertEquals(1, durablePorts.replyOpenCalls.get(), "assembly must call replyInboxOpener once");
|
||||
assertEquals(1, durablePorts.mailboxOpenCalls.get(), "assembly must call leadMailboxOpener once");
|
||||
assertSame(durablePorts.replyInbox, durable.replyInbox(),
|
||||
"the durable reply report must describe the exact inbox the runtime owns");
|
||||
assertSame(durablePorts.leadMailbox, durable.leadMailbox(),
|
||||
"the coordination-on report must describe the exact mailbox the runtime owns");
|
||||
assertNotNull(durable.leadCoordLoop(), "a durable mailbox must start lead coordination");
|
||||
assertTrue(reportContains(reports, "reply inbox: AMQP broker (durable)"));
|
||||
assertTrue(reportContains(reports, "lead coordination: ON as coord-id " + COORD_ID));
|
||||
} finally {
|
||||
durable.close();
|
||||
}
|
||||
|
||||
reports.list.clear();
|
||||
RecordingPorts fallbackPorts = new RecordingPorts(false);
|
||||
FleetdRuntime fallback = assemble(dir.resolve("fallback"), fallbackPorts);
|
||||
try {
|
||||
assertEquals(1, fallbackPorts.replyOpenCalls.get(), "assembly must call the failing reply opener once");
|
||||
assertEquals(1, fallbackPorts.mailboxOpenCalls.get(), "assembly must call the failing mailbox opener once");
|
||||
assertTrue(fallback.replyInbox() instanceof InMemoryReplyInbox,
|
||||
"a failed reply opener must make the runtime own the in-memory fallback");
|
||||
assertNull(fallback.leadMailbox(), "a failed mailbox opener must leave coordination off");
|
||||
assertNull(fallback.leadCoordLoop(), "coordination must not start without a mailbox");
|
||||
assertTrue(reportContains(reports, "reply inbox: in-memory (soft-state)"));
|
||||
assertTrue(reportContains(reports, "lead-to-lead messaging is OFF"));
|
||||
} finally {
|
||||
fallback.close();
|
||||
}
|
||||
} finally {
|
||||
logger.detachAppender(reports);
|
||||
logger.setLevel(oldLevel);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,230 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.herdr.PaneLocator;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
|
||||
/**
|
||||
* fleetd #612 step 2, unit B2 (CB-185, identity half). Replaces the deleted
|
||||
* {@code FleetdConnectionIdentityConstructionTest}, which pinned this claim by reading {@code
|
||||
* Fleetd.java}'s source text for {@code "new PaneLocator(herdr, memberHerdr)"}. That claim moved
|
||||
* to {@code FleetdAssembly.java} (fleetd #612 Unit A) and is pinned here instead, by driving the
|
||||
* real {@link ConnectionIdentity} — via {@code runtime.mcp().identity()}, not a copy — that {@link
|
||||
* FleetdAssembly#assembleAndStart} built.
|
||||
*
|
||||
* <p><strong>What this guards against</strong> (from the deleted test's own javadoc): pinning
|
||||
* {@code PaneLocator} to {@code memberHerdr} alone leaves every LEAD's own MCP connection
|
||||
* unresolvable ({@code callerTerminal == null}) the moment {@code memberHerdrSocket} names a
|
||||
* second daemon, which breaks {@code fleet_reply}/{@code fleet_ask}/{@code fleet_whoami} for a
|
||||
* lead. {@code PaneLocatorTest} already proves {@link PaneLocator} itself can search two clients
|
||||
* given two — the gap this pins is that the assembly actually passes it two, and in the right
|
||||
* order (lead first).
|
||||
*
|
||||
* <p><strong>Why this cannot be driven through a real MCP/HTTP round trip.</strong> The natural
|
||||
* way to observe {@code ConnectionIdentity} would be a real {@code fleet_whoami} call over the
|
||||
* built {@code FleetMcp}, the way {@code FleetMcpContextExtractorTest} drives its own
|
||||
* hand-built one. That does not work for the REAL assembly, because {@code FleetdAssembly} wires
|
||||
* {@code ConnectionIdentity} with a hardcoded {@code new LsofPeerPidLookup()} (see {@code
|
||||
* FleetdAssembly.java:444}), and {@code LsofPeerPidLookup} explicitly excludes its own PID — see
|
||||
* its javadoc: "we exclude our own PID and take the other end". In a JUnit test the HTTP client
|
||||
* and the daemon under test run in the very same JVM, so the "client" and "server" ends of the
|
||||
* loopback connection ARE the same PID, and {@code pidForLocalPort} always returns {@code -1}
|
||||
* before {@link PaneLocator} is ever reached — proving nothing about which daemon(s) got searched.
|
||||
* This test instead reaches the real {@link PaneLocator} the assembly built (through {@link
|
||||
* ConnectionIdentity#panes()}, added for exactly this) and drives it with a chosen pid directly,
|
||||
* bypassing the OS-dependent PID lookup entirely — a legitimate substitute, since the pid lookup
|
||||
* is not what CB-185 is about.
|
||||
*/
|
||||
class FleetdAssemblyConnectionIdentityTest {
|
||||
|
||||
/** Same shape as {@code FleetdAssemblyLifecycleTest}'s fake, but keys {@code connectHerdr} by
|
||||
* socket path so the lead and member daemons can be two DIFFERENT {@link FakeHerdr}s. */
|
||||
private static final class TwoHerdrResourcePorts implements ResourcePorts {
|
||||
|
||||
final Map<Path, HerdrClient> herdrsBySocket = new LinkedHashMap<>();
|
||||
final CopyOnWriteArrayList<ScheduledExecutorService> schedulers = new CopyOnWriteArrayList<>();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
HerdrClient client = herdrsBySocket.get(socketPath);
|
||||
if (client == null) {
|
||||
throw new IllegalStateException("no fake herdr registered for socket " + socketPath);
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> new dev.ltms.fleet.msg.InMemoryReplyInbox();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
schedulers.add(scheduler);
|
||||
return scheduler;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Deliberately never bind — this test never issues a real HTTP request.
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private FleetdRuntime runtime;
|
||||
private TwoHerdrResourcePorts ports;
|
||||
|
||||
@AfterEach
|
||||
void tearDown() {
|
||||
if (ports != null && ports.shutdownHook != null) {
|
||||
ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
|
||||
private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock");
|
||||
private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock");
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: "%s"
|
||||
memberHerdrSocket: "%s"
|
||||
lifecycle:
|
||||
idleTtlSeconds: 600
|
||||
health:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
""".formatted(LEAD_SOCKET, MEMBER_SOCKET));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
private FleetdRuntime assemble(Path dir, FakeHerdr lead, FakeHerdr member) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
ports = new TwoHerdrResourcePorts();
|
||||
ports.herdrsBySocket.put(LEAD_SOCKET, lead);
|
||||
ports.herdrsBySocket.put(MEMBER_SOCKET, member);
|
||||
runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
return runtime;
|
||||
}
|
||||
|
||||
/**
|
||||
* The pin. {@code lead} carries the one pane {@link FakeHerdr}'s canned {@code
|
||||
* pane.process_info} ties to {@link FakeHerdr#WORKER_PID} (pane {@code w2:p7}); {@code member}
|
||||
* reports NO panes at all ({@link FakeHerdr#withNoPanes()}) — modelling a second daemon that
|
||||
* simply does not host the caller's pane, exactly the CB-185 javadoc's scenario for a lead's
|
||||
* own connection. If {@code PaneLocator} only ever searches the member daemon (the bug), this
|
||||
* pid resolves to nothing, because the pane that owns it lives on the LEAD daemon the bug
|
||||
* skips.
|
||||
*/
|
||||
@Test
|
||||
void connectionIdentitySearchesTheLeadDaemonNotJustTheMemberOne(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
FakeHerdr member = new FakeHerdr().withNoPanes();
|
||||
|
||||
assemble(dir, lead, member);
|
||||
|
||||
PaneLocator panes = runtime.mcp().identity().panes();
|
||||
PaneLocator.Lookup lookup = panes.terminalForPid(FakeHerdr.WORKER_PID);
|
||||
|
||||
assertEquals("term_a", lookup.terminal(),
|
||||
"the pane owning WORKER_PID lives on the LEAD daemon only (the member fake reports "
|
||||
+ "no panes) — PaneLocator must still find it, which is only possible if it "
|
||||
+ "searches the lead client and not just the member one");
|
||||
}
|
||||
|
||||
/**
|
||||
* The mirror control: when the pane instead lives ONLY on the member daemon (the lead reports
|
||||
* no panes), the lookup must still find it — proving the member client is genuinely searched
|
||||
* too, not merely tolerated as a second, always-losing argument.
|
||||
*/
|
||||
@Test
|
||||
void connectionIdentityAlsoSearchesTheMemberDaemon(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr().withNoPanes();
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
|
||||
assemble(dir, lead, member);
|
||||
|
||||
PaneLocator panes = runtime.mcp().identity().panes();
|
||||
PaneLocator.Lookup lookup = panes.terminalForPid(FakeHerdr.WORKER_PID);
|
||||
|
||||
assertEquals("term_a", lookup.terminal(),
|
||||
"the pane owning WORKER_PID lives on the MEMBER daemon only — PaneLocator must "
|
||||
+ "find it there too");
|
||||
}
|
||||
|
||||
/** Sanity control: a pid nobody owns resolves to nothing on either daemon. */
|
||||
@Test
|
||||
void aPidNoPaneOwnsResolvesToNoTerminalOnEitherDaemon(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
|
||||
assemble(dir, lead, member);
|
||||
|
||||
PaneLocator panes = runtime.mcp().identity().panes();
|
||||
PaneLocator.Lookup lookup = panes.terminalForPid(999_999L);
|
||||
|
||||
assertNull(lookup.terminal());
|
||||
}
|
||||
}
|
||||
@@ -140,6 +140,14 @@ class FleetdAssemblyCoordinatorLifecycleTest {
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// No real HTTP bind in a unit test.
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private static final class SentinelReplyInbox implements ReplyInbox {
|
||||
|
||||
@@ -0,0 +1,267 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.Timeout;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.net.URI;
|
||||
import java.net.http.HttpClient;
|
||||
import java.net.http.HttpRequest;
|
||||
import java.net.http.HttpResponse;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #612 step 2, unit B2 (CB-185, {@code FleetApp} half). Replaces the deleted {@code
|
||||
* FleetdFleetAppConstructionTest}, which pinned this claim by reading {@code Fleetd.java}'s
|
||||
* source text for {@code "new FleetApp(herdr, memberHerdr, workers,"}. That claim moved to {@code
|
||||
* FleetdAssembly.java} (fleetd #612 Unit A) and is pinned here instead, by driving the real {@code
|
||||
* Javalin} app — via {@code runtime.app()}, not a copy — that {@link
|
||||
* FleetdAssembly#assembleAndStart} built and handed to {@link FleetdRuntime}.
|
||||
*
|
||||
* <p><strong>What this guards against</strong> (from the deleted test's own javadoc): constructing
|
||||
* {@code FleetApp} with the lead-only {@code herdr} client (dropping {@code memberHerdr}) makes
|
||||
* {@code GET /healthz} report green while the MEMBER daemon is down — so every spawn fails
|
||||
* invisibly — and silently drops every member workspace from {@code GET /sessions}. {@code
|
||||
* FleetAppTwoDaemonTest} already proves {@code FleetApp} itself merges/gates correctly given two
|
||||
* clients; the gap this pins is that the assembly actually passes it two.
|
||||
*
|
||||
* <p><strong>Both directions, not just one</strong> (fleetd #612 issue comment 17525): the deleted
|
||||
* guard's positive assertion required the exact pair {@code "new FleetApp(herdr, memberHerdr,
|
||||
* workers,"}, which does not survive EITHER daemon being dropped. An earlier version of this class
|
||||
* only proved the member-dropped direction, which left {@code new FleetApp(memberHerdr,
|
||||
* memberHerdr, ...)} — the symmetric bug, {@code /healthz} green while the LEAD daemon is down —
|
||||
* an undetected regression. {@link #healthzGoesRedWhenTheLeadDaemonIsDownEvenThoughTheMemberIsUp}
|
||||
* closes that.
|
||||
*
|
||||
* <p>Unlike the {@code ConnectionIdentity} half of CB-185 ({@code
|
||||
* FleetdAssemblyConnectionIdentityTest}), {@code /healthz} needs no caller identity at all, so
|
||||
* this test can bind {@link FleetdRuntime#app()} to a REAL ephemeral port (exactly {@code
|
||||
* FleetAppTwoDaemonTest} does for its own hand-built {@code FleetApp}) and drive it with a real
|
||||
* {@code HttpClient} — no accessor needed for this half.
|
||||
*
|
||||
* <p><strong>{@code GET /sessions} could not be driven the same way</strong>, so this class does
|
||||
* not pin the merge half of the deleted test's javadoc. {@code /sessions} requires
|
||||
* {@code Authz.Action.READ}, which — through the REAL assembly's real {@code
|
||||
* CallerResolver}/{@code ConnectionIdentity} (built with a hardcoded {@code
|
||||
* new LsofPeerPidLookup()}) — needs {@code Caller.resolved()}, i.e. a real positive pid from
|
||||
* {@code lsof}. {@code LsofPeerPidLookup} excludes its own pid (see its javadoc), and a JUnit
|
||||
* test's HTTP client and the daemon under test share one JVM pid, so the resolved pid is always
|
||||
* {@code -1} and every such request is refused as {@code ANONYMOUS} (fleetd #317's fail-closed
|
||||
* rule) before the route handler — and its {@code memberHerdr} merge — is ever reached. Verified
|
||||
* directly: driving {@code GET /sessions} here returns {@code 401 unauthenticated}, not the
|
||||
* merged body. {@code FleetAppTwoDaemonTest} avoids this because it builds {@code FleetApp} with
|
||||
* {@code callers: null}, which is not what the real assembly passes. The {@code /healthz} pin
|
||||
* below is what this class relies on for CB-185's {@code FleetApp} half; {@code
|
||||
* FleetAppTwoDaemonTest} remains the full behavioural proof that {@code FleetApp} itself merges
|
||||
* {@code /sessions} correctly once handed two clients.
|
||||
*
|
||||
* <p><strong>fleetd #629 follow-up.</strong> The fix below (see {@link TwoHerdrResourcePorts})
|
||||
* makes {@link #healthzGoesRedWhenTheLeadDaemonIsDownEvenThoughTheMemberIsUp}'s fake {@code
|
||||
* nanoClock()} frozen unless {@code herdrPollWait()} itself advances it. That is a sharper pin
|
||||
* than an assertion — if a future edit to {@code FleetdAssembly} ever bypasses {@code
|
||||
* ports.herdrPollWait()} again (e.g. reverting to a hardcoded {@code Thread.sleep}), the clock
|
||||
* never advances, {@code Fleetd#awaitHerdr}'s deadline is never reached, and this test hangs
|
||||
* forever instead of failing — proven by deliberately reintroducing that exact regression while
|
||||
* fixing this ticket. {@code @Timeout} turns that silent hang into a bounded, named test failure:
|
||||
* {@code SEPARATE_THREAD} so JUnit's timeout governor can actually interrupt a thread stuck in a
|
||||
* real {@code Thread.sleep} loop (the default {@code SAME_THREAD} mode cannot — it only measures
|
||||
* elapsed time after the test method returns on its own, which never happens here). 10 seconds is
|
||||
* roughly 150x the real passing times measured here (~0.06s), so a slow CI machine has no reason
|
||||
* to flake, and it is still 3x faster than discovering the regression by burning a CI job's whole
|
||||
* wall-clock budget.
|
||||
*/
|
||||
@Timeout(value = 10, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
|
||||
class FleetdAssemblyFleetAppTest {
|
||||
|
||||
private static final class TwoHerdrResourcePorts implements ResourcePorts {
|
||||
|
||||
final Map<Path, HerdrClient> herdrsBySocket = new LinkedHashMap<>();
|
||||
final CopyOnWriteArrayList<ScheduledExecutorService> schedulers = new CopyOnWriteArrayList<>();
|
||||
// fleetd #629: a fake, advanceable clock — NOT System::nanoTime. awaitHerdr's poll wait
|
||||
// (herdrPollWait() below) advances this on every poll instead of sleeping for real, so the
|
||||
// down-lead test below reaches awaitHerdr's deadline without burning real wall-clock time.
|
||||
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); // arbitrary non-zero start
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
HerdrClient client = herdrsBySocket.get(socketPath);
|
||||
if (client == null) {
|
||||
throw new IllegalStateException("no fake herdr registered for socket " + socketPath);
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> new dev.ltms.fleet.msg.InMemoryReplyInbox();
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return nowNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// fleetd #629: advance the fake clock instead of a real Thread.sleep, so awaitHerdr's
|
||||
// deadline is reached in real time regardless of the configured poll interval.
|
||||
return () -> nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(1));
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
schedulers.add(scheduler);
|
||||
return scheduler;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Deliberately never bind here — this test binds runtime.app() itself, for real, below.
|
||||
}
|
||||
}
|
||||
|
||||
private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock");
|
||||
private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock");
|
||||
|
||||
private final HttpClient http = HttpClient.newHttpClient();
|
||||
private FleetdRuntime runtime;
|
||||
private TwoHerdrResourcePorts ports;
|
||||
private Javalin boundApp;
|
||||
|
||||
@AfterEach
|
||||
void tearDown() {
|
||||
if (boundApp != null) {
|
||||
boundApp.stop();
|
||||
}
|
||||
if (ports != null && ports.shutdownHook != null) {
|
||||
ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: "%s"
|
||||
memberHerdrSocket: "%s"
|
||||
lifecycle:
|
||||
idleTtlSeconds: 600
|
||||
health:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
""".formatted(LEAD_SOCKET, MEMBER_SOCKET));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
/** Assembles the real graph, then binds the real {@code Javalin app} to an ephemeral port. */
|
||||
private int assembleAndBind(Path dir, FakeHerdr lead, FakeHerdr member) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
ports = new TwoHerdrResourcePorts();
|
||||
ports.herdrsBySocket.put(LEAD_SOCKET, lead);
|
||||
ports.herdrsBySocket.put(MEMBER_SOCKET, member);
|
||||
runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
boundApp = runtime.app().start("127.0.0.1", 0);
|
||||
return boundApp.port();
|
||||
}
|
||||
|
||||
private HttpResponse<String> get(int port, String path) throws Exception {
|
||||
HttpRequest req = HttpRequest.newBuilder(URI.create("http://127.0.0.1:" + port + path)).GET().build();
|
||||
return http.send(req, HttpResponse.BodyHandlers.ofString());
|
||||
}
|
||||
|
||||
/**
|
||||
* The pin. The MEMBER daemon is down; the LEAD daemon is healthy. If the assembly built
|
||||
* {@code FleetApp} with only the lead client (the bug: passing {@code herdr} where {@code
|
||||
* memberHerdr} is expected), the down member is invisible and {@code /healthz} stays 200.
|
||||
*/
|
||||
@Test
|
||||
void healthzGoesRedWhenTheMemberDaemonIsDownEvenThoughTheLeadIsUp(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
FakeHerdr member = new FakeHerdr().healthy(false);
|
||||
|
||||
int port = assembleAndBind(dir, lead, member);
|
||||
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
assertEquals(503, res.statusCode(),
|
||||
"a down MEMBER daemon must not be masked by a healthy lead: " + res.body());
|
||||
}
|
||||
|
||||
/** Sanity control: both daemons healthy must still be green through the real assembly. */
|
||||
@Test
|
||||
void healthzIsGreenWhenBothDaemonsAreUp(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
|
||||
int port = assembleAndBind(dir, lead, member);
|
||||
|
||||
assertEquals(200, get(port, "/healthz").statusCode());
|
||||
}
|
||||
|
||||
/**
|
||||
* The symmetric pin (fleetd #612 issue comment 17525): the LEAD daemon is down; the MEMBER
|
||||
* daemon is healthy. If the assembly built {@code FleetApp} with only the member client
|
||||
* (dropping {@code herdr} — the mirror of the bug above, {@code new FleetApp(memberHerdr,
|
||||
* memberHerdr, ...)}), the down LEAD is invisible and {@code /healthz} stays 200. Without this
|
||||
* case the pair above is one-directional and does not cover the deleted guard's positive
|
||||
* assertion (it required BOTH {@code herdr,} and {@code memberHerdr,} in that order).
|
||||
*/
|
||||
@Test
|
||||
void healthzGoesRedWhenTheLeadDaemonIsDownEvenThoughTheMemberIsUp(@TempDir Path dir) throws Exception {
|
||||
FakeHerdr lead = new FakeHerdr().healthy(false);
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
|
||||
int port = assembleAndBind(dir, lead, member);
|
||||
|
||||
HttpResponse<String> res = get(port, "/healthz");
|
||||
assertEquals(503, res.statusCode(),
|
||||
"a down LEAD daemon must not be masked by a healthy member: " + res.body());
|
||||
}
|
||||
}
|
||||
+198
@@ -0,0 +1,198 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.health.FleetHealthMonitor;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 rank 7 — {@code FleetdAssembly.java:429} wires {@link FleetHealthMonitor}'s {@code
|
||||
* failTarget} callback with {@code Fleetd.healthFailTarget(messages)}. {@link
|
||||
* FleetdHealthFailTargetWiringTest} already pins that the FACTORY itself delegates to {@code
|
||||
* messages::abandon}, but it calls {@code Fleetd.healthFailTarget} directly — it never drives {@code
|
||||
* FleetdAssembly.assembleAndStart} and so cannot see whether the real call site at {@code :429}
|
||||
* still passes it the real, assembled {@link MessageService}. Swapping that argument for a no-op
|
||||
* {@code (a, b) -> {}} compiles clean and leaves the whole suite — including the factory-level test
|
||||
* — green: a dead member's waiting ticket then sits {@code PENDING} for the full 30-minute async
|
||||
* timeout instead of failing immediately.
|
||||
*
|
||||
* <p>This test assembles the real daemon with {@code health.enabled: true}, pulls the REAL {@code
|
||||
* failTarget} {@link BiConsumer} out of the REAL, assembled {@link FleetHealthMonitor} (via
|
||||
* reflection — the field is package-private to {@code dev.ltms.fleet.health}, and nothing public
|
||||
* exposes it; {@code StatusPollerResilienceTest} already uses the same technique in this suite), and
|
||||
* invokes it directly against the REAL {@link MessageService} {@link FleetdRuntime#messages()}
|
||||
* returns. A no-op lambda swapped in at the call site leaves the ticket {@code PENDING} forever,
|
||||
* which this test catches; the real one fails it.
|
||||
*/
|
||||
class FleetdAssemblyHealthFailTargetBehaviouralTest {
|
||||
|
||||
private static final String TARGET = "term_a";
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> {
|
||||
throw new UnsupportedOperationException("no broker: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException("no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
health:
|
||||
enabled: true
|
||||
intervalSeconds: 30
|
||||
""");
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled FleetHealthMonitor's failTarget reaches the real "
|
||||
+ "MessageService.abandon, not a no-op")
|
||||
void assembledHealthFailTargetReachesRealMessages(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
|
||||
|
||||
// Surefire runs the whole suite in one JVM fork, so the scheduler/loops this assembly starts
|
||||
// (SessionReaper, StatusPoller, the health monitor) must be torn down here, on the failure
|
||||
// path too — hence the try/finally, not just a statement at the end of the happy path.
|
||||
try {
|
||||
FleetHealthMonitor healthMonitor = runtime.healthMonitor();
|
||||
assertNotNull(healthMonitor, "health.enabled: true in this test's config, so "
|
||||
+ "FleetdAssembly.assembleAndStart must have built a real FleetHealthMonitor");
|
||||
|
||||
Field field = FleetHealthMonitor.class.getDeclaredField("failTarget");
|
||||
field.setAccessible(true);
|
||||
BiConsumer<String, String> failTarget = (BiConsumer<String, String>) field.get(healthMonitor);
|
||||
assertNotNull(failTarget, "FleetHealthMonitor's failTarget must never be null — the "
|
||||
+ "constructor itself requires it");
|
||||
|
||||
MessageService messages = runtime.messages();
|
||||
|
||||
// --- loud control: prove the assembled MessageService is actually wired up and a ticket is
|
||||
// genuinely PENDING before failTarget ever runs. If this fails, the test below would pass
|
||||
// vacuously on a MessageService that never got a ticket in the first place. TARGET has no
|
||||
// live agent behind it (no session was ever acquired), so nothing resolves this ticket on
|
||||
// its own — it stays PENDING until failTarget (or a timeout) ends it.
|
||||
String ticket = messages.sendAsync(TARGET, "long task");
|
||||
MessageService.TaskView before = messages.poll(ticket);
|
||||
assertEquals(MessageService.Phase.PENDING, before.phase(),
|
||||
"control: the async ticket must be PENDING before failTarget runs");
|
||||
|
||||
failTarget.accept(TARGET, "member unreachable (health monitor)");
|
||||
|
||||
MessageService.TaskView after = awaitTerminal(messages, ticket);
|
||||
assertEquals(MessageService.Phase.FAILED, after.phase(),
|
||||
"FleetdAssembly.java:429 must pass Fleetd.healthFailTarget(messages) built from the "
|
||||
+ "SAME assembled MessageService — a no-op BiConsumer at that call site leaves "
|
||||
+ "this ticket PENDING for the full 30-minute async timeout instead of failing it");
|
||||
assertTrue(after.detail() != null && after.detail().contains("member unreachable"),
|
||||
"the failure reason passed to failTarget.accept must reach MessageService.abandon and "
|
||||
+ "end up in the ticket's detail");
|
||||
} finally {
|
||||
// Proof the teardown actually ran, not just an assurance that a finally was added: the
|
||||
// captured shutdown hook's close order (FleetdAssemblyLifecycleTest) closes the herdr
|
||||
// client last, so ports.herdr.closed flips to true only if this hook really executed.
|
||||
ports.shutdownHook.run();
|
||||
assertTrue(ports.herdr.closed, "the captured shutdown hook must have run and closed herdr — "
|
||||
+ "proof this test's assembled background loops/scheduler were torn down");
|
||||
}
|
||||
}
|
||||
|
||||
private static MessageService.TaskView awaitTerminal(MessageService messages, String ticket)
|
||||
throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 5000;
|
||||
MessageService.TaskView view = messages.poll(ticket);
|
||||
while (view.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
view = messages.poll(ticket);
|
||||
}
|
||||
return view;
|
||||
}
|
||||
}
|
||||
@@ -126,6 +126,14 @@ class FleetdAssemblyLifecycleTest {
|
||||
ledger.add("startHttp");
|
||||
this.startedApp = app;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
/** A fake {@link ReplyInbox} that is also {@link AutoCloseable}, so the ledger can prove it closes. */
|
||||
|
||||
@@ -0,0 +1,232 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import dev.ltms.fleet.msg.ReplyPushLoop;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 rank 6 — {@code FleetdAssembly.java:447} wires {@code
|
||||
* sessions.onRelease(Fleetd.releaseCleanup(messages, replyInbox, primaryRegistry))}. {@link
|
||||
* FleetdReleaseCleanupWiringTest} already pins that the FACTORY {@code Fleetd.releaseCleanup}
|
||||
* itself reaches all three collaborators — but it calls the factory directly, never {@code
|
||||
* FleetdAssembly.assembleAndStart}, so it cannot see whether the real call site at {@code :447}
|
||||
* still registers it (as opposed to a no-op {@code detail -> { }}) or still passes it the REAL,
|
||||
* assembled {@code messages}/{@code replyInbox}/{@code primaryRegistry}. Swapping the registered
|
||||
* listener for a no-op at that call site compiles clean and leaves the whole suite — including the
|
||||
* factory-level test — green: EVERY teardown then leaks a stuck rendezvous waiter, an unreleased
|
||||
* reply-inbox consumer, and a stale lead binding, all three at once.
|
||||
*
|
||||
* <p>This test assembles the real daemon with {@code idleSleepGuard.enabled: false} — the ONLY
|
||||
* other {@code onRelease} registration in {@code FleetdAssembly} (see {@code
|
||||
* dev.ltms.fleet.power.IdleSleepGuard}'s own wiring at {@code FleetdAssembly.java:228}) — so the
|
||||
* real {@link SessionManager}'s release-listener list holds exactly the one listener this call site
|
||||
* registers. It pulls that REAL listener out via reflection (the list itself is private, like
|
||||
* {@code StatusPollerResilienceTest}'s use of the same technique elsewhere in this suite), invokes
|
||||
* it directly, and asserts all three collaborator effects against the REAL, assembled {@link
|
||||
* MessageService} ({@link FleetdRuntime#messages()}), the REAL {@link ReplyInbox} ({@link
|
||||
* FleetdRuntime#replyInbox()}), and the REAL {@link PrimaryRegistry} — reached through {@link
|
||||
* FleetdRuntime#pushLoop()}, the only other accessor that was handed the same {@code
|
||||
* primaryRegistry} instance ({@code FleetdAssembly.java:382}), since {@code FleetMcp} never exposes
|
||||
* it.
|
||||
*/
|
||||
class FleetdAssemblyReleaseCleanupBehaviouralTest {
|
||||
|
||||
private static final String TARGET = "term_a";
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> {
|
||||
throw new UnsupportedOperationException("no broker: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException("no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
""");
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled release listener reaches messages.abandon, "
|
||||
+ "replyInbox.release, AND primaryRegistry.forgetDelegation — all three leaks at once")
|
||||
void assembledReleaseListenerReachesAllThreeCollaborators(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
|
||||
|
||||
// Surefire runs the whole suite in one JVM fork, so the scheduler/loops this assembly starts
|
||||
// must be torn down here, on the failure path too — hence the try/finally, not just a
|
||||
// statement at the end of the happy path.
|
||||
try {
|
||||
// --- reach into SessionManager's private release-listener list. idleSleepGuard.enabled:
|
||||
// false above means FleetdAssembly.java:228 never registers, so this list must hold EXACTLY
|
||||
// the one listener :447 registers.
|
||||
Field listenersField = SessionManager.class.getDeclaredField("releaseListeners");
|
||||
listenersField.setAccessible(true);
|
||||
List<Consumer<SessionManager.ReleaseDetail>> releaseListeners =
|
||||
(List<Consumer<SessionManager.ReleaseDetail>>) listenersField.get(runtime.sessions());
|
||||
assertEquals(1, releaseListeners.size(), "control: with idleSleepGuard.enabled: false, "
|
||||
+ "FleetdAssembly.java:447 must be the ONLY onRelease registration — a different "
|
||||
+ "count means this test is no longer isolating the call site it claims to pin");
|
||||
Consumer<SessionManager.ReleaseDetail> releaseListener = releaseListeners.get(0);
|
||||
|
||||
MessageService messages = runtime.messages();
|
||||
ReplyInbox replyInbox = runtime.replyInbox();
|
||||
|
||||
// primaryRegistry is never exposed by FleetdRuntime directly — ReplyPushLoop is the other
|
||||
// collaborator FleetdAssembly.java:382 hands the SAME instance to, so reach it from there.
|
||||
Field primaryRegistryField = ReplyPushLoop.class.getDeclaredField("primaryRegistry");
|
||||
primaryRegistryField.setAccessible(true);
|
||||
PrimaryRegistry primaryRegistry = (PrimaryRegistry) primaryRegistryField.get(runtime.pushLoop());
|
||||
assertNotNull(primaryRegistry, "control: the assembled ReplyPushLoop must hold a real "
|
||||
+ "PrimaryRegistry instance");
|
||||
|
||||
// --- loud controls: set up the "before" state each collaborator's effect is measured
|
||||
// against, against the REAL assembled objects. If any of these three fails, the test below
|
||||
// would pass vacuously because the subject it claims to observe never existed in the first
|
||||
// place.
|
||||
String ticket = messages.sendAsync(TARGET, "long task");
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
|
||||
"control: the async ticket must be PENDING before the release listener runs");
|
||||
|
||||
replyInbox.own(TARGET);
|
||||
replyInbox.publish(TARGET, "msg-1", "hello");
|
||||
assertEquals(1, replyInbox.peek(TARGET).size(),
|
||||
"control: the reply inbox must own TARGET and hold one message before the release "
|
||||
+ "listener runs");
|
||||
|
||||
primaryRegistry.recordDelegation(TARGET, "lead-1");
|
||||
assertEquals("lead-1", primaryRegistry.nudgeTargetFor(TARGET).orElse(null),
|
||||
"control: the delegation must be recorded before the release listener runs");
|
||||
|
||||
// --- the one call under test: invoke the REAL, assembled release listener directly, the
|
||||
// same way SessionManager.release(...) would on a real teardown.
|
||||
releaseListener.accept(new SessionManager.ReleaseDetail(TARGET, null, null, null, null));
|
||||
|
||||
MessageService.TaskView after = awaitTerminal(messages, ticket);
|
||||
assertEquals(MessageService.Phase.FAILED, after.phase(),
|
||||
"FleetdAssembly.java:447 must register a listener that calls messages.abandon(...) "
|
||||
+ "on the SAME assembled MessageService — an inert listener leaves this "
|
||||
+ "ticket PENDING for the full 30-minute async timeout");
|
||||
assertTrue(after.detail() != null && after.detail().contains("released"),
|
||||
"the abandon reason must say the worker session was released");
|
||||
|
||||
assertTrue(replyInbox.peek(TARGET).isEmpty(),
|
||||
"FleetdAssembly.java:447 must register a listener that calls replyInbox.release(...) "
|
||||
+ "— an inert listener leaves the inbox still owning TARGET with its message");
|
||||
|
||||
assertTrue(primaryRegistry.nudgeTargetFor(TARGET).isEmpty(),
|
||||
"FleetdAssembly.java:447 must register a listener that calls "
|
||||
+ "primaryRegistry.forgetDelegation(...) — an inert listener leaves the stale "
|
||||
+ "delegation in place");
|
||||
} finally {
|
||||
// Proof the teardown actually ran, not just an assurance that a finally was added: the
|
||||
// captured shutdown hook's close order (FleetdAssemblyLifecycleTest) closes the herdr
|
||||
// client last, so ports.herdr.closed flips to true only if this hook really executed.
|
||||
ports.shutdownHook.run();
|
||||
assertTrue(ports.herdr.closed, "the captured shutdown hook must have run and closed herdr — "
|
||||
+ "proof this test's assembled background loops/scheduler were torn down");
|
||||
}
|
||||
}
|
||||
|
||||
private static MessageService.TaskView awaitTerminal(MessageService messages, String ticket)
|
||||
throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 5000;
|
||||
MessageService.TaskView view = messages.poll(ticket);
|
||||
while (view.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
view = messages.poll(ticket);
|
||||
}
|
||||
return view;
|
||||
}
|
||||
}
|
||||
+241
@@ -0,0 +1,241 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.lead.LeadContextGauge;
|
||||
import dev.ltms.fleet.msg.LeadHeartbeatLoop;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.lang.reflect.Method;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #630, and fleetd #612's ranks for this call site. {@code FleetdAssembly.java:402}
|
||||
* computes {@code requireOperatorConfirm} from the effective {@code leadRollover.requireOperatorConfirm}
|
||||
* config, and {@code :409} threads it as the 14th argument into the full {@link LeadHeartbeatLoop}
|
||||
* constructor. Measured on 26f1986: dropping that one argument so the 13-argument overload is
|
||||
* selected instead (it delegates with {@code true} hardcoded — see that overload's own javadoc,
|
||||
* fleetd #621) compiles with 0 errors and leaves all 1883 tests green, both with and without the
|
||||
* argument. In production this means the daemon keeps starting and keeps nudging, but the
|
||||
* context-high notice silently goes back to telling EVERY lead to ask the operator before a
|
||||
* context roll — on a host that set {@code requireOperatorConfirm: false} specifically so it would
|
||||
* not have to. That is the operator's own fix silently reverting, with a fully green suite.
|
||||
*
|
||||
* <p>{@code LeadHeartbeatLoopTest} already proves {@link LeadHeartbeatLoop}'s package-private
|
||||
* {@code contextNotice(boolean, LeadContextGauge.Reading, boolean, boolean)} branches correctly on
|
||||
* its own {@code requireOperatorConfirm} argument — that the METHOD works. It says nothing about
|
||||
* which value {@code FleetdAssembly} actually passes into the constructed loop, so it is not
|
||||
* reused here as coverage for the call site.
|
||||
*
|
||||
* <p>This test assembles the real daemon TWICE — once with {@code leadRollover.requireOperatorConfirm:
|
||||
* false}, once with {@code true} — pulls the REAL {@code requireOperatorConfirm} field out of the
|
||||
* REAL, assembled {@link LeadHeartbeatLoop} each time (reflection: the field, and {@code
|
||||
* contextNotice} itself, are package-private to {@code dev.ltms.fleet.msg}, and nothing public
|
||||
* exposes either — the same technique {@code StatusPollerResilienceTest} already uses in this
|
||||
* suite), and calls the REAL {@code contextNotice} method with that field's value to produce the
|
||||
* actual notice text the assembled loop would append to a nudge. Both directions are asserted: a
|
||||
* one-directional test here would pass on a constant.
|
||||
*/
|
||||
class FleetdAssemblyRequireOperatorConfirmBehaviouralTest {
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> {
|
||||
throw new UnsupportedOperationException("no broker: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException("no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir, boolean requireOperatorConfirm) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
leadHeartbeat:
|
||||
idleAfterSeconds: 600
|
||||
backoffMs: 15000
|
||||
quietNudgeCap: 5
|
||||
leadRollover:
|
||||
handoverPath: handover.md
|
||||
requireOperatorConfirm: %s
|
||||
""".formatted(requireOperatorConfirm));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
/** Carries both the assembled loop under test AND its {@link RecordingResourcePorts}, so the
|
||||
* caller can tear the assembly down (this test assembles the real daemon TWICE — see the class
|
||||
* javadoc — and each assembly needs its own teardown, not just the last one). */
|
||||
private record Assembled(LeadHeartbeatLoop heartbeat, RecordingResourcePorts ports) {
|
||||
}
|
||||
|
||||
private static Assembled assembleHeartbeat(Path dir, boolean requireOperatorConfirm) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir, requireOperatorConfirm);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
assertNotNull(ports.shutdownHook, "FleetdAssembly must have registered a shutdown hook");
|
||||
|
||||
LeadHeartbeatLoop heartbeat = runtime.heartbeat();
|
||||
assertNotNull(heartbeat, "control: leadHeartbeat: is configured, so FleetdAssembly.assembleAndStart "
|
||||
+ "must have built a real LeadHeartbeatLoop");
|
||||
return new Assembled(heartbeat, ports);
|
||||
}
|
||||
|
||||
/** Pulls the REAL {@code requireOperatorConfirm} field off the REAL, assembled loop. */
|
||||
private static boolean assembledRequireOperatorConfirm(LeadHeartbeatLoop heartbeat) throws Exception {
|
||||
Field field = LeadHeartbeatLoop.class.getDeclaredField("requireOperatorConfirm");
|
||||
field.setAccessible(true);
|
||||
return field.getBoolean(heartbeat);
|
||||
}
|
||||
|
||||
/** Calls the REAL, package-private {@code contextNotice(boolean, Reading, boolean, boolean)} via reflection. */
|
||||
private static String contextNotice(boolean enabled, LeadContextGauge.Reading reading, boolean alreadyNotified,
|
||||
boolean requireOperatorConfirm) throws Exception {
|
||||
Method method = LeadHeartbeatLoop.class.getDeclaredMethod("contextNotice", boolean.class,
|
||||
LeadContextGauge.Reading.class, boolean.class, boolean.class);
|
||||
method.setAccessible(true);
|
||||
return (String) method.invoke(null, enabled, reading, alreadyNotified, requireOperatorConfirm);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled LeadHeartbeatLoop's context-high notice tracks "
|
||||
+ "leadRollover.requireOperatorConfirm — BOTH directions")
|
||||
void assembledRequireOperatorConfirmControlsNoticeWording(@TempDir Path dir) throws Exception {
|
||||
LeadContextGauge.Reading highReading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH,
|
||||
250_000L, 2);
|
||||
|
||||
String noticeFalse;
|
||||
String noticeTrue;
|
||||
|
||||
// --- direction 1: requireOperatorConfirm: false -----------------------------------------
|
||||
Path falseDir = dir.resolve("false");
|
||||
Files.createDirectories(falseDir);
|
||||
Assembled assembledFalse = assembleHeartbeat(falseDir, false);
|
||||
// Surefire runs the whole suite in one JVM fork, so each assembly's scheduler/loops must be
|
||||
// torn down here, on the failure path too — hence try/finally per assembly (this test
|
||||
// assembles TWICE, so both need their own teardown, not just the last one).
|
||||
try {
|
||||
boolean fieldFalse = assembledRequireOperatorConfirm(assembledFalse.heartbeat());
|
||||
assertFalse(fieldFalse, "FleetdAssembly.java:402/:409 must thread leadRollover."
|
||||
+ "requireOperatorConfirm: false into the assembled LeadHeartbeatLoop's own field — "
|
||||
+ "dropping the 14th constructor argument selects the 13-argument overload, which "
|
||||
+ "hardcodes true regardless of config (fleetd #621), and this would read true instead");
|
||||
|
||||
noticeFalse = contextNotice(true, highReading, false, fieldFalse);
|
||||
assertTrue(noticeFalse.contains("Decide for yourself when to confirm"),
|
||||
"with requireOperatorConfirm: false, the assembled loop's own notice must tell the "
|
||||
+ "lead it can decide for itself — got: " + noticeFalse);
|
||||
assertFalse(noticeFalse.contains("ask the operator") || noticeFalse.contains("Only the operator"),
|
||||
"with requireOperatorConfirm: false, the assembled loop's own notice must NOT ask the "
|
||||
+ "operator — got: " + noticeFalse);
|
||||
} finally {
|
||||
// Proof the teardown actually ran, not just an assurance that a finally was added: the
|
||||
// captured shutdown hook's close order (FleetdAssemblyLifecycleTest) closes the herdr
|
||||
// client last, so ports.herdr.closed flips to true only if this hook really executed.
|
||||
assembledFalse.ports().shutdownHook.run();
|
||||
assertTrue(assembledFalse.ports().herdr.closed, "the captured shutdown hook must have run "
|
||||
+ "and closed herdr — proof this assembly's background loops/scheduler were torn down");
|
||||
}
|
||||
|
||||
// --- direction 2: requireOperatorConfirm: true -------------------------------------------
|
||||
Path trueDir = dir.resolve("true");
|
||||
Files.createDirectories(trueDir);
|
||||
Assembled assembledTrue = assembleHeartbeat(trueDir, true);
|
||||
try {
|
||||
boolean fieldTrue = assembledRequireOperatorConfirm(assembledTrue.heartbeat());
|
||||
assertTrue(fieldTrue, "FleetdAssembly.java:402/:409 must thread leadRollover."
|
||||
+ "requireOperatorConfirm: true into the assembled LeadHeartbeatLoop's own field");
|
||||
|
||||
noticeTrue = contextNotice(true, highReading, false, fieldTrue);
|
||||
assertTrue(noticeTrue.contains("ask the operator") && noticeTrue.contains("Only the operator can approve the roll"),
|
||||
"with requireOperatorConfirm: true, the assembled loop's own notice must ask the "
|
||||
+ "operator — got: " + noticeTrue);
|
||||
assertFalse(noticeTrue.contains("Decide for yourself when to confirm"),
|
||||
"with requireOperatorConfirm: true, the assembled loop's own notice must NOT tell "
|
||||
+ "the lead it can decide for itself — got: " + noticeTrue);
|
||||
} finally {
|
||||
assembledTrue.ports().shutdownHook.run();
|
||||
assertTrue(assembledTrue.ports().herdr.closed, "the captured shutdown hook must have run "
|
||||
+ "and closed herdr — proof this assembly's background loops/scheduler were torn down");
|
||||
}
|
||||
|
||||
// --- the two directions must actually differ: a constant return would pass both assertion
|
||||
// blocks above vacuously if they happened to share wording, so compare them directly too.
|
||||
assertTrue(!noticeFalse.equals(noticeTrue),
|
||||
"the two directions must produce genuinely different notice text — got the same "
|
||||
+ "text for both: " + noticeFalse);
|
||||
}
|
||||
}
|
||||
@@ -98,6 +98,14 @@ class FleetdAssemblyRoleFallbackBoundaryTest {
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// No real HTTP bind in a unit test.
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private static final class SentinelReplyInbox implements ReplyInbox {
|
||||
|
||||
@@ -0,0 +1,207 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import io.javalin.Javalin;
|
||||
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 java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.OptionalLong;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — replaces {@code FleetdBackendQuarantineWiringTest} (fleetd #466), a source-text
|
||||
* test that scraped {@code Fleetd.java} (now {@code FleetdAssembly.java}, moved there by fleetd #612
|
||||
* Unit A) for the {@code BackendQuarantine.withEscalation(...)} call, and separately asserted the
|
||||
* flat two-argument constructor's text was ABSENT. That proves the right method NAME appears in
|
||||
* source; it proves nothing about what the constructed object actually DOES.
|
||||
*
|
||||
* <p>This test instead drives the REAL {@link BackendQuarantine} the real {@link
|
||||
* FleetdAssembly#assembleAndStart} builds — reached through {@link
|
||||
* dev.ltms.fleet.mcp.FleetMcp#quarantineSource()} on the real, live {@code FleetMcp} {@code
|
||||
* FleetdRuntime} owns — and asserts the ONE behavioural difference {@code withEscalation} and the
|
||||
* flat constructor actually produce (see {@link BackendQuarantine}'s own class doc, "Mechanism"):
|
||||
* quarantining the same credential twice in a row, within one base cooldown of the first deadline,
|
||||
* must escalate the second cooldown past the first. A flat instance reports the identical cooldown
|
||||
* both times.
|
||||
*/
|
||||
class FleetdBackendQuarantineAssemblyTest {
|
||||
|
||||
/** Base cooldown used throughout — long enough that rounding never blurs the 2x escalation. */
|
||||
private static final int COOLDOWN_SECONDS = 100;
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
final AtomicLong clockNanos = new AtomicLong(0L);
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
// Controllable: the SAME LongSupplier instance BackendQuarantine.withEscalation(...) is
|
||||
// built with, so advancing clockNanos after assembly moves the quarantine tracker's own
|
||||
// clock, with no real sleep needed to observe escalation.
|
||||
return clockNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return clockNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
quarantineCooldownSeconds: %d
|
||||
""".formatted(COOLDOWN_SECONDS));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled BackendQuarantine escalates a repeated exhaustion, "
|
||||
+ "which the flat two-argument constructor can never do")
|
||||
void assembledQuarantineEscalatesOnARepeatedExhaustion(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
|
||||
BackendQuarantine quarantine = runtime.mcp().quarantineSource().quarantine();
|
||||
|
||||
// First exhaustion, at clock=0: a fresh occurrence, blocked for exactly the base cooldown.
|
||||
quarantine.quarantine("cred-x");
|
||||
BackendQuarantine.Status first = quarantine.status("cred-x").orElseThrow(
|
||||
() -> new AssertionError("credential must be quarantined immediately after quarantine()"));
|
||||
assertEquals(1, first.repeatCount(), "the first call is repeat #1");
|
||||
assertEquals(COOLDOWN_SECONDS, first.remainingSeconds(),
|
||||
"a fresh quarantine blocks for exactly the base cooldown");
|
||||
|
||||
// Second exhaustion, arriving just after the first deadline — well within one base cooldown
|
||||
// of it, so this is a CONTINUATION of the same streak (repeat #2), not a fresh occurrence.
|
||||
long firstDeadlineNanos = COOLDOWN_SECONDS * 1_000_000_000L;
|
||||
ports.clockNanos.set(firstDeadlineNanos + 1);
|
||||
quarantine.quarantine("cred-x");
|
||||
BackendQuarantine.Status second = quarantine.status("cred-x").orElseThrow(
|
||||
() -> new AssertionError("credential must be quarantined immediately after the second "
|
||||
+ "quarantine() call"));
|
||||
assertEquals(2, second.repeatCount(), "the second call, arriving within one base cooldown of "
|
||||
+ "the first deadline, continues the streak as repeat #2");
|
||||
|
||||
// The one behavioural difference: withEscalation doubles the cooldown on repeat #2 (capped
|
||||
// well above this at 12x base), the flat two-argument constructor never grows past the base
|
||||
// cooldown no matter how many times quarantine() is called in a row.
|
||||
assertEquals(2 * COOLDOWN_SECONDS, second.remainingSeconds(),
|
||||
"withEscalation's default backoff doubles the cooldown on the second consecutive "
|
||||
+ "exhaustion — this is the exact call FleetdAssembly.java makes at the "
|
||||
+ "BackendQuarantine.withEscalation(...) call site");
|
||||
assertTrue(second.remainingSeconds() > first.remainingSeconds(),
|
||||
"the flat two-argument BackendQuarantine constructor would report the SAME remaining "
|
||||
+ "seconds both times — this inequality is what a mutation to the flat "
|
||||
+ "constructor at that call site must fail");
|
||||
|
||||
// Also confirm isQuarantined/remainingSeconds agree, exercising the accessors a real caller
|
||||
// (fleet_profiles / fleet_list, per BackendQuarantine's own class doc) actually reads.
|
||||
assertTrue(quarantine.isQuarantined("cred-x"));
|
||||
OptionalLong remaining = quarantine.remainingSeconds("cred-x");
|
||||
assertTrue(remaining.isPresent());
|
||||
assertEquals(2 * COOLDOWN_SECONDS, remaining.getAsLong());
|
||||
}
|
||||
}
|
||||
@@ -1,65 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #466 follow-up: {@code Fleetd.main} builds the daemon's one {@code BackendQuarantine}
|
||||
* from {@link dev.ltms.fleet.placement.BackendQuarantine#withEscalation(java.util.function.LongSupplier,
|
||||
* long)} — the escalating factory — rather than the plain two-argument constructor, which is still a
|
||||
* flat cooldown (kept for backward compatibility, see that class's doc). {@code
|
||||
* BackendQuarantineTest} proves {@code withEscalation} itself escalates, is ceilinged, and resets;
|
||||
* it says nothing about which one {@code main} actually calls.
|
||||
*
|
||||
* <p>Measured directly: reverting {@code main} to {@code new BackendQuarantine(System::nanoTime,
|
||||
* TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()))} — the pre-#466 flat call — compiles
|
||||
* with 0 errors and leaves the entire 1608-test suite (including every {@code BackendQuarantineTest}
|
||||
* case) green, because no other test constructs its {@code BackendQuarantine} through {@code main};
|
||||
* every one of them builds its own instance directly. That silent regression is exactly the shape
|
||||
* {@link FleetdLeadSeatWiringTest} and {@link FleetdCompletionResolverWiringTest} already guard
|
||||
* against for their own constructor arguments — this is the same class of gap for fleetd #466's
|
||||
* factory choice, following their approach.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
|
||||
* BackendQuarantine} and never runs {@code main} — a green result here proves only that the exact
|
||||
* text {@code main} calls {@code BackendQuarantine.withEscalation(...)} rather than the flat
|
||||
* constructor. It does not prove that call actually executes at startup (no test here starts the
|
||||
* daemon), and it does not prove the escalation reaches a real backend or credential — only
|
||||
* {@code BackendQuarantineTest} proves the factory's own behaviour, and only a live daemon proves
|
||||
* the wiring runs.
|
||||
*/
|
||||
class FleetdBackendQuarantineWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] main's BackendQuarantine local is still built from BackendQuarantine.withEscalation(...)")
|
||||
void mainStillWiresTheEscalatingQuarantineFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"BackendQuarantine quarantine = BackendQuarantine.withEscalation(System::nanoTime,\n"
|
||||
+ " TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));"),
|
||||
"Fleetd.main's BackendQuarantine local must still be built from "
|
||||
+ "BackendQuarantine.withEscalation(System::nanoTime, "
|
||||
+ "TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds())). Reverting to the flat "
|
||||
+ "two-argument constructor (fleetd #466's measured regression) compiles with 0 errors "
|
||||
+ "and leaves the whole suite green, including every BackendQuarantineTest case that "
|
||||
+ "proves the escalation itself works — this source check is what must go red instead. "
|
||||
+ "A reverted daemon would go back to retrying a weekly subscription limit on every "
|
||||
+ "flat ~30-minute cooldown, about 336 times across the week.");
|
||||
|
||||
// Negative form of the same check: the pre-#466 flat call, if it ever reappears at this
|
||||
// declaration, must not be mistaken for the escalating one by a looser positive-only check.
|
||||
assertFalse(source.contains(
|
||||
"BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,\n"
|
||||
+ " TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));"),
|
||||
"main's BackendQuarantine local must never regress to the flat two-argument constructor");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,315 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import dev.ltms.fleet.placement.PlacementException;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.WorktreeRequest;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 step 2, Unit B1 — replaces {@code FleetdCompletionResolverWiringTest} (deleted in
|
||||
* this same commit), whose four tests read {@code Fleetd.java}'s source text and asserted the
|
||||
* {@code CompletionResolver} construction call still named the right arguments. That proved the
|
||||
* call site's spelling, never that the assembled resolver actually behaves differently when an
|
||||
* argument is dropped.
|
||||
*
|
||||
* <p>These tests drive {@link FleetdAssembly#assembleAndStart} — the real boot composition,
|
||||
* fleetd #612 Unit A — and read {@link FleetdRuntime#completion()}: the exact {@link
|
||||
* CompletionResolver} instance the assembled daemon uses, never a copy built alongside it for the
|
||||
* test's benefit. Two behaviours are pinned, matching the ticket's own two measured mutations:
|
||||
*
|
||||
* <ul>
|
||||
* <li>the 8th constructor argument ({@code Fleetd.worktreeBranchLookup(sessions::roster)}) —
|
||||
* {@link #assembledResolverReportsTheMembersWorktreeAndBranchInAFallbackReport}; and</li>
|
||||
* <li>the 5th/6th arguments ({@code backendErrorPatterns}, {@code backendErrorSink}, both
|
||||
* assigned from {@code Fleetd}'s extracted factories rather than an inline lambda) —
|
||||
* {@link #assembledResolverClassifiesAndCoolsOffOnAConfiguredBackendErrorPattern}.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>Both tests bypass {@link dev.ltms.fleet.inject.StatusPoller} and drive {@link
|
||||
* CompletionResolver#onDelivered} / {@link CompletionResolver#resolveBeforePostAction} directly —
|
||||
* the same public, synchronous entry points {@code CompletionResolverTest} uses — with a
|
||||
* hand-built {@link CompletableFuture} waiter, so no real poller loop or herdr status poll is
|
||||
* needed. The pane scrape comes from {@link FakeHerdr#readText}; the elapsed-time floor
|
||||
* ({@code CompletionResolver.MIN_TURN_NANOS}) is controlled via a fake, advanceable {@link
|
||||
* ResourcePorts#nanoClock()} rather than a real sleep.
|
||||
*/
|
||||
class FleetdCompletionResolverAssemblyTest {
|
||||
|
||||
/** Same shape as {@code FleetdAssemblyLifecycleTest}'s fake, plus a nanoClock this test can advance. */
|
||||
private static final class ControllableResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr;
|
||||
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); // arbitrary non-zero start
|
||||
Runnable shutdownHook;
|
||||
|
||||
ControllableResourcePorts(FakeHerdr herdr) {
|
||||
this.herdr = herdr;
|
||||
}
|
||||
|
||||
void advanceSeconds(long seconds) {
|
||||
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
// Never invoked: this test's config has no `broker:` block, so Fleetd.selectReplyInbox
|
||||
// returns the in-memory inbox before calling the opener at all.
|
||||
return (uri, prefetch) -> {
|
||||
throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
// Never invoked: no `coordinator:` block configured either.
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException("leadMailboxOpener must not be called — no coordinator: block");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return nowNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return nowNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Deliberately never bind a real port.
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir, String profilesYaml, String extraGuardHost,
|
||||
String worktreeRootYamlLine) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
%s
|
||||
profiles:
|
||||
%s
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- %s
|
||||
""".formatted(worktreeRootYamlLine == null ? "" : worktreeRootYamlLine, profilesYaml, extraGuardHost));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
private static void gitQuiet(Path cwd, String... args) throws Exception {
|
||||
List<String> cmd = new java.util.ArrayList<>(List.of("git"));
|
||||
cmd.addAll(List.of(args));
|
||||
Process p = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true).start();
|
||||
String out = new String(p.getInputStream().readAllBytes());
|
||||
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: git " + String.join(" ", args));
|
||||
assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out);
|
||||
}
|
||||
|
||||
private static Path initRepo(Path dir) throws Exception {
|
||||
Files.createDirectories(dir);
|
||||
gitQuiet(dir, "init", "-q", "-b", "main");
|
||||
gitQuiet(dir, "config", "user.email", "test@example.invalid");
|
||||
gitQuiet(dir, "config", "user.name", "Test");
|
||||
Files.writeString(dir.resolve("README.md"), "seed\n");
|
||||
gitQuiet(dir, "add", "README.md");
|
||||
gitQuiet(dir, "commit", "-q", "-m", "seed");
|
||||
return dir;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248's first measured mutation: replacing {@code CompletionResolver}'s 8th constructor
|
||||
* argument with the inert {@code _ -> null} compiles clean and leaves every existing test green
|
||||
* — it silently drops fleetd #241's fallback-report location. This drives the real assembled
|
||||
* resolver through a member echoing its own injected brief back (no {@code fleet_reply}), which
|
||||
* resolves via {@code noReportMessage(target)}, and proves the real member's {@code branch} —
|
||||
* only obtainable via {@code Fleetd.worktreeBranchLookup(sessions::roster)} reading the real,
|
||||
* worktree-provisioned {@link MemberSession} — appears in the reported text.
|
||||
*/
|
||||
@Test
|
||||
void assembledResolverReportsTheMembersWorktreeAndBranchInAFallbackReport(@TempDir Path dir) throws Exception {
|
||||
Path repo = initRepo(dir.resolve("repo"));
|
||||
FleetConfig cfg = writeConfig(dir, """
|
||||
wtprofile:
|
||||
baseUrl: http://wthost.local:8000
|
||||
model: sonnet
|
||||
""", "wthost.local", "worktreeRoot: " + dir.resolve("wts"));
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr());
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
try {
|
||||
MemberSession session = runtime.sessions().acquire("wtprofile", repo.toString(), repo.toString(),
|
||||
null, new WorktreeRequest("fleetd-612-b1", null));
|
||||
String target = session.terminalId();
|
||||
String branch = session.branch();
|
||||
assertTrue(branch != null && branch.startsWith("worker/"),
|
||||
"sanity: a worktree-provisioned session must carry a real branch, got: " + branch);
|
||||
|
||||
CompletionResolver completion = runtime.completion();
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = new CompletableFuture<>();
|
||||
String echoedBrief = "z".repeat(450); // >= CompletionResolver.ECHO_MIN_CHARS normalised chars
|
||||
|
||||
ports.herdr.readText("idle, nothing yet");
|
||||
completion.onDelivered(target, new TurnToken(target, waiter, echoedBrief));
|
||||
|
||||
ports.herdr.readText(echoedBrief); // the pane just echoes the injected brief back — no real report
|
||||
ports.advanceSeconds(3); // clear CompletionResolver.MIN_TURN_NANOS (2s) without a real sleep
|
||||
completion.resolveBeforePostAction(target);
|
||||
|
||||
Rendezvous.Resolution resolution = waiter.getNow(null);
|
||||
assertTrue(resolution != null, "the waiter must have resolved synchronously");
|
||||
assertEquals(Rendezvous.Kind.COMPLETION, resolution.kind());
|
||||
assertTrue(resolution.text().contains(CompletionResolver.NO_REPORT_PREFIX),
|
||||
"sanity: must have gone down the noReportMessage sub-path: " + resolution.text());
|
||||
assertTrue(resolution.text().contains("branch=" + branch),
|
||||
"the assembled resolver must report the member's real branch (fleetd #241 via "
|
||||
+ "fleetd #248's worktreeBranchLookup wiring); got: " + resolution.text());
|
||||
} finally {
|
||||
if (ports.shutdownHook != null) ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248's second measured mutation, and fleetd#201 Unit 5's own gap: replacing {@code
|
||||
* backendErrorPatterns}/{@code backendErrorSink} with {@code BackendErrorPatternLookup.legacy()}
|
||||
* / {@code BackendErrorSink.none()} compiles clean and leaves every existing behavioural test
|
||||
* green.
|
||||
*
|
||||
* <p>Classification proof: this test's profile configures {@code errorPattern: "credential
|
||||
* outage"} — text the built-in {@code (?i)\bAPI Error\s*:} fallback ({@code legacy()}'s only
|
||||
* behaviour) never matches. So a real {@code Fleetd.backendErrorPatternLookup(...)} wiring
|
||||
* classifies the send as {@code FAILED}; {@code legacy()} would fall through to the plain
|
||||
* completion path instead ({@code Kind.COMPLETION}).
|
||||
*
|
||||
* <p>Cool-off proof: two distinct targets on the same profile/credential each classified as a
|
||||
* backend error inside the 60s window must cool the credential off ({@link
|
||||
* dev.ltms.fleet.placement.BackendOutagePolicy}, fleetd#201 Unit 5) — observable two ways: (1)
|
||||
* the real {@code Fleetd.backendErrorSink(...)} marks each session {@code BACKEND_ERROR} (only
|
||||
* the real sink calls {@code sessions.onBackendError}; {@code BackendErrorSink.none()} never
|
||||
* does), and (2) a third explicit-profile spawn attempt is refused with a {@link
|
||||
* PlacementException} naming the cool-off — only reachable because the real sink's {@code
|
||||
* outagePolicy.record(...)} call actually ran.
|
||||
*/
|
||||
@Test
|
||||
void assembledResolverClassifiesAndCoolsOffOnAConfiguredBackendErrorPattern(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir, """
|
||||
coolprofile:
|
||||
baseUrl: http://coolhost.local:8000
|
||||
model: sonnet
|
||||
errorPattern: "credential outage"
|
||||
""", "coolhost.local", null);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr());
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
try {
|
||||
MemberSession session1 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null);
|
||||
MemberSession session2 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null);
|
||||
String target1 = session1.terminalId();
|
||||
String target2 = session2.terminalId();
|
||||
assertTrue(!target1.equals(target2), "sanity: the two spawns must be distinct targets");
|
||||
|
||||
CompletionResolver completion = runtime.completion();
|
||||
|
||||
CompletableFuture<Rendezvous.Resolution> waiter1 = new CompletableFuture<>();
|
||||
ports.herdr.readText("idle 1");
|
||||
completion.onDelivered(target1, new TurnToken(target1, waiter1, null));
|
||||
ports.herdr.readText("credential outage: upstream 503");
|
||||
ports.advanceSeconds(3);
|
||||
completion.resolveBeforePostAction(target1);
|
||||
Rendezvous.Resolution resolution1 = waiter1.getNow(null);
|
||||
assertTrue(resolution1 != null, "target1's waiter must have resolved synchronously");
|
||||
assertEquals(Rendezvous.Kind.FAILED, resolution1.kind(),
|
||||
"a configured errorPattern the built-in fallback never matches must classify as "
|
||||
+ "a backend error, not a plain completion; got: " + resolution1);
|
||||
assertTrue(resolution1.text().contains("credential outage: upstream 503"), resolution1.text());
|
||||
|
||||
CompletableFuture<Rendezvous.Resolution> waiter2 = new CompletableFuture<>();
|
||||
ports.herdr.readText("idle 2");
|
||||
completion.onDelivered(target2, new TurnToken(target2, waiter2, null));
|
||||
ports.herdr.readText("credential outage: upstream 503 again");
|
||||
ports.advanceSeconds(3);
|
||||
completion.resolveBeforePostAction(target2);
|
||||
Rendezvous.Resolution resolution2 = waiter2.getNow(null);
|
||||
assertTrue(resolution2 != null, "target2's waiter must have resolved synchronously");
|
||||
assertEquals(Rendezvous.Kind.FAILED, resolution2.kind());
|
||||
|
||||
List<MemberSession> roster = runtime.sessions().roster();
|
||||
assertTrue(roster.stream().anyMatch(s -> target1.equals(s.terminalId())
|
||||
&& s.state() == MemberSession.State.BACKEND_ERROR),
|
||||
"the real backendErrorSink must have transitioned target1 to BACKEND_ERROR: " + roster);
|
||||
assertTrue(roster.stream().anyMatch(s -> target2.equals(s.terminalId())
|
||||
&& s.state() == MemberSession.State.BACKEND_ERROR),
|
||||
"the real backendErrorSink must have transitioned target2 to BACKEND_ERROR: " + roster);
|
||||
|
||||
PlacementException coolOff = assertThrows(PlacementException.class,
|
||||
() -> runtime.sessions().acquire("coolprofile", null, dir.toString(), null),
|
||||
"two distinct targets classified within the 60s window must cool the credential "
|
||||
+ "off (BackendOutagePolicy), refusing a third explicit-profile spawn");
|
||||
assertTrue(coolOff.getMessage().contains("cooling off"), coolOff.getMessage());
|
||||
} finally {
|
||||
if (ports.shutdownHook != null) ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,89 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #248: this is the test that was actually missing. {@code Fleetd.main} builds its {@code
|
||||
* CompletionResolver} from an 8-argument constructor, and the ticket's own measurement proved two
|
||||
* ways to silently unwire it — both compiled with 0 errors and left every existing test green:
|
||||
*
|
||||
* <ul>
|
||||
* <li>replacing the worktree/branch argument (the 8th) with {@code _ -> null} — drops
|
||||
* fleetd#241's fallback-report location entirely;</li>
|
||||
* <li>replacing {@code backendErrorPatterns, backendErrorSink} (5th/6th) with {@code
|
||||
* BackendErrorPatternLookup.legacy(), BackendErrorSink.none()} — drops fleetd#201 Unit 5's
|
||||
* backend-error classification and cool-off entirely.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>Neither mutation could be caught by any test that constructs its own {@code
|
||||
* CompletionResolver} (every test before this one did exactly that) or by a test of {@link
|
||||
* Fleetd#worktreeBranchLookup}, {@link Fleetd#backendErrorPatternLookup}, or {@link
|
||||
* Fleetd#backendErrorSink} in isolation (see {@code FleetdWorktreeBranchLookupTest}, {@code
|
||||
* FleetdBackendErrorPatternLookupTest}, {@code FleetdBackendErrorSinkTest}) — those prove the
|
||||
* factories work, never that {@code main} still calls them. This class is a plain source-text
|
||||
* assertion on {@code Fleetd.java} — crude, but honest about what it checks, and it turns red the
|
||||
* instant the wiring is dropped, mirroring the same fallback shape {@link
|
||||
* FleetdFleetAppConstructionTest} already uses for a different constructor argument.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
|
||||
* CompletionResolver} and never runs {@code main}.
|
||||
*/
|
||||
class FleetdCompletionResolverWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] CompletionResolver's construction call still names backendErrorPatterns and backendErrorSink")
|
||||
void backendErrorArgumentsAreStillNamedAtTheCallSite() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"exhaustionSink, backendErrorPatterns, backendErrorSink, System::nanoTime,"),
|
||||
"CompletionResolver's construction call must still pass backendErrorPatterns and "
|
||||
+ "backendErrorSink as its 5th/6th arguments. Replacing them with "
|
||||
+ "BackendErrorPatternLookup.legacy()/BackendErrorSink.none() (fleetd #248's measured "
|
||||
+ "mutation) compiles with 0 errors and leaves every behavioural test green — this "
|
||||
+ "source check is what must go red instead.");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] CompletionResolver's construction call still passes worktreeBranchLookup(sessions::roster)")
|
||||
void worktreeBranchLookupIsStillPassedAtTheCallSite() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains("worktreeBranchLookup(sessions::roster)"),
|
||||
"CompletionResolver's construction call must still pass worktreeBranchLookup(sessions::roster) "
|
||||
+ "as its 8th (last) argument. Replacing it with the inert `_ -> null` (fleetd #248's "
|
||||
+ "other measured mutation) compiles with 0 errors and leaves every behavioural test "
|
||||
+ "green — this source check is what must go red instead.");
|
||||
assertFalse(source.contains("System::nanoTime,\n _ -> null"),
|
||||
"the worktree/branch argument must never regress to the inert `_ -> null` literal");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] backendErrorPatterns is assigned from the extracted backendErrorPatternLookup(...) factory")
|
||||
void backendErrorPatternsComesFromTheFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"BackendErrorPatternLookup backendErrorPatterns = backendErrorPatternLookup(sessions::roster,"),
|
||||
"backendErrorPatterns must be assigned from Fleetd.backendErrorPatternLookup(...), not an "
|
||||
+ "inline lambda that a source check on the CompletionResolver call alone cannot see "
|
||||
+ "through");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] backendErrorSink is assigned from the extracted backendErrorSink(...) factory")
|
||||
void backendErrorSinkComesFromTheFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"BackendErrorSink backendErrorSink = backendErrorSink(sessions, () -> config.get().profiles(),"),
|
||||
"backendErrorSink must be assigned from Fleetd.backendErrorSink(...), not an inline lambda "
|
||||
+ "that a source check on the CompletionResolver call alone cannot see through");
|
||||
}
|
||||
}
|
||||
@@ -1,31 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-185: {@code ConnectionIdentity} must resolve a caller's pane on EITHER herdr daemon (a
|
||||
* lead's MCP connection resolves against the lead daemon; a member's against the member daemon).
|
||||
* Pinning {@code PaneLocator} to {@code memberHerdr} alone — the bug this guards against — leaves
|
||||
* every lead's own connection unresolvable ({@code callerTerminal == null}) the moment
|
||||
* {@code memberHerdrSocket} names a second daemon, which breaks {@code fleet_reply}/{@code
|
||||
* fleet_ask} and {@code fleet_whoami} for a lead. A unit test on {@link
|
||||
* dev.ltms.fleet.herdr.PaneLocator} alone (see {@code PaneLocatorTest}) proves the class CAN
|
||||
* search two clients, but not that {@code Fleetd.main} actually wires it that way — hence this
|
||||
* source-level assertion, the same technique {@code FleetdHerdrControlConstructionTest} uses.
|
||||
*/
|
||||
class FleetdConnectionIdentityConstructionTest {
|
||||
@Test
|
||||
void connectionIdentitySearchesBothDaemonsNotJustTheMemberOne() throws Exception {
|
||||
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
assertFalse(source.contains("new PaneLocator(memberHerdr)"),
|
||||
"PaneLocator must not be pinned to the member daemon alone — a lead's own "
|
||||
+ "connection resolves against the LEAD daemon and would never be found");
|
||||
assertTrue(source.contains("new PaneLocator(herdr, memberHerdr)"),
|
||||
"PaneLocator must search the lead daemon first, then the member daemon");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,218 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import io.javalin.Javalin;
|
||||
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 java.util.Map;
|
||||
import java.util.OptionalLong;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 step 4, ranks 1 and 2 (publish side) — {@link FleetdAssembly} lines
|
||||
* {@code liveExhaustedPatterns}/{@code exhaustedPatterns} (CB-578 stage A, the ticket's own "worst
|
||||
* consequence in the whole sweep": a genuine usage-limit refusal handed back to a waiting caller
|
||||
* AS REAL COMPLETED WORK) and {@code Fleetd.publishExhaustionSink(...)} (CB-578 stage B: the
|
||||
* credential that hit the limit is never quarantined). None of these three lines is driven by an
|
||||
* existing test through the real assembly: {@code FleetdExhaustedPatternLookupWiringTest} and
|
||||
* {@code FleetdLiveExhaustedPatternsWiringTest} (fleetd #589) call {@code Fleetd.liveExhaustedPatterns}
|
||||
* / {@code Fleetd.exhaustedPatternLookup} directly as factories, never through {@link
|
||||
* FleetdAssembly#assembleAndStart} — they prove the FACTORY classifies correctly, never that THIS
|
||||
* call site is the one that actually got wired into the running {@link CompletionResolver}. {@link
|
||||
* FleetdBackendQuarantineAssemblyTest} drives {@code BackendQuarantine.withEscalation(...)}
|
||||
* directly, a different call site from {@code publishExhaustionSink} here.
|
||||
*
|
||||
* <p>This test drives the REAL assembled {@link CompletionResolver} ({@link
|
||||
* FleetdRuntime#completion()}) with a profile carrying a configured {@code exhaustedPattern},
|
||||
* through a pane scrape that matches it, and asserts both halves of the production consequence:
|
||||
* (1) the resolution is {@link Rendezvous.Kind#BACKEND_EXHAUSTED}, never a plain completion handed
|
||||
* back as real work, and (2) the profile's credential is actually quarantined afterward, through
|
||||
* the REAL {@link BackendQuarantine} the same assembly built ({@link
|
||||
* FleetdRuntime#mcp()}{@code .quarantineSource().quarantine()}) — never a copy.
|
||||
*
|
||||
* <p>Same {@code ControllableResourcePorts} shape as {@code FleetdCompletionResolverAssemblyTest}:
|
||||
* a fake, advanceable {@code nanoClock} so {@code CompletionResolver.MIN_TURN_NANOS} clears without
|
||||
* a real sleep, and {@link FakeHerdr#readText} to drive the pane scrape.
|
||||
*/
|
||||
class FleetdExhaustedPatternAssemblyTest {
|
||||
|
||||
private static final class ControllableResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr;
|
||||
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); // arbitrary non-zero start
|
||||
Runnable shutdownHook;
|
||||
|
||||
ControllableResourcePorts(FakeHerdr herdr) {
|
||||
this.herdr = herdr;
|
||||
}
|
||||
|
||||
void advanceSeconds(long seconds) {
|
||||
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> {
|
||||
throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException("leadMailboxOpener must not be called — no coordinator: block");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return nowNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return nowNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Deliberately never bind a real port.
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir, int cooldownSeconds) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
quarantineCooldownSeconds: %d
|
||||
profiles:
|
||||
exhaustprofile:
|
||||
baseUrl: http://exhausthost.local:8000
|
||||
model: sonnet
|
||||
exhaustedPattern: "usage limit reached"
|
||||
guard:
|
||||
offSubscriptionHosts:
|
||||
- exhausthost.local
|
||||
""".formatted(cooldownSeconds));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589's own description of this gap ({@code Fleetd#exhaustedPatternLookup}'s javadoc):
|
||||
* "the worst consequence in the whole #589 sweep" — a genuine usage-limit refusal stops being
|
||||
* classified as {@code BACKEND_EXHAUSTED} and is handed back to a waiting {@code fleet_send} as
|
||||
* if it were real completed work. Pins {@code FleetdAssembly}'s {@code liveExhaustedPatterns}
|
||||
* AND {@code exhaustedPatterns} lines (rank 1) together with {@code publishExhaustionSink}
|
||||
* (rank 2, the non-OpenCode half) in one flow: classify, then quarantine.
|
||||
*/
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] a scrape matching the profile's exhaustedPattern resolves "
|
||||
+ "BACKEND_EXHAUSTED (never a plain completion) and quarantines the credential")
|
||||
void assembledResolverClassifiesExhaustionAndQuarantinesTheCredential(@TempDir Path dir) throws Exception {
|
||||
int cooldownSeconds = 120;
|
||||
FleetConfig cfg = writeConfig(dir, cooldownSeconds);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr());
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
try {
|
||||
MemberSession session = runtime.sessions().acquire("exhaustprofile", null, dir.toString(), null);
|
||||
String target = session.terminalId();
|
||||
|
||||
CompletionResolver completion = runtime.completion();
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = new CompletableFuture<>();
|
||||
|
||||
ports.herdr.readText("idle, nothing yet");
|
||||
completion.onDelivered(target, new TurnToken(target, waiter, null));
|
||||
// The matched text must START the pane line (CompletionResolver.startsWithExhaustion) —
|
||||
// no preceding sentence — for the quarantine side-effect to fire, same as production.
|
||||
ports.herdr.readText("usage limit reached: try again in a few hours");
|
||||
ports.advanceSeconds(3); // clear CompletionResolver.MIN_TURN_NANOS (2s), no real sleep
|
||||
completion.resolveBeforePostAction(target);
|
||||
|
||||
// CONTROL: the waiter must have resolved synchronously at all — if the assembled
|
||||
// CompletionResolver were never actually driven (e.g. a wiring break upstream silently
|
||||
// left the resolver unreachable), this fails loudly before the real assertions below
|
||||
// ever run, rather than passing on an untouched waiter.
|
||||
Rendezvous.Resolution resolution = waiter.getNow(null);
|
||||
assertTrue(resolution != null, "CONTROL: the waiter must have resolved synchronously — "
|
||||
+ "if this is null, the assembled resolver was never actually exercised");
|
||||
|
||||
assertEquals(Rendezvous.Kind.BACKEND_EXHAUSTED, resolution.kind(),
|
||||
"a scrape matching the profile's configured exhaustedPattern must classify as "
|
||||
+ "BACKEND_EXHAUSTED, not a plain completion handed back as real work — "
|
||||
+ "replacing FleetdAssembly's liveExhaustedPatterns/exhaustedPatterns "
|
||||
+ "lines with their inert forms (Map.of() / target -> null) must fail "
|
||||
+ "this assertion; got: " + resolution);
|
||||
assertTrue(resolution.text().contains("usage limit reached"), resolution.text());
|
||||
|
||||
BackendQuarantine quarantine = runtime.mcp().quarantineSource().quarantine();
|
||||
assertTrue(quarantine.isQuarantined("exhaustprofile"),
|
||||
"the real publishExhaustionSink-built sink must have quarantined the profile's "
|
||||
+ "credential (effectiveCredentialId() == the profile name here, no "
|
||||
+ "credentialId configured) — replacing FleetdAssembly's "
|
||||
+ "publishExhaustionSink call site with a hardcoded ExhaustionSink.none() "
|
||||
+ "must fail this assertion, since nothing would ever call "
|
||||
+ "quarantine.quarantine(...)");
|
||||
OptionalLong remaining = quarantine.remainingSeconds("exhaustprofile");
|
||||
assertTrue(remaining.isPresent() && remaining.getAsLong() > 0
|
||||
&& remaining.getAsLong() <= cooldownSeconds,
|
||||
"a fresh quarantine must block for at most the configured base cooldown: " + remaining);
|
||||
} finally {
|
||||
if (ports.shutdownHook != null) ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,29 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-185: {@code FleetApp} must be constructed with BOTH herdr clients (the lead's and the
|
||||
* member's), never the raw lead-only {@code herdr}. Passing only {@code herdr} — the bug this
|
||||
* guards against — makes {@code GET /healthz} green while the member daemon is down (so every
|
||||
* spawn fails invisibly) and silently drops every member workspace from {@code GET /sessions}.
|
||||
* A behavioural test on {@code FleetApp} alone (see {@code FleetAppTwoDaemonTest}) proves the
|
||||
* class merges/gates correctly when given two clients, but not that {@code Fleetd.main} actually
|
||||
* passes it two — hence this source-level assertion, mirroring
|
||||
* {@code FleetdHerdrControlConstructionTest}.
|
||||
*/
|
||||
class FleetdFleetAppConstructionTest {
|
||||
@Test
|
||||
void fleetAppIsConstructedWithBothHerdrDaemons() throws Exception {
|
||||
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
assertFalse(source.contains("new FleetApp(herdr, workers,"),
|
||||
"FleetApp must not be constructed with the lead-only herdr client");
|
||||
assertTrue(source.contains("new FleetApp(herdr, memberHerdr, workers,"),
|
||||
"FleetApp must be constructed with both the lead and the member herdr client");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,293 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
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.HerdrClient;
|
||||
import dev.ltms.fleet.lead.LeadRollover;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
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 java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
import static org.junit.jupiter.api.Assertions.fail;
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — replaces {@code FleetdLeadRolloverWiringTest} (fleetd #480). That class was a
|
||||
* source-text test scraping {@code Fleetd.java} (now {@code FleetdAssembly.java}, moved there by
|
||||
* fleetd #612 Unit A) with three methods: {@code unrelatedAnchorStillPresent} (a scaffold anchor,
|
||||
* not an independent claim — needs no replacement of its own), {@code
|
||||
* mainStillCallsTheLeadRolloverFactory} (the call-site pin replaced by {@link
|
||||
* #assembledLeadRolloverRunsTheRealClearAndBootstrapSequence}), and {@code
|
||||
* factoryGatesOnConfigPresence} (the absent-config claim replaced by {@link
|
||||
* #absentLeadRolloverConfigMeansNoRolloverIsBuilt} — a claim this ticket found was NOT actually
|
||||
* covered behaviourally anywhere else: {@code LeadRolloverTest}'s only related assertion is
|
||||
* vacuous, {@code assertNull(null)}, and never calls the real factory).
|
||||
*
|
||||
* <p><strong>fleetd #612 B3 correction (ticket comment 17553):</strong> the first version of this
|
||||
* test configured a single shared {@link FakeHerdr} for both the lead and member herdr sockets.
|
||||
* {@code FleetdAssembly.java:140-142} falls back to {@code memberHerdr = herdr} whenever no
|
||||
* distinct {@code memberHerdrSocket} is configured, so with one fake, {@code
|
||||
* router.leadAgents()} and {@code router.memberAgents()} wrapped the identical client — a
|
||||
* mutation swapping {@code Fleetd.leadRollover(cfg, router.leadAgents(), config, leads)} for
|
||||
* {@code ..., router.memberAgents(), ...} at {@code FleetdAssembly.java:408} was therefore
|
||||
* invisible to this test, even though the two are genuinely different daemons in production. This
|
||||
* version configures two distinct sockets and two distinct {@link FakeHerdr} instances (the same
|
||||
* pattern {@code FleetdAssemblyConnectionIdentityTest}, fleetd #612 B2, already uses to separate
|
||||
* lead from member) and asserts the roll's {@code /clear}/bootstrap sends land on the LEAD fake
|
||||
* and never on the MEMBER one.
|
||||
*/
|
||||
class FleetdLeadRolloverAssemblyTest {
|
||||
|
||||
private static final Path LEAD_SOCKET = Path.of("/fake/lead-herdr.sock");
|
||||
private static final Path MEMBER_SOCKET = Path.of("/fake/member-herdr.sock");
|
||||
|
||||
/** Keys {@code connectHerdr} by socket path so the lead and member daemons can be two
|
||||
* DIFFERENT {@link FakeHerdr}s — same shape as B2's {@code FleetdAssemblyConnectionIdentityTest
|
||||
* .TwoHerdrResourcePorts}. */
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final Map<Path, HerdrClient> herdrsBySocket = new LinkedHashMap<>();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
HerdrClient client = herdrsBySocket.get(socketPath);
|
||||
if (client == null) {
|
||||
throw new IllegalStateException("no fake herdr registered for socket " + socketPath);
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir, Path leadCwd) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
herdrSocket: "%s"
|
||||
memberHerdrSocket: "%s"
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
tab: "lead: opus"
|
||||
cwd: "%s"
|
||||
leadRollover:
|
||||
handoverPath: handover.md
|
||||
requireOperatorConfirm: false
|
||||
""".formatted(LEAD_SOCKET, MEMBER_SOCKET, leadCwd.toString()));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled LeadRollover runs the full open/confirm/continuation "
|
||||
+ "sequence — /clear, then bootstrapText — through the real herdr router")
|
||||
void assembledLeadRolloverRunsTheRealClearAndBootstrapSequence(@TempDir Path dir) throws Exception {
|
||||
Path leadCwd = dir.resolve("lead-workspace");
|
||||
Files.createDirectories(leadCwd);
|
||||
FleetConfig cfg = writeConfig(dir, leadCwd);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
// Two DISTINCT fakes — one per configured socket — so leadAgents()/memberAgents() wrap
|
||||
// genuinely different clients, exactly like production when memberHerdrSocket is set.
|
||||
FakeHerdr lead = new FakeHerdr();
|
||||
lead.withTab("w2", "w2:t7", "lead: opus");
|
||||
FakeHerdr member = new FakeHerdr();
|
||||
ports.herdrsBySocket.put(LEAD_SOCKET, lead);
|
||||
ports.herdrsBySocket.put(MEMBER_SOCKET, member);
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
|
||||
LeadRollover rollover = runtime.mcp().leadRollover();
|
||||
assertNotNull(rollover, "leadRollover: is present in this test's config, so "
|
||||
+ "FleetdAssembly.assembleAndStart must have built a real LeadRollover through the "
|
||||
+ "Fleetd.leadRollover(...) call site — a mutation to `LeadRollover leadRollover = "
|
||||
+ "null;` at that call site can never pass this");
|
||||
|
||||
LeadRollover.PendingRollover pending = rollover.open("term_a", "fleetd #612 B3 test");
|
||||
String expectedHandoverPath = leadCwd.resolve("handover.md").normalize().toString();
|
||||
assertEquals(expectedHandoverPath, pending.handoverPath());
|
||||
|
||||
// Ensure the handover file's mtime lands strictly AFTER open()'s requestedAtMillis —
|
||||
// LeadRollover.checkHandover refuses on mtime <= requestedAt (HANDOVER_STALE).
|
||||
Thread.sleep(50);
|
||||
Files.writeString(Path.of(pending.handoverPath()), "handover content for fleetd #612 B3");
|
||||
|
||||
LeadRollover.RollDecision decision = rollover.confirm("term_a", pending.token(), true);
|
||||
assertTrue(decision.accepted(), "confirm() must approve: requireOperatorConfirm is false, "
|
||||
+ "the caller terminal matches open()'s, and the handover file exists, is non-empty "
|
||||
+ "and fresh — got: " + decision);
|
||||
|
||||
// The production LeadRollover constructor always runs the post-confirm continuation on a
|
||||
// real virtual thread (see Fleetd.leadRollover, which never passes the package-private test
|
||||
// constructor), so this polls the real FleetMcp.leadRollover() instance's status(token)
|
||||
// until the real continuation finishes.
|
||||
LeadRollover.RollStatus status = pollUntilTerminal(rollover, pending.token());
|
||||
assertEquals(LeadRollover.RollState.ROLLED, status.state(),
|
||||
"the full happy path must complete: FakeHerdr's default agent status is 'idle', so "
|
||||
+ "the turn-boundary wait settles immediately and the post-/clear wait "
|
||||
+ "releases via its pickup-grace path — detail: " + status.detail());
|
||||
|
||||
// Prove the real herdr router actually sent BOTH messages, in order, to the real LEAD
|
||||
// pane — this is the one thing a source-text pin on the call site could never show.
|
||||
List<FakeHerdr.Call> prompts = lead.calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt"))
|
||||
.toList();
|
||||
assertTrue(prompts.size() >= 2, "expected at least a /clear send and a bootstrapText send "
|
||||
+ "on the LEAD daemon, got " + prompts.size() + " agent.prompt calls: " + prompts);
|
||||
assertEquals("/clear", ((Map<String, Object>) prompts.get(0).params()).get("text"),
|
||||
"the first send must be the literal /clear housekeeping command");
|
||||
Object secondText = ((Map<String, Object>) prompts.get(1).params()).get("text");
|
||||
assertTrue(secondText instanceof String && ((String) secondText).contains(expectedHandoverPath),
|
||||
"the second send must be the default bootstrapText naming the resolved handover "
|
||||
+ "path, got: " + secondText);
|
||||
|
||||
// fleetd #612 B3 correction: prove the roll never touches the MEMBER daemon. A mutation
|
||||
// swapping router.leadAgents() for router.memberAgents() at the real call site would move
|
||||
// both sends above onto `member` instead, which this assertion catches — the thing the
|
||||
// single-fake version of this test could never see, because both wrapped the same client.
|
||||
List<FakeHerdr.Call> memberPrompts = member.calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt"))
|
||||
.toList();
|
||||
assertTrue(memberPrompts.isEmpty(), "the roll must be wired to the LEAD daemon only — got "
|
||||
+ memberPrompts.size() + " agent.prompt call(s) on the MEMBER daemon instead: "
|
||||
+ memberPrompts);
|
||||
}
|
||||
|
||||
private static LeadRollover.RollStatus pollUntilTerminal(LeadRollover rollover, String token)
|
||||
throws InterruptedException {
|
||||
long deadline = System.nanoTime() + java.util.concurrent.TimeUnit.SECONDS.toNanos(10);
|
||||
while (System.nanoTime() < deadline) {
|
||||
LeadRollover.RollStatus status = rollover.status(token);
|
||||
if (status.state() != LeadRollover.RollState.PENDING
|
||||
&& status.state() != LeadRollover.RollState.IN_PROGRESS) {
|
||||
return status;
|
||||
}
|
||||
Thread.sleep(50);
|
||||
}
|
||||
fail("the real continuation did not reach a terminal state within 10s — last status: "
|
||||
+ rollover.status(token));
|
||||
throw new AssertionError("unreachable");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] Fleetd.leadRollover(...) returns null when leadRollover: is absent "
|
||||
+ "from config — the opt-in gate FleetdLeadRolloverWiringTest's "
|
||||
+ "factoryGatesOnConfigPresence pinned by source text alone")
|
||||
void absentLeadRolloverConfigMeansNoRolloverIsBuilt(@TempDir Path dir) throws Exception {
|
||||
Path yaml = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(yaml, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
""");
|
||||
ConfigRef config = new ConfigRef(yaml, FleetConfig.load(yaml));
|
||||
AgentControl agents = new AgentControl(new FakeHerdr());
|
||||
|
||||
LeadRollover rollover = Fleetd.leadRollover(config.get(), agents, config, Map::of);
|
||||
|
||||
assertNull(rollover, "leadRollover: is absent from this config, so the factory's opt-in "
|
||||
+ "gate (`if (cfg.leadRollover() == null) return null;`) must fire and no "
|
||||
+ "LeadRollover must be constructed at all");
|
||||
}
|
||||
}
|
||||
@@ -1,89 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #480 Unit A, hard requirement 6: pin {@code Fleetd.main}'s construction of {@link
|
||||
* dev.ltms.fleet.lead.LeadRollover} with a source-text assertion, mirroring {@code
|
||||
* FleetdCompletionResolverWiringTest}'s pattern — five log-only reporters in {@code Fleetd.main}
|
||||
* already survived mutation batteries this exact way (fleetd #415's extraction antidote note).
|
||||
*
|
||||
* <p>What this class still covers, and what it never claimed to. {@code LeadRolloverTest}
|
||||
* constructs its own {@code LeadRollover} directly (as every prior test of an extracted factory
|
||||
* does) with a hand-built lookup, so a mutation that deletes the {@code leadRollover(...)} call
|
||||
* from {@code main} — or replaces one of its arguments with something that still compiles, e.g.
|
||||
* {@code router.leadAgents()} swapped for {@code null}, or the whole assignment swapped for a bare
|
||||
* {@code null} literal — leaves every behavioural test green. This is a plain string read, guarded
|
||||
* by an unrelated anchor assertion so a broken or empty file read cannot pass as a real change.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a {@code
|
||||
* LeadRollover} and never runs {@code main}. It pins the {@code leadRollover(...)} CALL SITE's
|
||||
* argument list — that {@code main} still passes {@code leads} at all — never what the factory
|
||||
* DOES with that argument once inside its own body.
|
||||
*
|
||||
* <p><b>Correction (fleetd #480 relative-handover-path follow-up): that gap used to be real, and
|
||||
* now is not — but not here.</b> This class's javadoc previously claimed "no behavioural test can
|
||||
* catch this wiring dropping out" for the whole factory, including the lambda {@code
|
||||
* leadRollover(...)} builds internally (terminal → lead name → {@code Leader.cwd()}). That claim
|
||||
* was proven true at the time — mutating that lambda's body to {@code String leadName = null;}
|
||||
* (always "no lead found", which silently reintroduces the daemon-cwd bug this ticket fixes) left
|
||||
* the full suite green, {@code Tests run: 1669, Failures: 0}. It is no longer true: {@code
|
||||
* FleetdLeadRolloverWorkspaceLookupTest} now calls {@code Fleetd.leadRollover(...)} directly with a
|
||||
* real {@link dev.ltms.fleet.config.ConfigRef} built from a temp {@code fleetd.yaml}, and fails
|
||||
* against that exact one-line mutation. So: THIS class still covers only the call site's argument
|
||||
* list; {@code FleetdLeadRolloverWorkspaceLookupTest} is what now covers the lambda's body. Neither
|
||||
* one subsumes the other — keep both.
|
||||
*/
|
||||
class FleetdLeadRolloverWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] unrelated anchor: Fleetd.java still declares the Fleetd class")
|
||||
void unrelatedAnchorStillPresent() throws Exception {
|
||||
// Guards the two assertions below: without this, a bad read (empty string, wrong file,
|
||||
// truncated file) could vacuously fail to contain the leadRollover(...) call too, and a
|
||||
// test that only asserts "contains X" would report a false pass for the wrong reason if X
|
||||
// happened to match. Asserting an unrelated, structurally distant string first proves the
|
||||
// read actually pulled real file content.
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains("public final class Fleetd"),
|
||||
"sanity anchor failed — the file read did not return real Fleetd.java source; the "
|
||||
+ "leadRollover(...) wiring assertions below cannot be trusted until this passes");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] main still constructs LeadRollover via the leadRollover(...) factory, exactly as heartbeat is constructed")
|
||||
void mainStillCallsTheLeadRolloverFactory() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains(
|
||||
"LeadRollover leadRollover = leadRollover(cfg, router.leadAgents(), config, leads);"),
|
||||
"Fleetd.main must still assign `LeadRollover leadRollover = leadRollover(cfg, "
|
||||
+ "router.leadAgents(), config, leads);`. Dropping this call, or swapping one of "
|
||||
+ "its arguments for something that still compiles (e.g. null in place of "
|
||||
+ "router.leadAgents()), leaves every behavioural test green — this source check is "
|
||||
+ "what must go red instead. fleetd #480 correction 2 deliberately dropped "
|
||||
+ "primaryRegistry from this call — see LeadRollover's class javadoc for why a "
|
||||
+ "single-slot lookup was wrong here. The fleetd #480 relative-handover-path "
|
||||
+ "follow-up added `leads` (terminal → lead name) so the factory can resolve a "
|
||||
+ "relative handoverPath against the calling lead's own workspace.");
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] the leadRollover(...) factory itself gates construction on cfg.leadRollover() != null")
|
||||
void factoryGatesOnConfigPresence() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains("if (cfg.leadRollover() == null) {"),
|
||||
"Fleetd.leadRollover(...) must refuse to construct a LeadRollover when the "
|
||||
+ "leadRollover: block is absent — an upgraded daemon must never silently acquire "
|
||||
+ "the ability to clear the lead's own pane. See LeadHeartbeatLoop's construction "
|
||||
+ "gate (cfg.leadHeartbeat() != null) for the pattern this mirrors.");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,180 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.mcp.FleetMcp;
|
||||
import dev.ltms.fleet.msg.ReplyInbox;
|
||||
import io.javalin.Javalin;
|
||||
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 java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #612 B3 — replaces {@code FleetdLeadSeatWiringTest} (fleetd #176), a source-text test that
|
||||
* scraped {@code Fleetd.java} (now {@code FleetdAssembly.java}, moved there by fleetd #612 Unit A)
|
||||
* for the exact {@code new FleetMcp.LeadSeatSource(Fleetd.leadSeatLookup(...))} constructor-call
|
||||
* text. That proves the right symbols appear in source; it proves nothing about what the daemon's
|
||||
* live {@code fleet_list} actually reports.
|
||||
*
|
||||
* <p>This test instead drives the REAL {@link FleetMcp.LeadSeatSource} the real {@link
|
||||
* FleetdAssembly#assembleAndStart} builds — including the REAL {@code LeadTabScanner} it wires
|
||||
* {@code Fleetd.leadSeatLookup} through — reached via {@link FleetMcp#leadSeatSource()} on the
|
||||
* live {@code FleetMcp} {@code FleetdRuntime} owns. It seeds one FakeHerdr tab labelled to match a
|
||||
* configured {@code fleet.leaders.opus.tab}, with a live agent already in it (FakeHerdr's own
|
||||
* default {@code agent.list}/{@code pane.list} entries for {@code term_a}/{@code w2:p7}/{@code
|
||||
* w2:t7} — no FakeHerdr change needed), and asserts the assembled seat source reports exactly the
|
||||
* seat {@link FleetMcp.LeadSeatSource#none()} (the inert stand-in) could never produce: 1, not 0.
|
||||
*/
|
||||
class FleetdLeadSeatAssemblyTest {
|
||||
|
||||
private static final class RecordingResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final SentinelReplyInbox replyInbox = new SentinelReplyInbox();
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> replyInbox;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException(
|
||||
"leadMailboxOpener must not be called — no coordinator: block is configured");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return System::nanoTime;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
private static final class SentinelReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
@Override
|
||||
public void own(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void release(String target) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean ack(String target, String msgId) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
broker:
|
||||
uri: "amqp://fake-test-broker/vh"
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
tab: "lead: opus"
|
||||
profile: sonnet
|
||||
profiles:
|
||||
sonnet:
|
||||
subscription: true
|
||||
argv: ["ccs", "sonnet"]
|
||||
""");
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled LeadSeatSource, backed by the real LeadTabScanner, "
|
||||
+ "reports a live lead's seat against its own subscription profile")
|
||||
void assembledLeadSeatSourceReportsALiveLeadsSeat(@TempDir Path dir) throws Exception {
|
||||
FleetConfig cfg = writeConfig(dir);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
RecordingResourcePorts ports = new RecordingResourcePorts();
|
||||
// Label FakeHerdr's own default pane's tab (term_a / w2:p7 / w2:t7, already carrying a live
|
||||
// agent) to match fleet.leaders.opus.tab exactly — no FakeHerdr change needed at all.
|
||||
ports.herdr.withTab("w2", "w2:t7", "lead: opus");
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
|
||||
FleetMcp.LeadSeatSource seatSource = runtime.mcp().leadSeatSource();
|
||||
assertEquals(1, seatSource.seatsFor().apply("sonnet"),
|
||||
"the real LeadTabScanner recognises the labelled tab as a live 'opus' lead on "
|
||||
+ "profile 'sonnet' (same credential, matched by Fleetd.leadSeatLookup), so "
|
||||
+ "subscription profile 'sonnet' must be charged one seat — "
|
||||
+ "FleetMcp.LeadSeatSource.none() (the inert stand-in this test's mutation "
|
||||
+ "swaps the call site for) always reports 0, whatever the input");
|
||||
|
||||
// A profile no lead is running on gets no seat charged — the same seat source, applied to
|
||||
// an input that must stay at the inert answer even on the real, non-inert instance.
|
||||
assertEquals(0, seatSource.seatsFor().apply("no-such-profile"));
|
||||
}
|
||||
}
|
||||
@@ -1,43 +0,0 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #176: {@code Fleetd.main} builds its {@code FleetMcp} from a 14-argument constructor whose
|
||||
* last argument is a {@code FleetMcp.LeadSeatSource} wrapping {@link Fleetd#leadSeatLookup}. That
|
||||
* argument is exactly the kind of wiring fleetd #248 warned about: dropping it (or swapping it for
|
||||
* the inert {@code FleetMcp.LeadSeatSource.none()}) compiles with 0 errors and leaves every test
|
||||
* that builds its own {@code FleetMcp}/{@code CapacitySource} directly — every test that predates
|
||||
* this ticket — green, because none of them go through {@code main} at all.
|
||||
*
|
||||
* <p>{@link FleetdLeadSeatLookupTest} proves the factory's own matching logic; this class is the
|
||||
* plain source-text assertion that proves {@code main} still passes its result in, mirroring
|
||||
* {@code FleetdCompletionResolverWiringTest}'s approach for the same class of gap.
|
||||
*
|
||||
* <p><b>This test checks source text, not runtime behaviour.</b> It never constructs a
|
||||
* {@code FleetMcp} and never runs {@code main}.
|
||||
*/
|
||||
class FleetdLeadSeatWiringTest {
|
||||
|
||||
private static String fleetdSource() throws Exception {
|
||||
return Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[SOURCE TEXT] FleetMcp's construction call still passes a LeadSeatSource built from leadSeatLookup(...)")
|
||||
void fleetMcpConstructionStillWiresLeadSeatLookup() throws Exception {
|
||||
String source = fleetdSource();
|
||||
assertTrue(source.contains("new FleetMcp.LeadSeatSource(leadSeatLookup(() -> config.get().profiles(), "
|
||||
+ "leaders, leads))"),
|
||||
"FleetMcp's construction call must still pass a LeadSeatSource built from "
|
||||
+ "Fleetd.leadSeatLookup(...). Dropping it or swapping in "
|
||||
+ "FleetMcp.LeadSeatSource.none() (fleetd #176's would-be silent regression, the same "
|
||||
+ "shape as fleetd #248's measured mutations) compiles with 0 errors and leaves every "
|
||||
+ "existing behavioural test green — this source check is what must go red instead.");
|
||||
}
|
||||
}
|
||||
+197
@@ -0,0 +1,197 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.config.ConfigRef;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import dev.ltms.fleet.guard.SubscriptionGuard;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.member.CompositePeerLauncher;
|
||||
import dev.ltms.fleet.member.HerdrPeerLauncher;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.lang.reflect.Field;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #612 step 4, rank 2 (OpenCode half) — {@link FleetdAssembly}'s {@code
|
||||
* forwardingExhaustionSink} line ({@code Fleetd.forwardingExhaustionSink(exhaustionSinkRef)}),
|
||||
* handed to {@link dev.ltms.fleet.member.OpenCodeLauncher} so its fleetd #175 model-mismatch check
|
||||
* can quarantine a credential before {@code sessions} exists to build the real sink (the
|
||||
* construction-order cycle documented at that call site). The ticket calls this independent from
|
||||
* {@code publishExhaustionSink} (pinned by {@link FleetdExhaustedPatternAssemblyTest}): a credential
|
||||
* that hits a usage limit through THIS path is never quarantined if {@code forwardingExhaustionSink}
|
||||
* is swapped for a hardcoded {@link ExhaustionSink#none()} at that call site — the OpenCode
|
||||
* launcher's own quarantine check keeps compiling and keeps "running", but it permanently talks to
|
||||
* a sink that does nothing, independent of whatever {@code publishExhaustionSink} does later.
|
||||
*
|
||||
* <p>{@code FleetdExhaustionSinkForwardingWiringTest} (fleetd #589) already proves {@code
|
||||
* Fleetd.forwardingExhaustionSink(ref)} forwards to whatever {@code ref} holds — as a bare factory
|
||||
* call, never through {@link FleetdAssembly#assembleAndStart}. It proves nothing about whether
|
||||
* THIS call site is the one FleetdAssembly actually wires into the real {@code OpenCodeLauncher}
|
||||
* it builds, which is exactly the #602/#606-shaped gap this ticket exists to close.
|
||||
*
|
||||
* <p>No accessor on {@link FleetdRuntime} reaches the adapter instances (by design — see that
|
||||
* class's own javadoc: only the final collaborators it owns directly are exposed), so this test
|
||||
* reaches the REAL, assembled {@code OpenCodeLauncher}'s {@code exhaustionSink} field the same way
|
||||
* {@code SessionManager}/{@code CompositePeerLauncher} wire it internally: a short, targeted
|
||||
* reflective walk ({@code SessionManager.launcher} → {@code CompositePeerLauncher.byProfile} →
|
||||
* {@code OpenCodeLauncher.exhaustionSink}) onto the exact object the assembly built — never a copy,
|
||||
* and never a read of the source text. Reflection is used the same way elsewhere in this suite
|
||||
* (e.g. {@code StatusPollerWatchdogTest}) to reach a private collaborator a production constructor
|
||||
* intentionally does not expose a public accessor for.
|
||||
*/
|
||||
class FleetdOpenCodeExhaustionForwardingAssemblyTest {
|
||||
|
||||
private static final class ControllableResourcePorts implements ResourcePorts {
|
||||
|
||||
final FakeHerdr herdr = new FakeHerdr();
|
||||
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L);
|
||||
Runnable shutdownHook;
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return Map.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
return herdr;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
return (uri, prefetch) -> {
|
||||
throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
return (uri, selfCoordId, prefetch) -> {
|
||||
throw new UnsupportedOperationException("leadMailboxOpener must not be called — no coordinator: block");
|
||||
};
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
|
||||
return () -> {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
|
||||
};
|
||||
}
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
return nowNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
return nowNanos::get;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
return Executors.newSingleThreadScheduledExecutor();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
this.shutdownHook = hook;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
// Deliberately never bind a real port.
|
||||
}
|
||||
}
|
||||
|
||||
private static FleetConfig writeConfig(Path dir, int cooldownSeconds) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
idleSleepGuard:
|
||||
enabled: false
|
||||
quarantineCooldownSeconds: %d
|
||||
profiles:
|
||||
gemini:
|
||||
kind: opencode
|
||||
model: google/gemini-2.5-pro
|
||||
""".formatted(cooldownSeconds));
|
||||
return FleetConfig.load(f);
|
||||
}
|
||||
|
||||
/** Reach a declared field by name on {@code target}'s runtime class, bypassing the access check. */
|
||||
private static Object readField(Object target, Class<?> declaringClass, String fieldName) throws Exception {
|
||||
Field field = declaringClass.getDeclaredField(fieldName);
|
||||
field.setAccessible(true);
|
||||
return field.get(target);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("[BEHAVIOURAL] the real assembled OpenCodeLauncher's exhaustionSink field forwards "
|
||||
+ "an onExhausted call into the real daemon's BackendQuarantine")
|
||||
void assembledOpenCodeLauncherExhaustionSinkQuarantinesTheCredential(@TempDir Path dir) throws Exception {
|
||||
int cooldownSeconds = 90;
|
||||
FleetConfig cfg = writeConfig(dir, cooldownSeconds);
|
||||
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
||||
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
||||
ControllableResourcePorts ports = new ControllableResourcePorts();
|
||||
|
||||
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
||||
try {
|
||||
PeerLauncher launcherField = (PeerLauncher) readField(runtime.sessions(),
|
||||
runtime.sessions().getClass(), "launcher");
|
||||
// CONTROL: the composite launcher must actually be the real production type with a
|
||||
// "gemini" -> OpenCodeLauncher entry — if this fails, nothing below exercised the real
|
||||
// assembly at all, rather than silently passing on an empty/wrong object.
|
||||
assertTrue(launcherField instanceof CompositePeerLauncher,
|
||||
"CONTROL: SessionManager.launcher must be the real CompositePeerLauncher the "
|
||||
+ "assembly built, got: " + launcherField);
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<String, HerdrPeerLauncher> byProfile = (Map<String, HerdrPeerLauncher>)
|
||||
readField(launcherField, CompositePeerLauncher.class, "byProfile");
|
||||
HerdrPeerLauncher adapter = byProfile.get("gemini");
|
||||
assertTrue(adapter != null && adapter.getClass().getSimpleName().equals("OpenCodeLauncher"),
|
||||
"CONTROL: the 'gemini' profile must resolve to a real OpenCodeLauncher adapter, "
|
||||
+ "got: " + adapter);
|
||||
|
||||
ExhaustionSink sink = (ExhaustionSink) readField(adapter, adapter.getClass(), "exhaustionSink");
|
||||
assertTrue(sink != null, "CONTROL: OpenCodeLauncher.exhaustionSink must never be null");
|
||||
|
||||
// The exact call OpenCodeLauncher.SessionAwareHandle#checkModelMatch makes on a real
|
||||
// model mismatch (fleetd #175): target, reason, and its own already-known profile name.
|
||||
sink.onExhausted("term_gemini_1", "opencode model mismatch (test)", "gemini");
|
||||
|
||||
BackendQuarantine quarantine = runtime.mcp().quarantineSource().quarantine();
|
||||
assertTrue(quarantine.isQuarantined("gemini"),
|
||||
"the real forwardingExhaustionSink-wired field must have delegated into the "
|
||||
+ "published production sink, which quarantines the profile's credential "
|
||||
+ "('gemini' here — no credentialId configured) — replacing "
|
||||
+ "FleetdAssembly's forwardingExhaustionSink call site with a hardcoded "
|
||||
+ "ExhaustionSink.none() must fail this assertion, since the field read "
|
||||
+ "above would then BE the inert no-op and nothing would ever reach "
|
||||
+ "quarantine.quarantine(...)");
|
||||
assertEquals(cooldownSeconds, quarantine.remainingSeconds("gemini").orElseThrow(
|
||||
() -> new AssertionError("credential must report a remaining cooldown")));
|
||||
} finally {
|
||||
if (ports.shutdownHook != null) ports.shutdownHook.run();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,202 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.guard.GuardException;
|
||||
import dev.ltms.fleet.herdr.HerdrClient;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.function.LongSupplier;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #625: pins {@link dev.ltms.fleet.guard.SubscriptionGuard#assertPrimaryClean}'s call site
|
||||
* in {@link Fleetd#main(String[])} — the ONE place it runs at startup, and the check behind the
|
||||
* bridge charter's invariant 1 (never let the primary carry {@code ANTHROPIC_BASE_URL}). Nothing
|
||||
* pinned it before this ticket: deleting {@code guard.assertPrimaryClean(...)} from {@code main}
|
||||
* left the full suite green, because the call site read the real process environment ({@code
|
||||
* System.getenv()}), which a test cannot taint from inside the JVM.
|
||||
*
|
||||
* <p>{@link Fleetd#main(String[], ResourcePorts)} (added by this ticket) is the literal production
|
||||
* sequence — not a copy of it — driven here with a {@link ResourcePorts} whose {@link
|
||||
* ResourcePorts#environment()} is a plain {@code Map} a test controls. The guard itself was
|
||||
* already pinned by {@code SubscriptionGuardTest}, directly, with a {@code Map} — that proves the
|
||||
* method's behaviour, not that {@code main} still calls it at the right point. This class pins the
|
||||
* call site and, separately, the ORDER: the guard must still run before {@code cfg.validateAll()}
|
||||
* and before {@code FleetdAssembly.assembleAndStart} touches a socket, a broker, or HTTP — not just
|
||||
* be present somewhere in {@code main}.
|
||||
*
|
||||
* <p>Presence alone is not enough (a fix that pins only presence trades an invisible deletion for
|
||||
* an invisible reordering), so each test below is built so that EITHER deleting the guard call OR
|
||||
* moving it later makes the <em>same</em> test fail — with a different exception type than the one
|
||||
* asserted, not a vacuous pass. See each test's own javadoc for how.
|
||||
*/
|
||||
class FleetdSubscriptionGuardOrderingTest {
|
||||
|
||||
private static final Map<String, String> TAINTED_ENV =
|
||||
Map.of("ANTHROPIC_BASE_URL", "http://tainted.example");
|
||||
private static final Map<String, String> CLEAN_ENV = Map.of("PATH", "/usr/bin");
|
||||
|
||||
/**
|
||||
* A {@link ResourcePorts} whose {@link #environment()} is fixed to whatever the test hands it,
|
||||
* and whose every other method refuses to be called at all. That refusal is the ordering pin:
|
||||
* if {@code main} ever reaches {@link FleetdAssembly#assembleAndStart} before the guard has had
|
||||
* a chance to throw, the very first thing the assembly does with {@code ports} is {@link
|
||||
* #connectHerdr} — so a test that expects {@link GuardException} and instead observes {@link
|
||||
* UnsupportedOperationException} has just caught the guard running too late (or not at all).
|
||||
*/
|
||||
private static final class FixedEnvPorts implements ResourcePorts {
|
||||
private final Map<String, String> env;
|
||||
|
||||
FixedEnvPorts(Map<String, String> env) {
|
||||
this.env = env;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Map<String, String> environment() {
|
||||
return env;
|
||||
}
|
||||
|
||||
@Override
|
||||
public HerdrClient connectHerdr(Path socketPath) {
|
||||
throw new UnsupportedOperationException("connectHerdr must not be called before the guard runs");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.AmqpOpener replyInboxOpener() {
|
||||
throw new UnsupportedOperationException("replyInboxOpener must not be called before the guard runs");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
||||
throw new UnsupportedOperationException("leadMailboxOpener must not be called before the guard runs");
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier nanoClock() {
|
||||
throw new UnsupportedOperationException("nanoClock must not be called before the guard runs");
|
||||
}
|
||||
|
||||
@Override
|
||||
public LongSupplier wallClockNanos() {
|
||||
throw new UnsupportedOperationException("wallClockNanos must not be called before the guard runs");
|
||||
}
|
||||
|
||||
@Override
|
||||
public ScheduledExecutorService newScheduler(String purpose) {
|
||||
throw new UnsupportedOperationException("newScheduler must not be called before the guard runs");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void addShutdownHook(Runnable hook) {
|
||||
throw new UnsupportedOperationException("addShutdownHook must not be called before the guard runs");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void startHttp(Javalin app, String host, int port) {
|
||||
throw new UnsupportedOperationException("startHttp must not be called before the guard runs");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Runnable herdrPollWait() {
|
||||
throw new UnsupportedOperationException("herdrPollWait must not be called before the guard runs");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Otherwise-invalid: {@code bind.host: 0.0.0.0} with no {@code auth.mode: token} fails {@code
|
||||
* cfg.validateAll()} (CB-501's auth-exposure check — the same fixture {@code
|
||||
* FleetdStartupValidationTest#mainRefusesANonLoopbackBindWithoutTokenMode} uses), with an
|
||||
* {@link IllegalStateException}. That is deliberate: it is what {@code main} would throw INSTEAD
|
||||
* of {@link GuardException} if the guard call were deleted, or moved to run after {@code
|
||||
* validateAll()} — a different, distinguishable exception type.
|
||||
*/
|
||||
private static Path invalidConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 0.0.0.0
|
||||
port: 8765
|
||||
""");
|
||||
return f;
|
||||
}
|
||||
|
||||
/** Passes {@code cfg.validateAll()} cleanly — nothing here trips any of its checks. */
|
||||
private static Path validConfig(Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
""");
|
||||
return f;
|
||||
}
|
||||
|
||||
/**
|
||||
* Pins the order against {@code cfg.validateAll()}. The environment is tainted and the config
|
||||
* is otherwise invalid (see {@link #invalidConfig}). If the guard runs first (the required
|
||||
* order), {@code main} throws {@link GuardException} before {@code validateAll()} is ever
|
||||
* reached. If the guard were deleted, or reordered to run after {@code validateAll()}, {@code
|
||||
* validateAll()} throws {@link IllegalStateException} instead and this assertion fails on the
|
||||
* wrong exception type.
|
||||
*/
|
||||
@Test
|
||||
void mainRefusesATaintedEnvironmentBeforeValidatingTheConfig(@TempDir Path dir) throws Exception {
|
||||
Path config = invalidConfig(dir);
|
||||
FixedEnvPorts ports = new FixedEnvPorts(TAINTED_ENV);
|
||||
|
||||
GuardException ex = assertThrows(GuardException.class,
|
||||
() -> Fleetd.main(new String[]{config.toString()}, ports));
|
||||
assertTrue(ex.getMessage().contains("tainted"),
|
||||
"expected the primary-taint message, got: " + ex.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* Pins the order against {@code FleetdAssembly.assembleAndStart}. The environment is tainted
|
||||
* and the config is otherwise VALID (see {@link #validConfig}), so {@code cfg.validateAll()}
|
||||
* passes silently and the next thing that could possibly run is the assembly's first socket
|
||||
* call. If the guard runs first (the required order), {@code main} throws {@link
|
||||
* GuardException} before assembly starts. If the guard were deleted, or reordered to run after
|
||||
* assembly begins touching {@code ports}, {@link FixedEnvPorts#connectHerdr} throws {@link
|
||||
* UnsupportedOperationException} instead and this assertion fails on the wrong exception type.
|
||||
*/
|
||||
@Test
|
||||
void mainRefusesATaintedEnvironmentBeforeAssemblyTouchesAnyPort(@TempDir Path dir) throws Exception {
|
||||
Path config = validConfig(dir);
|
||||
FixedEnvPorts ports = new FixedEnvPorts(TAINTED_ENV);
|
||||
|
||||
GuardException ex = assertThrows(GuardException.class,
|
||||
() -> Fleetd.main(new String[]{config.toString()}, ports));
|
||||
assertTrue(ex.getMessage().contains("tainted"),
|
||||
"expected the primary-taint message, got: " + ex.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* The CONTROL for the two tests above. Same otherwise-valid config, same {@link FixedEnvPorts}
|
||||
* whose every method but {@code environment()} refuses to be called — but a CLEAN environment.
|
||||
* Without this, a guard that always threw {@link GuardException} regardless of input (the
|
||||
* opposite bug — e.g. the check inverted) would make the two tests above pass for the wrong
|
||||
* reason: not because they actually drove a real taint through a real guard, but because
|
||||
* anything would have thrown {@code GuardException}. Here, with nothing to taint, the guard
|
||||
* must let {@code main} proceed into {@code cfg.validateAll()} and on into the real assembly,
|
||||
* which reaches {@code ports.connectHerdr} — and THAT throws. A loud, positive assertion: if
|
||||
* the boot path never actually ran this far, there is no {@link UnsupportedOperationException}
|
||||
* to catch, only a quiet, unexpected hang or an unrelated early failure.
|
||||
*/
|
||||
@Test
|
||||
void mainProceedsPastTheGuardOnACleanEnvironment(@TempDir Path dir) throws Exception {
|
||||
Path config = validConfig(dir);
|
||||
FixedEnvPorts ports = new FixedEnvPorts(CLEAN_ENV);
|
||||
|
||||
UnsupportedOperationException ex = assertThrows(UnsupportedOperationException.class,
|
||||
() -> Fleetd.main(new String[]{config.toString()}, ports));
|
||||
assertTrue(ex.getMessage().contains("connectHerdr"),
|
||||
"expected forward progress to reach the assembly's first port call, got: " + ex.getMessage());
|
||||
}
|
||||
}
|
||||
@@ -476,6 +476,26 @@ class LeadHeartbeatLoopTest {
|
||||
assertTrue(notice.contains("2 compactions"), notice);
|
||||
}
|
||||
|
||||
// ── fleetd #621: the notice must track the effective requireOperatorConfirm value ─────────────
|
||||
|
||||
@Test
|
||||
void contextNoticeKeepsAskingTheOperatorWhenRequireOperatorConfirmIsTrue() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||
String notice = LeadHeartbeatLoop.contextNotice(true, reading, false, true);
|
||||
assertTrue(notice.contains("ask the operator"), notice);
|
||||
assertTrue(notice.contains("Only the operator can approve the roll"), notice);
|
||||
}
|
||||
|
||||
@Test
|
||||
void contextNoticeDropsTheOperatorAskWhenRequireOperatorConfirmIsFalse() {
|
||||
var reading = new LeadContextGauge.Reading(LeadContextGauge.State.HIGH, 260_771L, 1);
|
||||
String notice = LeadHeartbeatLoop.contextNotice(true, reading, false, false);
|
||||
assertFalse(notice.contains("ask the operator"), notice);
|
||||
assertFalse(notice.contains("Only the operator can approve the roll"), notice);
|
||||
assertTrue(notice.contains("fleet_handover"), notice);
|
||||
assertTrue(notice.contains("maxDocAgeSeconds"), notice);
|
||||
}
|
||||
|
||||
// ── fleetd #609 review: the latch must mean "the notice reached the pane" ────────────────────
|
||||
//
|
||||
// These four drive LeadHeartbeatLoop.tick() directly (package-private, same reasoning as
|
||||
|
||||
@@ -1846,7 +1846,8 @@ class MessageServiceTest {
|
||||
* but backed by {@link ManualScheduler} instead of a real timer (fleetd #608): its tick never
|
||||
* fires on its own — a test drives it explicitly via {@link ManualScheduler#runDueTasks()}.
|
||||
*/
|
||||
private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler)
|
||||
private record ManualPushWiring(MessageService service, FakeHerdr leadHerdr, ManualScheduler scheduler,
|
||||
ReplyPushLoop pushLoop)
|
||||
implements AutoCloseable {
|
||||
@Override
|
||||
public void close() {
|
||||
@@ -1861,14 +1862,20 @@ class MessageServiceTest {
|
||||
* {@code anAlreadyCollectedTicketProducesNoNudge}, which sets it to 1 to prove exactly that.
|
||||
*/
|
||||
private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs) {
|
||||
return wireWithManualScheduler(maxReminders, backoffMs, System::nanoTime);
|
||||
}
|
||||
|
||||
/** As above, with an injectable clock for tests that exercise terminal-ticket pruning. */
|
||||
private ManualPushWiring wireWithManualScheduler(int maxReminders, long backoffMs,
|
||||
java.util.function.LongSupplier nowNanos) {
|
||||
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||
registry.recordDelegation(T, LEAD);
|
||||
FakeHerdr leadHerdr = new FakeHerdr();
|
||||
AgentControl leadAgents = new AgentControl(leadHerdr);
|
||||
ManualScheduler scheduler = new ManualScheduler();
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, leadAgents, inbox, scheduler, maxReminders, backoffMs);
|
||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, System::nanoTime);
|
||||
return new ManualPushWiring(service, leadHerdr, scheduler);
|
||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop, null, nowNanos);
|
||||
return new ManualPushWiring(service, leadHerdr, scheduler, pushLoop);
|
||||
}
|
||||
|
||||
private void awaitNudge(FakeHerdr leadHerdr) throws InterruptedException {
|
||||
@@ -1959,24 +1966,35 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void severalAsyncTicketsFinishingTogetherProduceOneCoalescedNudge() throws Exception {
|
||||
try (var wiring = wireWithPushLoop(1, 300)) { // wide backoff: both tickets land before the tick fires
|
||||
String first = wiring.service().sendAsync(T, "first task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "first done"));
|
||||
// Settle without polling: poll() itself marks a ticket collected (that's the point of
|
||||
// anAlreadyCollectedTicketProducesNoNudge above) — using it here to detect completion
|
||||
// would collect the ticket before the coalescing this test checks ever gets a chance.
|
||||
Thread.sleep(100);
|
||||
// A 1ms backoff is due immediately. ManualScheduler still cannot run it until this test
|
||||
// explicitly calls runDueTasks(), so both terminal tickets join one scheduled tick.
|
||||
try (var wiring = wireWithManualScheduler(1, 1)) {
|
||||
String first;
|
||||
String second;
|
||||
java.util.concurrent.CountDownLatch firstTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(firstTerminalReached::countDown);
|
||||
try {
|
||||
first = wiring.service().sendAsync(T, "first task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "first done"));
|
||||
assertTrue(firstTerminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the first ticket never reached its terminal phase");
|
||||
|
||||
String second = wiring.service().sendAsync(T, "second task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "second done"));
|
||||
Thread.sleep(100);
|
||||
java.util.concurrent.CountDownLatch secondTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(secondTerminalReached::countDown);
|
||||
second = wiring.service().sendAsync(T, "second task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "second done"));
|
||||
assertTrue(secondTerminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the second ticket never reached its terminal phase");
|
||||
} finally {
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
|
||||
}
|
||||
|
||||
awaitNudge(wiring.leadHerdr());
|
||||
Thread.sleep(200); // settle — nothing more should arrive beyond the one coalesced nudge
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(),
|
||||
"both terminal tickets must coalesce onto one scheduled tick");
|
||||
long nudgeCount = wiring.leadHerdr().calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt")).count();
|
||||
assertEquals(1, nudgeCount, "two tickets finishing together must produce ONE nudge, not two");
|
||||
@@ -2055,7 +2073,8 @@ class MessageServiceTest {
|
||||
|
||||
@Test
|
||||
void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception {
|
||||
try (var wiring = wireWithPushLoop(5, 50)) {
|
||||
// A 1ms backoff is due immediately, but ManualScheduler only ticks when this test asks it to.
|
||||
try (var wiring = wireWithManualScheduler(5, 1)) {
|
||||
String ticket = wiring.service().sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
@@ -2063,8 +2082,9 @@ class MessageServiceTest {
|
||||
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
|
||||
() -> wiring.service().ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING);
|
||||
awaitQuestionPendingOn(wiring.pushLoop(), LEAD, asking.turnId());
|
||||
|
||||
awaitNudge(wiring.leadHerdr());
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(), "the open question must have one scheduled tick");
|
||||
long callsBeforeAnswer = wiring.leadHerdr().calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt")).count();
|
||||
|
||||
@@ -2072,12 +2092,22 @@ class MessageServiceTest {
|
||||
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
java.util.concurrent.CountDownLatch terminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(terminalReached::countDown);
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
try {
|
||||
assertTrue(terminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the answered ticket never reached its terminal phase");
|
||||
} finally {
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
|
||||
}
|
||||
answer.get(5, TimeUnit.SECONDS);
|
||||
|
||||
// Let several more ticks (and the ticket's own now-legitimate terminal nudge) fire —
|
||||
// none of them may still name the question's turnId, which is closed.
|
||||
Thread.sleep(300);
|
||||
// Run the ticket's legitimate terminal nudge and every later scheduled tick through its
|
||||
// reminder cap. None may still name the closed question.
|
||||
for (int tick = 0; tick < 6; tick++) {
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(), "expected one scheduled reminder tick");
|
||||
}
|
||||
boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt"))
|
||||
.skip(callsBeforeAnswer)
|
||||
@@ -2160,35 +2190,45 @@ class MessageServiceTest {
|
||||
// decideTickets hits the cap and STOPs — activeLeads drops the lead, but (before the fix)
|
||||
// pendingTickets never drops the ticket. That is the exact "cap already STOPped" branch of
|
||||
// the bug report, reached deterministically rather than by timing it against a live tick.
|
||||
try (var wiring = wireWithPushLoop(1, 50, clock::get)) {
|
||||
String stale = wiring.service().sendAsync(T, "first task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "stale result"));
|
||||
// A 1ms backoff is due immediately, but ManualScheduler runs only the ticks below.
|
||||
try (var wiring = wireWithManualScheduler(1, 1, clock::get)) {
|
||||
String stale;
|
||||
String fresh;
|
||||
java.util.concurrent.CountDownLatch staleTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(staleTerminalReached::countDown);
|
||||
try {
|
||||
stale = wiring.service().sendAsync(T, "first task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "stale result"));
|
||||
assertTrue(staleTerminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the stale ticket never reached its terminal phase");
|
||||
|
||||
// Let the reminder loop fire its one nudge and hit the cap (STOP removes it from
|
||||
// activeLeads; pendingTickets is untouched either way — that asymmetry is the bug).
|
||||
awaitNudge(wiring.leadHerdr());
|
||||
Thread.sleep(300);
|
||||
assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale),
|
||||
"sanity: the stale ticket's own reminder must have fired first");
|
||||
// Fire the stale ticket's one nudge, then its cap tick (STOP removes it from activeLeads;
|
||||
// pendingTickets is untouched either way — that asymmetry is the bug).
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one nudge tick");
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(), "the stale ticket must have one cap tick");
|
||||
assertTrue(wiring.leadHerdr().lastCall("agent.prompt").params().toString().contains(stale),
|
||||
"sanity: the stale ticket's own reminder must have fired first");
|
||||
|
||||
// Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly.
|
||||
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
|
||||
// Cross the TTL — a real clock would need 10 minutes; the injected one does it instantly.
|
||||
clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1));
|
||||
|
||||
// A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes.
|
||||
String fresh = wiring.service().sendAsync(T, "second task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "fresh result"));
|
||||
|
||||
// The fresh ticket restarts the (now-dormant) reminder loop with its own nudge.
|
||||
long before = wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count();
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (wiring.leadHerdr().calls.stream().filter(c -> c.method().equals("agent.prompt")).count() <= before
|
||||
&& System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(10);
|
||||
// A second, unrelated ticket to the same target/lead reaches sendAsync, which prunes.
|
||||
java.util.concurrent.CountDownLatch freshTerminalReached = new java.util.concurrent.CountDownLatch(1);
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(freshTerminalReached::countDown);
|
||||
fresh = wiring.service().sendAsync(T, "second task");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "fresh result"));
|
||||
assertTrue(freshTerminalReached.await(5, TimeUnit.SECONDS),
|
||||
"the fresh ticket never reached its terminal phase");
|
||||
} finally {
|
||||
wiring.service().setAfterFinishAsyncTaskCompleteHookForTest(null);
|
||||
}
|
||||
|
||||
// The fresh ticket restarts the now-dormant reminder loop with its own nudge.
|
||||
assertEquals(1, wiring.scheduler().runDueTasks(), "the fresh ticket must have one nudge tick");
|
||||
String latestNudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
|
||||
assertTrue(latestNudge.contains(fresh), "the fresh ticket's nudge must still arrive: " + latestNudge);
|
||||
assertFalse(latestNudge.contains(stale),
|
||||
|
||||
Executable
+732
@@ -0,0 +1,732 @@
|
||||
#!/usr/bin/env bash
|
||||
#
|
||||
# The one auditable way to edit the live fleetd.yaml.
|
||||
#
|
||||
# fleetd ticket #635 — why this exists at all: fleetd.yaml is gitignored and holds the live
|
||||
# fleet's settings. A bad raw edit reaches a daemon that is already serving, so a direct `Edit`
|
||||
# on it is refused by policy. This script is the allow-listed alternative, and it is not just
|
||||
# convenience — it is the thing a raw file write can never give you: a backup, a parse check
|
||||
# BEFORE the file is installed, and the daemon's own reload verdict read back afterwards. An
|
||||
# edit to a live config is not finished when the bytes are written. It is finished when the
|
||||
# daemon has said what it did with them.
|
||||
#
|
||||
# What the daemon says, and how this script finds it — measured against `ConfigRef.java` on
|
||||
# fleetd commit 158a2a8, 2026-10-01:
|
||||
#
|
||||
# 1. `ConfigRef` re-reads fleetd.yaml only when the WATCHER sees the mtime move (every 10s by
|
||||
# default — read the real interval out of the daemon's own startup line, "config watch: ...
|
||||
# re-read when it changes (every Ns)"). So a verdict never appears before the next tick.
|
||||
# 2. `ConfigRef.Outcome.summary()` logs exactly one of five strings (ConfigRef.java:371-391):
|
||||
# config reload refused — <error message>
|
||||
# config reload refused — these keys cannot change under a running daemon: <keys>. ...
|
||||
# config reloaded
|
||||
# config reloaded; these changes need a restart to take effect: <keys>
|
||||
# config reloaded; partially live — <key: detail | ...>
|
||||
# A parse/validation failure logs a DIFFERENT line instead, before any summary ever runs
|
||||
# (ConfigRef.java:425): "config reload from <path> refused, keeping the running config:
|
||||
# <message>". This script recognises both shapes of refusal.
|
||||
# 3. The em dash in those strings is a real multi-byte character — match the stable prefix
|
||||
# "config reload refused" (or "...refused, keeping the running config" for the parse-failure
|
||||
# shape), never the dash itself.
|
||||
# 4. A cold-key change (bind/herdrSocket/memberHerdrSocket/broker/auth) throws away the WHOLE
|
||||
# reload — the running config keeps every old value, not only the cold one.
|
||||
# 5. A deferred/split change IS applied (current.set(fresh) runs) — "needs a restart" is a
|
||||
# SUCCESS with a follow-up, never a failure.
|
||||
#
|
||||
# Four outcomes, and they stay four (see the exit code table below). The one most likely to be
|
||||
# gotten wrong is "cannot tell" (exit 5): the daemon may be down, or the watcher may be stalled,
|
||||
# and folding that into either "refused" or "applied" is worse than never checking at all,
|
||||
# because a caller then acts on a verdict nobody actually read. So exit 5 never restores — a
|
||||
# visible, recoverable edit beats an invisible revert of a GOOD edit.
|
||||
#
|
||||
# Usage:
|
||||
# scripts/config-edit.sh --check
|
||||
# scripts/config-edit.sh --set <yq-path>=<value> [--set ...]
|
||||
# scripts/config-edit.sh --from <candidate.yaml>
|
||||
# scripts/config-edit.sh --dry-run --set <yq-path>=<value>
|
||||
# scripts/config-edit.sh --restore
|
||||
#
|
||||
# `--set .a.b=` (an empty value — a forgotten typo) is REFUSED, not accepted as "clear the
|
||||
# field": a null value falls back to its default rather than erroring, which is silent, not
|
||||
# safe. To clear a key on purpose, write a literal null: `--set .a.b=null`. Every other value
|
||||
# is always written as a YAML string (via yq's strenv(), never spliced into the expression), so
|
||||
# there is currently no --set spelling for the literal three-character STRING "null" itself — use
|
||||
# --from for that rare case.
|
||||
#
|
||||
# Overrides (so this is drivable with no daemon — see scripts/test-config-edit.sh):
|
||||
# --config <path> default: fleetd/fleetd.yaml
|
||||
# --log <path> default: fleetd/fleetd.out
|
||||
# --wait-seconds <n> default: 4x the watch interval this script reads out of --log (10 -> 40)
|
||||
#
|
||||
# Exit codes (the --check/--restore/usage-error paths are reported separately, see below):
|
||||
# 0 applied; verdict read; clean
|
||||
# 3 applied; verdict read; needs a restart (deferred or split keys named)
|
||||
# 4 REFUSED by the daemon; backup restored (and the restore's own verdict reported if seen)
|
||||
# 5 CANNOT TELL — no verdict line inside the wait window. Nothing is restored.
|
||||
#
|
||||
# Never prints a secret. fleetd.yaml keeps credentials out by indirection (broker.uriEnv,
|
||||
# gitTokenEnv) but this script does not rely on that staying true: every diff it prints is piped
|
||||
# through `redact`, which (a) blanks the userinfo of any `scheme://user:pass@host` and (b) masks
|
||||
# the whole value on any line whose key looks like a credential. See `redact` below.
|
||||
|
||||
set -euo pipefail
|
||||
|
||||
REPO="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
|
||||
SELF="$REPO/scripts/config-edit.sh"
|
||||
|
||||
CONFIG="$REPO/fleetd/fleetd.yaml"
|
||||
LOG="$REPO/fleetd/fleetd.out"
|
||||
WAIT_SECONDS_OVERRIDE=""
|
||||
FALLBACK_PORT=8765
|
||||
|
||||
MODE=""
|
||||
DRY_RUN=0
|
||||
SETS=()
|
||||
FROM_FILE=""
|
||||
|
||||
# fleetd #635 follow-up — a signal (or any early exit while a candidate is still uninstalled) must
|
||||
# not leave a `.config-edit.XXXXXX` file sitting beside the live config forever. CAND is global
|
||||
# (never a function-local) on purpose: this ONE trap, set once, covers every path that ever
|
||||
# creates a candidate — run_edit and dry_run_diff both assign it, and clear it back to "" once the
|
||||
# file is consumed (installed, or explicitly removed), so a later, unrelated exit never retries a
|
||||
# path that already served its purpose.
|
||||
CAND=""
|
||||
cleanup_candidate() { [ -n "$CAND" ] && rm -f "$CAND" 2>/dev/null; return 0; }
|
||||
trap cleanup_candidate EXIT INT TERM
|
||||
|
||||
say() { printf '\n\033[1m== %s\033[0m\n' "$*"; }
|
||||
ok() { printf ' ok %s\n' "$*"; }
|
||||
warn() { printf ' WARN %s\n' "$*"; }
|
||||
die() { printf '\n FAIL %s\n\n' "$*" >&2; exit 1; }
|
||||
|
||||
set_mode() {
|
||||
local new="$1"
|
||||
if [ -n "$MODE" ] && [ "$MODE" != "$new" ]; then
|
||||
die "cannot combine --$MODE and --$new in one invocation"
|
||||
fi
|
||||
MODE="$new"
|
||||
}
|
||||
|
||||
while [ $# -gt 0 ]; do
|
||||
case "$1" in
|
||||
--check) set_mode check; shift ;;
|
||||
--restore) set_mode restore; shift ;;
|
||||
--set)
|
||||
[ $# -ge 2 ] || die "--set requires <yq-path>=<value>"
|
||||
set_mode set
|
||||
SETS+=("$2")
|
||||
shift 2 ;;
|
||||
--from)
|
||||
[ $# -ge 2 ] || die "--from requires a candidate file path"
|
||||
set_mode from
|
||||
FROM_FILE="$2"
|
||||
shift 2 ;;
|
||||
--dry-run) DRY_RUN=1; shift ;;
|
||||
--config)
|
||||
[ $# -ge 2 ] || die "--config requires a path"
|
||||
CONFIG="$2"; shift 2 ;;
|
||||
--log)
|
||||
[ $# -ge 2 ] || die "--log requires a path"
|
||||
LOG="$2"; shift 2 ;;
|
||||
--wait-seconds)
|
||||
[ $# -ge 2 ] || die "--wait-seconds requires a number of seconds"
|
||||
WAIT_SECONDS_OVERRIDE="$2"; shift 2 ;;
|
||||
-h|--help) sed -n '3,70p' "$SELF"; exit 0 ;;
|
||||
*) echo "unknown option: $1 (try --help)" >&2; exit 2 ;;
|
||||
esac
|
||||
done
|
||||
|
||||
[ -n "$MODE" ] || die "no action given — use --check, --set, --from, or --restore (see --help)"
|
||||
|
||||
# ------------------------------------------------------------------------------------- redaction
|
||||
#
|
||||
# Two independent passes, applied to every diff this script ever prints:
|
||||
# 1. `scheme://user:pass@host` -> `scheme://<redacted>@host`, globally (the `g` flag matters —
|
||||
# a line can carry more than one URI).
|
||||
# 2. Any line whose key looks like TOKEN|SECRET|PASSWORD|PASSWD|PASSPHRASE|CREDENTIAL|URI|_KEY,
|
||||
# matched case-insensitively against the key text (uriEnv, gitTokenEnv, ... are camelCase,
|
||||
# not SCREAMING_CASE) has its whole value blanked, diff marker and indentation kept so the
|
||||
# shape of the change is still visible. Deliberately conservative: a false-positive
|
||||
# redaction on an unrelated line costs nothing, an unredacted secret is a security defect
|
||||
# (acceptance criterion 7).
|
||||
#
|
||||
# fleetd #635 follow-up (ticket comment 17670, defect 7) — a masked key line is not the whole
|
||||
# story: a YAML block scalar (`|`, `|-`, `>`, `>-`, ...) puts the VALUE on the lines that follow
|
||||
# the key, each indented deeper than it. The key-name match above only ever sees the key line
|
||||
# itself, so those continuation lines used to flow straight through unredacted while the key line
|
||||
# right above them printed a reassuring "<redacted>" — an incomplete redactor that looks complete
|
||||
# is worse than one that visibly does nothing, because it stops a reviewer from looking further.
|
||||
# The fix is structural, not another name to match: once a key line is masked, every following
|
||||
# line indented STRICTLY DEEPER than that key is masked too, by indentation alone, until the
|
||||
# indentation returns to the key's own level or shallower. This needs no knowledge of the key's
|
||||
# name, so it covers a block scalar under any masked key — but ONLY while that key's own line is
|
||||
# itself inside the hunk being printed. `diff -u` prints just three lines of context, so a block
|
||||
# scalar's body often reaches this function with its key line left out; there is then nothing to
|
||||
# anchor to, `masked` is never set, and the body prints in full. A blank line inside a block
|
||||
# scalar loses the anchor the same way, because a blank diff line measures as indent 0. Both are
|
||||
# measured and filed as fleetd #639 — do not read this paragraph as a guarantee that a masked
|
||||
# key's value can never be printed.
|
||||
#
|
||||
# `redact` is always fed `diff -u` output, and every line of a unified diff starts with exactly
|
||||
# one of ' ', '+', '-' (the three body markers; '@'/'-'/'+' for the three header-line kinds too).
|
||||
# That one leading character is NOT part of the YAML indentation, and must be stripped before
|
||||
# indentation is measured or a key is matched — otherwise a changed ('+' or '-') line reads one
|
||||
# column shallower than it really is, and either wrongly escapes a continuation mask or wrongly
|
||||
# ends one early. Tabs are out of scope: YAML forbids them for indentation, and this is a bounded
|
||||
# fix, not a YAML parser.
|
||||
redact() {
|
||||
local line prefix content indent lead key
|
||||
local masked=0 masked_indent=0 saved_nocasematch=0
|
||||
shopt -q nocasematch && saved_nocasematch=1
|
||||
shopt -s nocasematch
|
||||
sed -E 's#://[^@]*@#://<redacted>@#g' | while IFS= read -r line || [ -n "$line" ]; do
|
||||
case "$line" in
|
||||
[\ +-]*) prefix="${line:0:1}"; content="${line:1}" ;;
|
||||
*) prefix=""; content="$line" ;;
|
||||
esac
|
||||
|
||||
indent=0
|
||||
while [ "${content:$indent:1}" = " " ]; do indent=$((indent + 1)); done
|
||||
|
||||
if [ "$masked" = 1 ] && [ "$indent" -gt "$masked_indent" ]; then
|
||||
printf '%s%*s<redacted>\n' "$prefix" "$indent" ""
|
||||
continue
|
||||
fi
|
||||
masked=0
|
||||
|
||||
if [[ "$content" =~ ^([[:space:]]*)([A-Za-z0-9_.-]+:) ]]; then
|
||||
lead="${BASH_REMATCH[1]}"
|
||||
key="${BASH_REMATCH[2]}"
|
||||
if [[ "$key" =~ (TOKEN|SECRET|PASSWORD|PASSWD|PASSPHRASE|CREDENTIAL|URI|_KEY) ]]; then
|
||||
printf '%s%s%s <redacted>\n' "$prefix" "$lead" "$key"
|
||||
masked=1
|
||||
masked_indent="$indent"
|
||||
continue
|
||||
fi
|
||||
fi
|
||||
printf '%s\n' "$line"
|
||||
done
|
||||
[ "$saved_nocasematch" = 1 ] || shopt -u nocasematch
|
||||
}
|
||||
|
||||
# ------------------------------------------------------------------------------------- the probe
|
||||
#
|
||||
# Probe the SOCKET, never `pgrep`/`ps -f` — both print argv, and argv holds `NAME=value`, making
|
||||
# either a credential channel. The port comes from the config's own `bind.port`; 8765 is only a
|
||||
# fallback when that key is absent or the file does not parse yet.
|
||||
resolve_port() {
|
||||
local file="$1" port
|
||||
if [ -f "$file" ] && command -v yq >/dev/null 2>&1; then
|
||||
port="$(yq eval '.bind.port' "$file" 2>/dev/null || true)"
|
||||
else
|
||||
port=""
|
||||
fi
|
||||
case "$port" in
|
||||
''|null) echo "$FALLBACK_PORT" ;;
|
||||
*) echo "$port" ;;
|
||||
esac
|
||||
}
|
||||
|
||||
daemon_listening() {
|
||||
local port="$1"
|
||||
if command -v nc >/dev/null 2>&1; then
|
||||
nc -z -w1 127.0.0.1 "$port" 2>/dev/null
|
||||
else
|
||||
( exec 3<>"/dev/tcp/127.0.0.1/$port" ) 2>/dev/null
|
||||
fi
|
||||
}
|
||||
|
||||
# --------------------------------------------------------------------------------- the log marker
|
||||
#
|
||||
# Take the log's line count BEFORE touching anything. Every later read of "what did the daemon
|
||||
# say" starts strictly after this mark, so a refusal from hours ago can never be mistaken for
|
||||
# this edit's verdict. Same approach as scripts/redeploy-fleetd.sh's RESTART_MARK.
|
||||
log_mark() {
|
||||
local file="$1"
|
||||
if [ -f "$file" ]; then
|
||||
wc -l < "$file" 2>/dev/null || echo 0
|
||||
else
|
||||
echo 0
|
||||
fi
|
||||
}
|
||||
|
||||
read_verdict_after_marker() {
|
||||
local file="$1" mark="$2"
|
||||
[ -f "$file" ] || return 0
|
||||
tail -n "+$((mark + 1))" "$file" 2>/dev/null || true
|
||||
}
|
||||
|
||||
# Classifies one log LINE. Echoes one of: refused | clean | needs-restart | none. Always
|
||||
# succeeds (every branch ends in `echo`), so it is safe to call from inside `$( )`.
|
||||
classify_verdict_line() {
|
||||
local line="$1"
|
||||
case "$line" in
|
||||
*'config reload refused'*) echo refused ;;
|
||||
*'config reload from '*'refused, keeping the running config'*) echo refused ;;
|
||||
*'config reloaded'*)
|
||||
case "$line" in
|
||||
*'need a restart'*|*'partially live'*) echo needs-restart ;;
|
||||
*) echo clean ;;
|
||||
esac ;;
|
||||
*) echo none ;;
|
||||
esac
|
||||
return 0
|
||||
}
|
||||
|
||||
scan_region_for_verdict() {
|
||||
local region="$1" line kind
|
||||
[ -n "$region" ] || return 1
|
||||
while IFS= read -r line || [ -n "$line" ]; do
|
||||
kind="$(classify_verdict_line "$line")"
|
||||
if [ "$kind" != "none" ]; then
|
||||
VERDICT_KIND="$kind"
|
||||
VERDICT_LINE="$line"
|
||||
return 0
|
||||
fi
|
||||
done <<< "$region"
|
||||
return 1
|
||||
}
|
||||
|
||||
# Sets VERDICT_KIND/VERDICT_LINE and returns 0 on the first verdict line found after $mark;
|
||||
# returns 1 (VERDICT_KIND=none) if none appeared inside $wait_s seconds. Checks once before each
|
||||
# sleep AND once more after the last sleep, the same boundary idiom
|
||||
# scripts/redeploy-fleetd.sh's wait_for_daemon_exit/wait_for_new_pid already use.
|
||||
wait_for_verdict() {
|
||||
local log="$1" mark="$2" wait_s="$3" _i region
|
||||
VERDICT_KIND="none"
|
||||
VERDICT_LINE=""
|
||||
for _i in $(seq "$wait_s"); do
|
||||
region="$(read_verdict_after_marker "$log" "$mark")"
|
||||
scan_region_for_verdict "$region" && return 0
|
||||
sleep 1
|
||||
done
|
||||
region="$(read_verdict_after_marker "$log" "$mark")"
|
||||
scan_region_for_verdict "$region" && return 0
|
||||
return 1
|
||||
}
|
||||
|
||||
last_verdict_line() {
|
||||
local file="$1" line out=""
|
||||
[ -f "$file" ] || return 0
|
||||
while IFS= read -r line || [ -n "$line" ]; do
|
||||
if [ "$(classify_verdict_line "$line")" != "none" ]; then
|
||||
out="$line"
|
||||
fi
|
||||
done < "$file"
|
||||
printf '%s' "$out"
|
||||
}
|
||||
|
||||
default_wait_seconds() {
|
||||
local log="$1" interval=""
|
||||
if [ -f "$log" ]; then
|
||||
interval="$(grep -F 'config watch:' "$log" 2>/dev/null | tail -1 \
|
||||
| sed -E 's/.*\(every ([0-9]+)s\).*/\1/' || true)"
|
||||
fi
|
||||
case "$interval" in
|
||||
''|*[!0-9]*) interval=10 ;;
|
||||
esac
|
||||
echo $((interval * 4))
|
||||
}
|
||||
|
||||
# ----------------------------------------------------------------------------------- the backup
|
||||
#
|
||||
# Timestamped, never pruned — "keep backups" per the ticket. A pid suffix avoids a same-second
|
||||
# collision between two invocations.
|
||||
#
|
||||
# fleetd #635 follow-up — lands under a DEDICATED, gitignored directory beside the config
|
||||
# (<dir>/.config-backups/), never beside the config file itself. The whole reason fleetd.yaml is
|
||||
# gitignored is that it must never be committed, and a backup of it inherits that requirement — a
|
||||
# bare `fleetd.yaml.bak.*` next to a tracked directory is one `git add -A`/`git add .` away from
|
||||
# committing the live config. A directory beats a glob on its own: the glob only protects today's
|
||||
# naming, a location keeps working even if the naming changes later. (See .gitignore for the glob
|
||||
# kept anyway, as a backstop for a stray backup written the old way.)
|
||||
BACKUP_DIRNAME=".config-backups"
|
||||
|
||||
backup_dir_for() {
|
||||
local src="$1"
|
||||
printf '%s/%s' "$(dirname "$src")" "$BACKUP_DIRNAME"
|
||||
}
|
||||
|
||||
backup_config() {
|
||||
local src="$1" ts backup dir base
|
||||
dir="$(backup_dir_for "$src")"
|
||||
mkdir -p "$dir" \
|
||||
|| die "could not create the backup directory $dir — refusing to edit without a backup. The live config at $src was NOT touched."
|
||||
ts="$(date -u +%Y%m%dT%H%M%S)Z"
|
||||
base="$(basename "$src")"
|
||||
backup="${dir}/${base}.bak.${ts}.$$"
|
||||
cp "$src" "$backup" \
|
||||
|| die "could not create a backup at $backup — refusing to edit without one. The live config at $src was NOT touched."
|
||||
printf '%s' "$backup"
|
||||
}
|
||||
|
||||
newest_backup() {
|
||||
local cfg="$1" dir base
|
||||
dir="$(backup_dir_for "$cfg")"
|
||||
base="$(basename "$cfg")"
|
||||
ls -t "${dir}/${base}".bak.* 2>/dev/null | head -1 || true
|
||||
}
|
||||
|
||||
# ------------------------------------------------------------------------------------ the file mode
|
||||
#
|
||||
# fleetd #635 follow-up — `mv` from a mktemp candidate carries mktemp's 0600 onto the live path
|
||||
# forever (measured: 644 -> 600 after one --set), and a restore does not undo it either, because
|
||||
# `cp` onto an EXISTING file keeps the DESTINATION's mode, not the source's. Capture the live
|
||||
# file's mode before anything touches it, and reapply it to whatever lands on that path
|
||||
# afterwards — the candidate before install, and the config again after a restore — so an edit
|
||||
# changes the file's CONTENT only, never its permissions. BSD `stat -f '%Lp'` first (matches this
|
||||
# project's dev machine), GNU `stat -c '%a'` as the fallback. Prints nothing when the file does
|
||||
# not exist yet, so apply_mode then does nothing and a first-ever edit falls back to the normal
|
||||
# umask default rather than inventing a number.
|
||||
file_mode() {
|
||||
local file="$1"
|
||||
[ -f "$file" ] || return 0
|
||||
stat -f '%Lp' "$file" 2>/dev/null || stat -c '%a' "$file" 2>/dev/null || true
|
||||
}
|
||||
|
||||
apply_mode() {
|
||||
local file="$1" mode="$2"
|
||||
[ -n "$mode" ] || return 0
|
||||
chmod "$mode" "$file" 2>/dev/null || true
|
||||
}
|
||||
|
||||
# --------------------------------------------------------------------------- candidate builders
|
||||
#
|
||||
# Never edit the live file in place. Each builder fills $1 (a temp file already sitting in the
|
||||
# SAME directory as the live config, so the later `mv` install is a rename, not a cross-device
|
||||
# copy — see run_edit).
|
||||
# fleetd #635 follow-up — a forgotten value (`--set .a.b=`, a plausible typo) must never be
|
||||
# accepted as "clear the field". `*=*` alone cannot tell "--set .a.b=" from "--set .a.b=7" apart
|
||||
# — both contain an `=` — so the guard has to look at the VALUE, not the shape of the argument.
|
||||
# An empty value refuses outright: nothing is installed, and the message names the likely cause
|
||||
# AND the two ways to actually mean it (clear on purpose, or an intentional empty string via
|
||||
# --from). Measured against the real daemon loader: a quoted empty string reads back as a null
|
||||
# field (`quoted empty -> OK int=null`), and a null numeric field FALLS BACK TO ITS DEFAULT rather
|
||||
# than erroring — so this is not a cosmetic nit, it is the one shape of edit that widens capacity
|
||||
# silently instead of failing loudly, which is exactly what this script exists to catch.
|
||||
#
|
||||
# A deliberate clear needs its own spelling, because `""` and YAML `null` are NOT the same value
|
||||
# to the loader (`""` is a valid empty String; `null` means absent, and an Integer field reads
|
||||
# either the same way — null — but a String field would keep `""` as a real value). `--set
|
||||
# .a.b=null` is that spelling: it writes a literal, unquoted `null` via yq, never the string
|
||||
# "null" through strenv(). One consequence worth knowing: there is currently no --set spelling
|
||||
# for the three-character STRING "null" itself (it collides with the clear spelling) — use
|
||||
# --from for that rare case.
|
||||
apply_set_pairs() {
|
||||
local cand="$1" kv path value
|
||||
shift
|
||||
for kv in "$@"; do
|
||||
case "$kv" in
|
||||
*=*) : ;;
|
||||
*) die "--set expects <yq-path>=<value>, got: '$kv'" ;;
|
||||
esac
|
||||
path="${kv%%=*}"
|
||||
path="${path#.}"
|
||||
value="${kv#*=}"
|
||||
if [ -z "$value" ]; then
|
||||
die "--set '$kv' has an EMPTY value — refusing. Nothing was installed. A forgotten value
|
||||
would NULL the field, and a null value falls back to its default rather than erroring —
|
||||
silent, not safe. Did you mean --set .${path}=null to clear it on purpose, or --from a
|
||||
file if you need a genuinely empty string?"
|
||||
fi
|
||||
# fleetd #635 follow-up (ticket comment 17673, defect 8) — these two failure messages used to
|
||||
# echo the full "$kv" (path=value, exactly as the operator typed it), unredacted. The operator
|
||||
# already has the value, so a terminal is not where this leaks — the risk is where the output
|
||||
# goes NEXT: this fleet pastes command output into tickets, PRs and fleet_reply bodies, and a
|
||||
# failure is exactly when someone copies it to ask for help. Print the PATH, which is what's
|
||||
# needed to fix the command, and never the value. $kv is not key:value-shaped YAML, so piping
|
||||
# it through redact would just pass it straight through — a false sense of coverage, the same
|
||||
# mistake as defect 7.
|
||||
if [ "$value" = "null" ]; then
|
||||
yq eval -i ".${path} = null" "$cand" \
|
||||
|| die "yq could not clear --set '.${path}=null' — nothing was installed. The live config is unchanged."
|
||||
continue
|
||||
fi
|
||||
CONFIG_EDIT_SET_VALUE="$value" yq eval -i ".${path} = strenv(CONFIG_EDIT_SET_VALUE)" "$cand" \
|
||||
|| die "yq could not apply --set '.${path}=<value>' — nothing was installed. The live config is unchanged."
|
||||
done
|
||||
}
|
||||
|
||||
build_from_set() {
|
||||
local cand="$1"
|
||||
cp "$CONFIG" "$cand"
|
||||
apply_set_pairs "$cand" "${SETS[@]}"
|
||||
}
|
||||
|
||||
build_from_file() {
|
||||
local cand="$1"
|
||||
[ -f "$FROM_FILE" ] || die "--from file not found: $FROM_FILE"
|
||||
cp "$FROM_FILE" "$cand"
|
||||
}
|
||||
|
||||
parse_check() {
|
||||
yq eval '.' "$1" >/dev/null 2>&1
|
||||
}
|
||||
|
||||
install_candidate() {
|
||||
local cand="$1" live="$2"
|
||||
mv -f "$cand" "$live"
|
||||
}
|
||||
|
||||
# -------------------------------------------------------------------------------- the report path
|
||||
#
|
||||
# Prints the literal command the operator (or a test) can run to restore the backup by hand — the
|
||||
# absolute path to THIS script plus the overrides actually in force, so it works from any cwd.
|
||||
restore_command_line() {
|
||||
printf '%q --restore --config %q --log %q --wait-seconds %q' "$SELF" "$CONFIG" "$LOG" "$WAIT_SECONDS"
|
||||
}
|
||||
|
||||
# State 4 only: restore the pre-edit backup, then wait for a SECOND verdict confirming the
|
||||
# restore itself reloaded cleanly. Never claims a restore it did not observe — if the second wait
|
||||
# also times out, it says so plainly rather than reporting "restored" as though confirmed.
|
||||
restore_and_confirm() {
|
||||
local backup="$1" mark2 orig_mode
|
||||
orig_mode="$(file_mode "$CONFIG")"
|
||||
mark2="$(log_mark "$LOG")"
|
||||
cp "$backup" "$CONFIG" \
|
||||
|| die "could not restore $backup onto $CONFIG — the live config is left as the REFUSED edit. Fix this by hand immediately: cp \"$backup\" \"$CONFIG\""
|
||||
apply_mode "$CONFIG" "$orig_mode"
|
||||
ok "restored from $backup"
|
||||
if wait_for_verdict "$LOG" "$mark2" "$WAIT_SECONDS"; then
|
||||
case "$VERDICT_KIND" in
|
||||
refused) warn "the RESTORE was also refused by the daemon: $VERDICT_LINE" ;;
|
||||
*) ok "restore confirmed: $VERDICT_LINE" ;;
|
||||
esac
|
||||
else
|
||||
warn "the restore is on disk, but no confirming verdict line appeared within ${WAIT_SECONDS}s"
|
||||
warn "cannot confirm the restore reloaded cleanly — check $LOG by hand"
|
||||
fi
|
||||
return 0
|
||||
}
|
||||
|
||||
# The four-outcome decision. Echoed as a function so run_edit/restore_mode share one place that
|
||||
# can return 0/3/4/5 — never duplicated, never re-worded between the two callers.
|
||||
report_outcome() {
|
||||
local mark="$1" backup="$2" kind line
|
||||
|
||||
say "waiting for the daemon's verdict (up to ${WAIT_SECONDS}s)"
|
||||
if wait_for_verdict "$LOG" "$mark" "$WAIT_SECONDS"; then
|
||||
kind="$VERDICT_KIND"; line="$VERDICT_LINE"
|
||||
else
|
||||
kind="none"
|
||||
fi
|
||||
|
||||
case "$kind" in
|
||||
clean)
|
||||
ok "daemon verdict: $line"
|
||||
say "result: applied cleanly"
|
||||
return 0 ;;
|
||||
needs-restart)
|
||||
ok "daemon verdict: $line"
|
||||
say "result: applied — a restart is needed for the change(s) named above"
|
||||
return 3 ;;
|
||||
refused)
|
||||
warn "daemon verdict: $line"
|
||||
say "result: REFUSED — restoring the backup"
|
||||
restore_and_confirm "$backup"
|
||||
return 4 ;;
|
||||
none)
|
||||
warn "no verdict line appeared within ${WAIT_SECONDS}s after $LOG line $mark"
|
||||
warn "CANNOT TELL whether the daemon applied this edit, refused it, or is simply down."
|
||||
warn "Nothing was restored — the edit is still on disk at $CONFIG."
|
||||
echo
|
||||
echo " backup: $backup"
|
||||
echo " to restore it by hand:"
|
||||
echo " $(restore_command_line)"
|
||||
return 5 ;;
|
||||
esac
|
||||
}
|
||||
|
||||
# ------------------------------------------------------------------------------------- the modes
|
||||
check_mode() {
|
||||
say "config-edit --check"
|
||||
if [ -f "$CONFIG" ]; then
|
||||
if parse_check "$CONFIG"; then
|
||||
ok "config parses: $CONFIG"
|
||||
else
|
||||
warn "config does NOT parse as valid YAML: $CONFIG"
|
||||
fi
|
||||
else
|
||||
warn "no config file at $CONFIG"
|
||||
fi
|
||||
|
||||
local port
|
||||
port="$(resolve_port "$CONFIG")"
|
||||
if daemon_listening "$port"; then
|
||||
ok "daemon is listening on 127.0.0.1:$port"
|
||||
else
|
||||
warn "no daemon detected listening on 127.0.0.1:$port"
|
||||
fi
|
||||
|
||||
ok "watch interval assumed: $(( $(default_wait_seconds "$LOG") / 4 ))s (derives --wait-seconds default of $(default_wait_seconds "$LOG")s)"
|
||||
|
||||
local verdict
|
||||
verdict="$(last_verdict_line "$LOG")"
|
||||
if [ -n "$verdict" ]; then
|
||||
ok "last verdict in log: $verdict"
|
||||
else
|
||||
warn "no reload verdict line found in $LOG"
|
||||
fi
|
||||
|
||||
local backup
|
||||
backup="$(newest_backup "$CONFIG")"
|
||||
if [ -n "$backup" ]; then
|
||||
ok "newest backup: $backup"
|
||||
else
|
||||
warn "no backups found for $CONFIG"
|
||||
fi
|
||||
|
||||
if command -v yq >/dev/null 2>&1; then
|
||||
ok "yq: $(yq --version 2>&1)"
|
||||
else
|
||||
warn "yq not found on PATH"
|
||||
fi
|
||||
|
||||
return 0
|
||||
}
|
||||
|
||||
# Shared by --set and --from: backup, build, parse-check, redacted diff, install, await verdict.
|
||||
run_edit() {
|
||||
local builder="$1"
|
||||
[ -f "$CONFIG" ] || die "no config at $CONFIG — nothing to edit"
|
||||
|
||||
local mark orig_mode
|
||||
mark="$(log_mark "$LOG")"
|
||||
orig_mode="$(file_mode "$CONFIG")"
|
||||
|
||||
say "probe"
|
||||
local port
|
||||
port="$(resolve_port "$CONFIG")"
|
||||
if daemon_listening "$port"; then
|
||||
ok "daemon appears to be listening on 127.0.0.1:$port"
|
||||
else
|
||||
warn "no daemon detected listening on 127.0.0.1:$port — a verdict may never appear"
|
||||
fi
|
||||
|
||||
say "backup"
|
||||
local backup
|
||||
backup="$(backup_config "$CONFIG")"
|
||||
ok "backup: $backup"
|
||||
|
||||
say "candidate"
|
||||
local cand
|
||||
CAND="$(mktemp "$(dirname "$CONFIG")/.config-edit.XXXXXX")" \
|
||||
|| die "could not create a candidate temp file next to $CONFIG"
|
||||
cand="$CAND"
|
||||
if ! "$builder" "$cand"; then
|
||||
rm -f "$cand"; CAND=""
|
||||
die "could not build the candidate — nothing was installed. The live config at $CONFIG is unchanged."
|
||||
fi
|
||||
|
||||
if ! parse_check "$cand"; then
|
||||
rm -f "$cand"; CAND=""
|
||||
die "candidate does not parse as valid YAML — nothing was installed. The live config at $CONFIG is unchanged."
|
||||
fi
|
||||
ok "candidate parses"
|
||||
|
||||
apply_mode "$cand" "$orig_mode"
|
||||
|
||||
say "change (redacted)"
|
||||
diff -u "$backup" "$cand" | redact || true
|
||||
|
||||
say "install"
|
||||
install_candidate "$cand" "$CONFIG" \
|
||||
|| die "could not install the candidate onto $CONFIG — the live config was NOT changed. The validated candidate is sitting at $cand; investigate before retrying."
|
||||
CAND=""
|
||||
ok "installed: $CONFIG"
|
||||
|
||||
local rc=0
|
||||
report_outcome "$mark" "$backup" || rc=$?
|
||||
return "$rc"
|
||||
}
|
||||
|
||||
dry_run_diff() {
|
||||
local builder="$1"
|
||||
[ -f "$CONFIG" ] || die "no config at $CONFIG — nothing to diff against"
|
||||
local cand
|
||||
CAND="$(mktemp "$(dirname "$CONFIG")/.config-edit.XXXXXX")" \
|
||||
|| die "could not create a candidate temp file next to $CONFIG"
|
||||
cand="$CAND"
|
||||
if ! "$builder" "$cand"; then
|
||||
rm -f "$cand"; CAND=""
|
||||
die "could not build the candidate — this was a --dry-run, nothing would have been installed either"
|
||||
fi
|
||||
if ! parse_check "$cand"; then
|
||||
rm -f "$cand"; CAND=""
|
||||
die "candidate does not parse as valid YAML — this was a --dry-run, nothing would have been installed either"
|
||||
fi
|
||||
say "dry run — diff (redacted), nothing installed"
|
||||
diff -u "$CONFIG" "$cand" | redact || true
|
||||
rm -f "$cand"; CAND=""
|
||||
return 0
|
||||
}
|
||||
|
||||
restore_mode() {
|
||||
[ -f "$CONFIG" ] || die "no config at $CONFIG to restore onto"
|
||||
local backup dir base
|
||||
backup="$(newest_backup "$CONFIG")"
|
||||
if [ -z "$backup" ]; then
|
||||
# fleetd #635 follow-up (ticket comment 17664) — this message must name the directory the
|
||||
# code actually searches (backup_dir_for, same as newest_backup), not the old beside-the-
|
||||
# config glob. A backup written the OLD way is real and NOT searched any more — say so and
|
||||
# give the one-line recovery command — but do NOT make the search itself look there; that
|
||||
# would be a behaviour change nobody asked for. The message is the only thing being fixed.
|
||||
dir="$(backup_dir_for "$CONFIG")"
|
||||
base="$(basename "$CONFIG")"
|
||||
die "no backup found matching ${dir}/${base}.bak.* — nothing to restore.
|
||||
A backup written the OLD way, directly beside the config (${CONFIG}.bak.*), is NOT
|
||||
searched — that location was retired so a backup of a file that must never be committed
|
||||
cannot sit next to a tracked directory. If one exists there, recover it by hand:
|
||||
cp ${CONFIG}.bak.<timestamp>.<pid> $CONFIG"
|
||||
fi
|
||||
[ -f "$backup" ] || die "backup candidate $backup vanished"
|
||||
|
||||
say "restore"
|
||||
ok "restoring $backup onto $CONFIG"
|
||||
local mark orig_mode
|
||||
mark="$(log_mark "$LOG")"
|
||||
orig_mode="$(file_mode "$CONFIG")"
|
||||
cp "$backup" "$CONFIG" || die "could not copy $backup onto $CONFIG"
|
||||
apply_mode "$CONFIG" "$orig_mode"
|
||||
ok "installed: $CONFIG"
|
||||
|
||||
local rc=0
|
||||
report_outcome "$mark" "$backup" || rc=$?
|
||||
return "$rc"
|
||||
}
|
||||
|
||||
# -------------------------------------------------------------------------------------- dispatch
|
||||
|
||||
if [ -n "$WAIT_SECONDS_OVERRIDE" ]; then
|
||||
WAIT_SECONDS="$WAIT_SECONDS_OVERRIDE"
|
||||
else
|
||||
WAIT_SECONDS="$(default_wait_seconds "$LOG")"
|
||||
fi
|
||||
|
||||
RC=0
|
||||
case "$MODE" in
|
||||
check)
|
||||
check_mode || RC=$?
|
||||
;;
|
||||
set)
|
||||
[ "${#SETS[@]}" -gt 0 ] || die "--set requires at least one <yq-path>=<value>"
|
||||
if [ "$DRY_RUN" = 1 ]; then
|
||||
dry_run_diff build_from_set || RC=$?
|
||||
else
|
||||
run_edit build_from_set || RC=$?
|
||||
fi
|
||||
;;
|
||||
from)
|
||||
[ -n "$FROM_FILE" ] || die "--from requires a candidate file path"
|
||||
if [ "$DRY_RUN" = 1 ]; then
|
||||
dry_run_diff build_from_file || RC=$?
|
||||
else
|
||||
run_edit build_from_file || RC=$?
|
||||
fi
|
||||
;;
|
||||
restore)
|
||||
restore_mode || RC=$?
|
||||
;;
|
||||
esac
|
||||
|
||||
exit "$RC"
|
||||
Executable
+585
@@ -0,0 +1,585 @@
|
||||
#!/usr/bin/env bash
|
||||
# Self-contained checks for scripts/config-edit.sh — fleetd ticket #635.
|
||||
#
|
||||
# Drives the REAL config-edit.sh as a subprocess against a FIXTURE config and a FIXTURE log in a
|
||||
# throwaway temp directory this file creates and removes. Never touches fleetd/fleetd.yaml or
|
||||
# fleetd/fleetd.out, and never starts, stops, or contacts a daemon — there is no daemon here, so
|
||||
# each test PLAYS the daemon: it starts config-edit.sh in the background (it is waiting on the
|
||||
# log), appends the verdict line it wants, then collects the real exit code.
|
||||
|
||||
set -euo pipefail
|
||||
|
||||
ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
|
||||
EDIT="$ROOT/scripts/config-edit.sh"
|
||||
TMP="$(mktemp -d "$ROOT/.config-edit-test.XXXXXX")"
|
||||
trap 'rm -rf "$TMP"' EXIT
|
||||
|
||||
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]"
|
||||
}
|
||||
|
||||
assert_not_contains() {
|
||||
local needle="$1" text="$2" description="$3"
|
||||
if printf '%s' "$text" | grep -qF -- "$needle"; then
|
||||
fail "$description: must NOT contain [$needle], but it does"
|
||||
fi
|
||||
return 0
|
||||
}
|
||||
|
||||
# A fresh fixture pair per test: $1/fleetd.yaml (the config) and $1/fleetd.out (the log), plus a
|
||||
# small wait-seconds budget so no test takes long. Returns the fixture dir via stdout.
|
||||
new_fixture() {
|
||||
local dir
|
||||
dir="$(mktemp -d "$TMP/fixture.XXXXXX")"
|
||||
cat > "$dir/fleetd.yaml" <<'YAML'
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 19999
|
||||
broker:
|
||||
uri: amqp://user:hunter2@host/vhost
|
||||
profiles:
|
||||
sonnet:
|
||||
weight: 3
|
||||
maxLoad: 5
|
||||
YAML
|
||||
: > "$dir/fleetd.out"
|
||||
printf '%s' "$dir"
|
||||
}
|
||||
|
||||
# Runs config-edit.sh in the background against $dir's fixtures, with the given extra args, and
|
||||
# a short --wait-seconds. Sets RUN_PID. Caller appends to $dir/fleetd.out (or not, for the
|
||||
# silence test) and then calls collect_run to block for the exit code.
|
||||
start_run() {
|
||||
local dir="$1" wait_s="$2"; shift 2
|
||||
(
|
||||
# config-edit.sh deliberately exits 3/4/5 on several of these tests. This subshell inherits
|
||||
# the parent's `set -e`, and without disabling it here the FIRST nonzero exit would kill the
|
||||
# subshell before the `echo $? > rc` line ever ran — the real code would never reach the file.
|
||||
set +e
|
||||
"$EDIT" --config "$dir/fleetd.yaml" --log "$dir/fleetd.out" --wait-seconds "$wait_s" "$@" \
|
||||
> "$dir/stdout.log" 2>&1
|
||||
echo $? > "$dir/rc"
|
||||
) &
|
||||
RUN_PID=$!
|
||||
}
|
||||
|
||||
collect_run() {
|
||||
local dir="$1"
|
||||
# wait echoes back the backgrounded subshell's own exit status (here, deliberately 3/4/5 on
|
||||
# several tests) — under `set -e` a bare nonzero `wait` would abort this whole test script, so
|
||||
# it is neutralized with `|| true`; the real code is read from $dir/rc right after.
|
||||
wait "$RUN_PID" || true
|
||||
RUN_OUTPUT="$(cat "$dir/stdout.log")"
|
||||
RUN_RC="$(cat "$dir/rc")"
|
||||
}
|
||||
|
||||
# -------------------------------------------------------------- acceptance criterion 1: refusal
|
||||
test_refusal_restores_byte_for_byte() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
cp "$dir/fleetd.yaml" "$dir/pre-edit.yaml"
|
||||
|
||||
start_run "$dir" 5 --set '.broker.uri=amqp://changed@host/x'
|
||||
sleep 1
|
||||
printf 'config reload refused — these keys cannot change under a running daemon: broker. Restart fleetd to apply them.\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 4 "$RUN_RC" "refusal exit code"
|
||||
cmp -s "$dir/fleetd.yaml" "$dir/pre-edit.yaml" \
|
||||
|| fail "refusal must restore the config byte for byte onto the pre-edit backup"
|
||||
}
|
||||
|
||||
# -------------------------------------------------------------- acceptance criterion 2: clean
|
||||
test_clean_reload_keeps_the_edit() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.weight=7'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 0 "$RUN_RC" "clean reload exit code"
|
||||
assert_equals "7" "$(yq eval '.profiles.sonnet.weight' "$dir/fleetd.yaml")" "clean reload live value"
|
||||
}
|
||||
|
||||
# ----------------------------------------------------- acceptance criterion 3: deferred != clean
|
||||
test_deferred_reload_is_told_apart_from_clean() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.weight=9'
|
||||
sleep 1
|
||||
printf 'config reloaded; these changes need a restart to take effect: profiles\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 3 "$RUN_RC" "deferred reload exit code"
|
||||
[ "$RUN_RC" != 0 ] || fail "deferred reload must not report exit 0"
|
||||
assert_contains "profiles" "$RUN_OUTPUT" "deferred reload names the key"
|
||||
assert_contains "restart" "$RUN_OUTPUT" "deferred reload says a restart is needed"
|
||||
}
|
||||
|
||||
# -------------------------------------------------------------- acceptance criterion 4: silence
|
||||
test_silence_is_its_own_answer() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
|
||||
start_run "$dir" 2 --set '.profiles.sonnet.weight=11'
|
||||
# Feed the log nothing.
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 5 "$RUN_RC" "silence exit code"
|
||||
assert_equals "11" "$(yq eval '.profiles.sonnet.weight' "$dir/fleetd.yaml")" "the edited value must still be on disk"
|
||||
assert_contains '--restore' "$RUN_OUTPUT" "silence prints the --restore command"
|
||||
|
||||
local backup restore_cmd
|
||||
backup="$(ls -t "$dir"/.config-backups/fleetd.yaml.bak.* | head -1)"
|
||||
[ -n "$backup" ] || fail "silence must still have taken a backup"
|
||||
|
||||
restore_cmd="$(printf '%s\n' "$RUN_OUTPUT" | grep -F -- '--restore --config' | sed -E 's/^[[:space:]]*//')"
|
||||
[ -n "$restore_cmd" ] || fail "could not find the printed --restore invocation in the output"
|
||||
# Running this --restore invocation installs the backup, then itself waits for a confirming
|
||||
# verdict that this fixture never feeds — so it legitimately exits 5 ("cannot tell") here, same
|
||||
# as any edit with no daemon on the other end. Only a usage/internal error (1 or 2) is a real
|
||||
# failure of the command itself; the actual assertion is the byte-for-byte cmp below.
|
||||
local restore_rc=0
|
||||
eval "$restore_cmd" > "$dir/restore.log" 2>&1 || restore_rc=$?
|
||||
case "$restore_rc" in
|
||||
0|3|4|5) : ;;
|
||||
*) fail "the printed --restore command errored out (exit $restore_rc): $(cat "$dir/restore.log")" ;;
|
||||
esac
|
||||
|
||||
cmp -s "$dir/fleetd.yaml" "$backup" \
|
||||
|| fail "running the printed --restore command must put the file back to the original backup"
|
||||
}
|
||||
|
||||
# --------------------------------------------------------- acceptance criterion 5: bad candidate
|
||||
test_broken_candidate_never_reaches_live_path() {
|
||||
local dir rc=0
|
||||
dir="$(new_fixture)"
|
||||
printf 'foo: [unclosed\n' > "$dir/broken.yaml"
|
||||
|
||||
"$EDIT" --from "$dir/broken.yaml" --config "$dir/fleetd.yaml" --log "$dir/fleetd.out" --wait-seconds 2 \
|
||||
> "$dir/stdout.log" 2>&1 || rc=$?
|
||||
|
||||
[ "$rc" -ne 0 ] || fail "a broken --from candidate must exit non-zero"
|
||||
cmp -s "$dir/fleetd.yaml" <(new_fixture_yaml) \
|
||||
|| fail "the broken candidate must never reach the live fixture config"
|
||||
}
|
||||
new_fixture_yaml() {
|
||||
cat <<'YAML'
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 19999
|
||||
broker:
|
||||
uri: amqp://user:hunter2@host/vhost
|
||||
profiles:
|
||||
sonnet:
|
||||
weight: 3
|
||||
maxLoad: 5
|
||||
YAML
|
||||
}
|
||||
|
||||
# -------------------------------------------------------------------- acceptance criterion 6
|
||||
test_marker_skips_lines_before_it() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
printf 'config reload refused — something ancient\n' > "$dir/fleetd.out"
|
||||
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.weight=5'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 0 "$RUN_RC" "a stale refusal before the marker must not be read as this edit's verdict"
|
||||
}
|
||||
|
||||
# ------------------------------------------------------------------- acceptance criterion 7 (+13)
|
||||
# fleetd #635 follow-up (ticket comment 17659) — the two assertions below this comment were the
|
||||
# WHOLE test before the follow-up, and both are negative-only: they pass just as happily when the
|
||||
# diff is never printed at all as when it is printed and correctly redacted. A mutant that deletes
|
||||
# `diff -u "$backup" "$cand" | redact` from the edit path survives them, because an absent output
|
||||
# contains neither "hunter2" nor "user:" either — see the mutation-and-revert proof in the reply.
|
||||
# Criterion 13 is the fix: a LOUD positive control that only passes when a diff was demonstrably
|
||||
# printed AND the redaction demonstrably ran on real content, not merely that nothing leaked.
|
||||
test_redaction_holds() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.weight=4'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 0 "$RUN_RC" "redaction-case reload exit code"
|
||||
assert_not_contains "hunter2" "$RUN_OUTPUT" "full output must never contain the password"
|
||||
assert_not_contains "user:" "$RUN_OUTPUT" "full output must never contain the userinfo"
|
||||
# acceptance criterion 13 — positive control: the diff's default 3-line context around the
|
||||
# changed "weight" key also covers the fixture's "uri:" line, so a genuinely-printed, genuinely-
|
||||
# redacted diff must contain BOTH the redaction marker and the changed key's name. A test that
|
||||
# only ever asserts absence cannot tell "redacted" from "never printed" apart; this can.
|
||||
assert_contains "<redacted>" "$RUN_OUTPUT" "the redaction must be PROVEN to have run on real content, not merely absent"
|
||||
assert_contains "weight" "$RUN_OUTPUT" "a diff must have been demonstrably printed at all"
|
||||
}
|
||||
|
||||
# ------------------------------------------------------- acceptance criterion 9: forgotten value
|
||||
# `--set .a.b=` is a plausible typo (the value simply forgotten), and it must be refused outright
|
||||
# rather than silently nulling the field — a null numeric field falls back to its default, which
|
||||
# widens capacity instead of failing loudly. No background verdict feeder here: a refused --set
|
||||
# must never even reach the daemon, so this never starts a background run at all.
|
||||
test_forgotten_value_refuses_and_installs_nothing() {
|
||||
local dir rc=0
|
||||
dir="$(new_fixture)"
|
||||
cp "$dir/fleetd.yaml" "$dir/pre-edit.yaml"
|
||||
|
||||
"$EDIT" --set '.profiles.sonnet.maxLoad=' \
|
||||
--config "$dir/fleetd.yaml" --log "$dir/fleetd.out" --wait-seconds 2 \
|
||||
> "$dir/stdout.log" 2>&1 || rc=$?
|
||||
RUN_OUTPUT="$(cat "$dir/stdout.log")"
|
||||
|
||||
[ "$rc" -ne 0 ] || fail "an empty --set value must exit non-zero, got 0"
|
||||
cmp -s "$dir/fleetd.yaml" "$dir/pre-edit.yaml" \
|
||||
|| fail "an empty --set value must install nothing — the live fixture changed"
|
||||
assert_contains "EMPTY value" "$RUN_OUTPUT" "the refusal must name the empty value"
|
||||
}
|
||||
|
||||
# ---------------------------------------------------------- acceptance criterion 10: explicit null
|
||||
# `--set .a.b=null` is the deliberate-clear spelling, and it must write a REAL yaml null, never
|
||||
# the string "''" — those are different values to the daemon's loader (fleetd ticket #635's
|
||||
# follow-up comment measured `""` reading back as a null field anyway, which is exactly why the
|
||||
# two forms must not collapse onto each other: `--set path=` refuses instead of silently reaching
|
||||
# this same null outcome through the back door). Read the RAW line with grep, never only through
|
||||
# `yq` — `yq eval` reports `null` for both an actual null and a missing/absent key, so it cannot
|
||||
# tell "wrote null" apart from "wrote nothing"; only the literal line on disk can.
|
||||
test_explicit_null_writes_bare_null_not_empty_string() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.maxLoad=null'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 0 "$RUN_RC" "explicit null clear exit code"
|
||||
local raw_line
|
||||
raw_line="$(grep -E 'maxLoad' "$dir/fleetd.yaml")"
|
||||
assert_contains "null" "$raw_line" "the installed line must spell a bare null"
|
||||
assert_not_contains '""' "$raw_line" "the installed line must NOT be a quoted empty string"
|
||||
}
|
||||
|
||||
# ------------------------------------------------------- acceptance criterion 11: backup never committable
|
||||
# A backup of fleetd.yaml inherits fleetd.yaml's own "never commit this" requirement (fleetd #635
|
||||
# follow-up, ticket comment 17655). Proves two things: the backup lands somewhere `git
|
||||
# check-ignore` reports as ignored (equivalently, a path `git status --porcelain` never lists as
|
||||
# untracked), AND that --restore still finds and uses it from that location.
|
||||
test_backup_is_never_committable() {
|
||||
local dir backup
|
||||
dir="$(new_fixture)"
|
||||
cp "$dir/fleetd.yaml" "$dir/pre-edit.yaml"
|
||||
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.weight=55'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
assert_equals 0 "$RUN_RC" "setup edit exit code"
|
||||
|
||||
backup="$(ls -t "$dir"/.config-backups/fleetd.yaml.bak.* 2>/dev/null | head -1)"
|
||||
[ -n "$backup" ] || fail "no backup found under .config-backups/ — did the location change?"
|
||||
|
||||
git -C "$ROOT" check-ignore -q -- "$backup" \
|
||||
|| fail "the backup at $backup is NOT gitignored — it would survive a git add -A"
|
||||
if git -C "$ROOT" status --porcelain -- "$backup" 2>/dev/null | grep -q '^??'; then
|
||||
fail "git status still lists the backup as untracked: $backup"
|
||||
fi
|
||||
|
||||
start_run "$dir" 5 --restore
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
assert_equals 0 "$RUN_RC" "--restore after the backup-location change exit code"
|
||||
cmp -s "$dir/fleetd.yaml" "$dir/pre-edit.yaml" \
|
||||
|| fail "--restore from the new backup location must still put the file back byte for byte"
|
||||
}
|
||||
|
||||
# --------------------------------------------------------- acceptance criterion 12: file mode
|
||||
# `mv` from a mktemp candidate carries mktemp's 0600 forever, and a plain `cp` onto an existing
|
||||
# file keeps the DESTINATION's mode rather than the source's, so a restore does not undo the
|
||||
# narrowing either (fleetd #635 follow-up, ticket comment 17657). Proves the mode survives an edit
|
||||
# AND a subsequent restore, from two different starting points — 644 is the common case, 600
|
||||
# proves the fix PRESERVES whatever mode was there rather than hardcoding 644.
|
||||
test_file_mode_survives_edit_and_restore() {
|
||||
local dir want got
|
||||
for want in 644 600; do
|
||||
dir="$(new_fixture)"
|
||||
chmod "$want" "$dir/fleetd.yaml"
|
||||
|
||||
start_run "$dir" 5 --set ".profiles.sonnet.weight=${want}"
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
assert_equals 0 "$RUN_RC" "mode-preservation setup edit exit code ($want)"
|
||||
got="$(stat -f '%Lp' "$dir/fleetd.yaml" 2>/dev/null || stat -c '%a' "$dir/fleetd.yaml")"
|
||||
assert_equals "$want" "$got" "mode must survive a --set ($want)"
|
||||
|
||||
start_run "$dir" 5 --restore
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
assert_equals 0 "$RUN_RC" "mode-preservation restore exit code ($want)"
|
||||
got="$(stat -f '%Lp' "$dir/fleetd.yaml" 2>/dev/null || stat -c '%a' "$dir/fleetd.yaml")"
|
||||
assert_equals "$want" "$got" "mode must survive a --restore ($want)"
|
||||
done
|
||||
}
|
||||
|
||||
# ----------------------------------------- acceptance criterion 14: restore message names the real directory
|
||||
# fleetd #635 follow-up (ticket comment 17664, defect 6) — the --restore "no backup found"
|
||||
# message used to print the OLD beside-the-config glob even though newest_backup had already
|
||||
# moved to searching the managed directory. Proves BOTH directions: the not-found message names
|
||||
# the directory actually searched (not merely that it says SOMETHING), and that a real backup
|
||||
# sitting in that directory still lets --restore succeed — otherwise the fix could regress into
|
||||
# a message that is always printed regardless of whether a backup exists.
|
||||
test_restore_message_names_the_searched_directory() {
|
||||
local dir rc=0
|
||||
|
||||
# Direction 1: no backup anywhere — the message must name .config-backups/, not the bare
|
||||
# beside-the-config glob the OLD code printed.
|
||||
dir="$(new_fixture)"
|
||||
"$EDIT" --restore --config "$dir/fleetd.yaml" --log "$dir/fleetd.out" --wait-seconds 2 \
|
||||
> "$dir/stdout.log" 2>&1 || rc=$?
|
||||
RUN_OUTPUT="$(cat "$dir/stdout.log")"
|
||||
|
||||
assert_equals 1 "$rc" "--restore with no backup anywhere exit code"
|
||||
assert_contains ".config-backups/fleetd.yaml.bak.*" "$RUN_OUTPUT" \
|
||||
"the not-found message must name the directory actually searched, not the old beside-the-config glob"
|
||||
|
||||
# Direction 2: a real backup IS present in .config-backups/ — --restore must still succeed, so
|
||||
# the message fix cannot have turned into one that prints regardless of whether a backup exists.
|
||||
dir="$(new_fixture)"
|
||||
cp "$dir/fleetd.yaml" "$dir/pre-edit.yaml"
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.weight=77'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
assert_equals 0 "$RUN_RC" "setup edit exit code for criterion 14's second half"
|
||||
|
||||
start_run "$dir" 5 --restore
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
assert_equals 0 "$RUN_RC" "--restore with a real backup present must still succeed"
|
||||
cmp -s "$dir/fleetd.yaml" "$dir/pre-edit.yaml" \
|
||||
|| fail "--restore with a real backup present must put the file back byte for byte"
|
||||
}
|
||||
|
||||
# ------------------------- acceptance criterion 15a: block-scalar continuation lines are redacted
|
||||
# fleetd #635 follow-up (ticket comment 17670, defect 7) — redact() used to look only AT the key
|
||||
# line. A YAML block scalar (`|`) puts its value on the lines that FOLLOW the key, each indented
|
||||
# deeper than it, so the real secret flowed through untouched while the key line right above it
|
||||
# printed a reassuring "<redacted>" — worse than no redaction, because the marker stops a reader
|
||||
# from looking further. The edited key here ("retries") sits directly next to the block scalar,
|
||||
# well inside diff -u's default 3-line context window, so the printed hunk is GUARANTEED to
|
||||
# include the secret's lines — placing the edit further away would let this pass today even
|
||||
# without the fix, proving nothing (the ticket comment's own warning, from the lead's first
|
||||
# reproduction attempt). The positive control runs FIRST: without it, "the secret never entered
|
||||
# the diff at all" would pass identically to "it entered and was correctly redacted".
|
||||
new_fixture_block_scalar() {
|
||||
local dir
|
||||
dir="$(mktemp -d "$TMP/fixture.XXXXXX")"
|
||||
cat > "$dir/fleetd.yaml" <<'YAML'
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 19999
|
||||
broker:
|
||||
uri: amqp://user:hunter2@host/vhost
|
||||
auth:
|
||||
token: |
|
||||
FAKELEAK-BLOCK-SCALAR
|
||||
retries: 1
|
||||
profiles:
|
||||
sonnet:
|
||||
weight: 3
|
||||
maxLoad: 5
|
||||
YAML
|
||||
: > "$dir/fleetd.out"
|
||||
printf '%s' "$dir"
|
||||
}
|
||||
|
||||
test_block_scalar_continuation_is_redacted() {
|
||||
local dir
|
||||
dir="$(new_fixture_block_scalar)"
|
||||
|
||||
start_run "$dir" 5 --set '.auth.retries=2'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 0 "$RUN_RC" "block-scalar case reload exit code"
|
||||
# Positive control FIRST: the key's own (masked) line must really be in the printed diff, or the
|
||||
# negative assertion right after proves nothing — see the comment above this test.
|
||||
assert_contains "token:" "$RUN_OUTPUT" "block-scalar case: the key's line must be in the printed diff"
|
||||
assert_contains "<redacted>" "$RUN_OUTPUT" "block-scalar case: redaction must be proven to have run on real content"
|
||||
assert_not_contains "FAKELEAK-BLOCK-SCALAR" "$RUN_OUTPUT" "block-scalar case: the block scalar's VALUE must never leak"
|
||||
}
|
||||
|
||||
# ------------------------------------- acceptance criterion 15b: "passphrase" is also recognised
|
||||
# "passphrase" was in none of TOKEN|SECRET|PASSWORD|PASSWD|CREDENTIAL|URI|_KEY (ticket comment
|
||||
# 17670). This is a plain key:value line, not a block scalar — kept in its OWN fixture and OWN
|
||||
# function, separate from criterion 15a, so that a failure in one case can never mask a failure in
|
||||
# the other (a single combined test would abort under `set -e` at its first failing assertion,
|
||||
# and the second case would then never even run).
|
||||
new_fixture_passphrase() {
|
||||
local dir
|
||||
dir="$(mktemp -d "$TMP/fixture.XXXXXX")"
|
||||
cat > "$dir/fleetd.yaml" <<'YAML'
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 19999
|
||||
broker:
|
||||
uri: amqp://user:hunter2@host/vhost
|
||||
auth:
|
||||
passphrase: FAKELEAK-PASSPHRASE
|
||||
retries: 1
|
||||
profiles:
|
||||
sonnet:
|
||||
weight: 3
|
||||
maxLoad: 5
|
||||
YAML
|
||||
: > "$dir/fleetd.out"
|
||||
printf '%s' "$dir"
|
||||
}
|
||||
|
||||
test_passphrase_key_is_redacted() {
|
||||
local dir
|
||||
dir="$(new_fixture_passphrase)"
|
||||
|
||||
start_run "$dir" 5 --set '.auth.retries=2'
|
||||
sleep 1
|
||||
printf 'config reloaded\n' >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 0 "$RUN_RC" "passphrase case reload exit code"
|
||||
assert_contains "passphrase:" "$RUN_OUTPUT" "passphrase case: the key's line must be in the printed diff"
|
||||
assert_contains "<redacted>" "$RUN_OUTPUT" "passphrase case: redaction must be proven to have run on real content"
|
||||
assert_not_contains "FAKELEAK-PASSPHRASE" "$RUN_OUTPUT" "passphrase case: the passphrase VALUE must never leak"
|
||||
}
|
||||
|
||||
# ----------------------------------- acceptance criterion 16: a failing --set must not echo value
|
||||
# fleetd #635 follow-up (ticket comment 17673, defect 8) — apply_set_pairs used to echo the FULL
|
||||
# "$kv" (path=value, exactly as typed) in its yq-failure messages, so a broken --set with a
|
||||
# secret-looking value printed that value right back out. The path alone is what the positive
|
||||
# control proves is still there — it is what the operator needs to fix their command — and the
|
||||
# negative assertion proves the value itself never appears. Kept to exactly this one failure
|
||||
# shape (an invalid yq path/expression), matching the ticket's own reproduction.
|
||||
test_failing_set_does_not_echo_its_value() {
|
||||
local dir rc=0
|
||||
dir="$(new_fixture)"
|
||||
cp "$dir/fleetd.yaml" "$dir/pre-edit.yaml"
|
||||
|
||||
"$EDIT" --dry-run --set '.broker.["bad=FAKELEAK-SETVALUE' \
|
||||
--config "$dir/fleetd.yaml" --log "$dir/fleetd.out" --wait-seconds 2 \
|
||||
> "$dir/stdout.log" 2>&1 || rc=$?
|
||||
RUN_OUTPUT="$(cat "$dir/stdout.log")"
|
||||
|
||||
[ "$rc" -ne 0 ] || fail "a --set with an invalid yq expression must exit non-zero, got 0"
|
||||
cmp -s "$dir/fleetd.yaml" "$dir/pre-edit.yaml" \
|
||||
|| fail "a failing --set must install nothing — the live fixture changed"
|
||||
# Positive control FIRST: the path must still be in the message, or the negative assertion right
|
||||
# after proves nothing (the message could simply have disappeared entirely).
|
||||
assert_contains '.broker.["bad' "$RUN_OUTPUT" "the failure message must still name the PATH"
|
||||
assert_not_contains "FAKELEAK-SETVALUE" "$RUN_OUTPUT" "the failure message must NEVER echo the VALUE"
|
||||
}
|
||||
|
||||
# dry-run must never touch the live file and must still redact.
|
||||
test_dry_run_never_installs_and_redacts() {
|
||||
local dir before
|
||||
dir="$(new_fixture)"
|
||||
before="$(cat "$dir/fleetd.yaml")"
|
||||
|
||||
"$EDIT" --dry-run --set '.profiles.sonnet.weight=99' \
|
||||
--config "$dir/fleetd.yaml" --log "$dir/fleetd.out" --wait-seconds 2 \
|
||||
> "$dir/stdout.log" 2>&1
|
||||
local rc=$?
|
||||
RUN_OUTPUT="$(cat "$dir/stdout.log")"
|
||||
|
||||
assert_equals 0 "$rc" "dry-run exit code"
|
||||
assert_equals "$before" "$(cat "$dir/fleetd.yaml")" "dry-run must never write the live config"
|
||||
assert_not_contains "hunter2" "$RUN_OUTPUT" "dry-run diff must also be redacted"
|
||||
assert_contains "99" "$RUN_OUTPUT" "dry-run diff must show the candidate value"
|
||||
# Same positive-control reasoning as acceptance criterion 13, applied to the dry-run diff path.
|
||||
assert_contains "<redacted>" "$RUN_OUTPUT" "the dry-run diff's redaction must be PROVEN to have run, not merely absent"
|
||||
}
|
||||
|
||||
# --check is read-only and always exits 0, even against a dead "daemon".
|
||||
test_check_is_read_only_and_exits_zero() {
|
||||
local dir before rc=0
|
||||
dir="$(new_fixture)"
|
||||
before="$(cat "$dir/fleetd.yaml")"
|
||||
|
||||
"$EDIT" --check --config "$dir/fleetd.yaml" --log "$dir/fleetd.out" \
|
||||
> "$dir/stdout.log" 2>&1 || rc=$?
|
||||
|
||||
assert_equals 0 "$rc" "--check exit code"
|
||||
assert_equals "$before" "$(cat "$dir/fleetd.yaml")" "--check must never modify the config"
|
||||
}
|
||||
|
||||
test_refusal_shape_from_parse_failure_wording_is_recognised() {
|
||||
local dir
|
||||
dir="$(new_fixture)"
|
||||
|
||||
start_run "$dir" 5 --set '.profiles.sonnet.weight=6'
|
||||
sleep 1
|
||||
printf 'config reload from %s refused, keeping the running config: boom\n' "$dir/fleetd.yaml" >> "$dir/fleetd.out"
|
||||
collect_run "$dir"
|
||||
|
||||
assert_equals 4 "$RUN_RC" "the parse-failure refusal shape must also exit 4, not be read as silence"
|
||||
}
|
||||
|
||||
echo "== acceptance criterion 1: refusal restores byte for byte =="
|
||||
test_refusal_restores_byte_for_byte
|
||||
echo "== acceptance criterion 2: clean reload keeps the edit =="
|
||||
test_clean_reload_keeps_the_edit
|
||||
echo "== acceptance criterion 3: deferred reload told apart from clean =="
|
||||
test_deferred_reload_is_told_apart_from_clean
|
||||
echo "== acceptance criterion 4: silence is its own answer =="
|
||||
test_silence_is_its_own_answer
|
||||
echo "== acceptance criterion 5: broken candidate never reaches the live path =="
|
||||
test_broken_candidate_never_reaches_live_path
|
||||
echo "== acceptance criterion 6: the marker works =="
|
||||
test_marker_skips_lines_before_it
|
||||
echo "== acceptance criterion 7 (+13: redaction is proven to have run) =="
|
||||
test_redaction_holds
|
||||
echo "== acceptance criterion 9: a forgotten value refuses and installs nothing =="
|
||||
test_forgotten_value_refuses_and_installs_nothing
|
||||
echo "== acceptance criterion 10: an explicit clear writes a bare null =="
|
||||
test_explicit_null_writes_bare_null_not_empty_string
|
||||
echo "== acceptance criterion 11: a backup is never committable =="
|
||||
test_backup_is_never_committable
|
||||
echo "== acceptance criterion 12: the file mode survives an edit and a restore =="
|
||||
test_file_mode_survives_edit_and_restore
|
||||
echo "== acceptance criterion 14: the restore message names the directory actually searched =="
|
||||
test_restore_message_names_the_searched_directory
|
||||
echo "== acceptance criterion 15a: a block scalar's continuation lines are redacted =="
|
||||
test_block_scalar_continuation_is_redacted
|
||||
echo "== acceptance criterion 15b: a passphrase key is also recognised =="
|
||||
test_passphrase_key_is_redacted
|
||||
echo "== acceptance criterion 16: a failing --set must not echo its value =="
|
||||
test_failing_set_does_not_echo_its_value
|
||||
echo "== extra: dry-run never installs, and redacts =="
|
||||
test_dry_run_never_installs_and_redacts
|
||||
echo "== extra: --check is read-only and always exits 0 =="
|
||||
test_check_is_read_only_and_exits_zero
|
||||
echo "== extra: the parse-failure refusal shape is also recognised =="
|
||||
test_refusal_shape_from_parse_failure_wording_is_recognised
|
||||
|
||||
printf 'PASS: config-edit acceptance criteria\n'
|
||||
Reference in New Issue
Block a user