Compare commits
31 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| ea12107497 | |||
| 591df91de1 | |||
| 6a814176f0 | |||
| d11d1d157c | |||
| d057d56156 | |||
| d703ce1313 | |||
| b32a30fd47 | |||
| 464dbc0930 | |||
| a8cadd9150 | |||
| 57b8c0b56d | |||
| 147f50c19e | |||
| eee4d576a2 | |||
| b4f9d7f53a | |||
| 3aca53b967 | |||
| 4aa1fae296 | |||
| 0c865032f9 | |||
| ea41bbf6b9 | |||
| 7b918c51ff | |||
| 554395b104 | |||
| 823976c1b5 | |||
| c8388a7f92 | |||
| 02e6aef98c | |||
| 6d493bc7bb | |||
| e5cb51a90e | |||
| e545c08082 | |||
| efa0deb9b2 | |||
| fa1f49675b | |||
| 8426c3528f | |||
| 65f98ba910 | |||
| 667254df47 | |||
| 2926cd1784 |
@@ -22,19 +22,19 @@ import java.util.function.Supplier;
|
||||
* choice rather than an accident of where the field was initialised.
|
||||
*
|
||||
* <h2>Not every key can change under a running daemon</h2>
|
||||
* Keys fall into three classes, and the difference is about what already exists when the reload
|
||||
* Keys fall into four classes, and the difference is about what already exists when the reload
|
||||
* happens — not about how important the key is.
|
||||
*
|
||||
* <ul>
|
||||
* <li><strong>Hot</strong> — re-read per use, so a reload takes effect on the next spawn:
|
||||
* {@code fleet:} (every role pool, {@code charters}, and {@code tabLabel}),
|
||||
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Those
|
||||
* three are read through a supplier on {@code CompositePeerLauncher}, which is what makes
|
||||
* them hot — not the fact that they are config. <strong>This does NOT include
|
||||
* {@code fleet.leaders}</strong>: {@code Fleetd.main} reads {@code cfg.fleet().leaders()}
|
||||
* once at startup to build the {@code LeadTabScanner} and the {@code LeadLauncher}, and
|
||||
* neither is reconstructed on reload — so a lead added, removed, or re-{@code tab}'d under
|
||||
* {@code fleet.leaders} needs a restart, the same as any deferred key below.</li>
|
||||
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Both are
|
||||
* read through a supplier on {@code CompositePeerLauncher}, which is what makes them hot —
|
||||
* not the fact that they are config. Most of {@code fleet:} — every role pool
|
||||
* ({@code architects}/{@code developers}/{@code reviewers}), {@code charters}, and
|
||||
* {@code tabLabel} — is read the same live way, through the same supplier
|
||||
* ({@code () -> config.get().fleet()}). <strong>But {@code fleet:} as a whole is NOT in this
|
||||
* class</strong>: {@code fleet.leaders} inside the same key is frozen, which is exactly what
|
||||
* makes {@code fleet:} split rather than hot — see below.</li>
|
||||
* <li><strong>Deferred</strong> — accepted into the new snapshot, but the wiring built at startup
|
||||
* keeps the old value until a restart: {@code lifecycle:}, {@code leadHeartbeat:},
|
||||
* {@code spawnReadyTimeoutMs} / {@code spawnReadyPollMs}, {@code quarantineCooldownSeconds}
|
||||
@@ -42,7 +42,21 @@ import java.util.function.Supplier;
|
||||
* {@code guard:}, {@code worktreeRoot:} and {@code worktreeGroup:} (both baked once into the
|
||||
* {@code GitWorktrees} built at {@code Fleetd.java:251} and never rebuilt — fleetd #323
|
||||
* instance 2 found {@code worktreeGroup} missing from this list and from
|
||||
* {@link #changedDeferredKeys}), adding or removing a profile (a new backend needs its own launcher,
|
||||
* {@link #changedDeferredKeys}), {@code primary:} (fleetd #326 — {@code Fleetd.java:506, 519,
|
||||
* 520} read {@code cfg.primary()} only off the startup snapshot to build {@code
|
||||
* PrimaryRegistry} and size {@code ReplyPushLoop}'s reminder cap/backoff, and neither is
|
||||
* rebuilt on reload. Say the consequence exactly: {@code primary.terminal} is DEPRECATED
|
||||
* (CB-532, and {@code Fleetd.java:511} warns about it at startup) — a lead's identity comes
|
||||
* from {@code leaders:}/{@code leadScan:}, so changing this pin does not demote or promote a
|
||||
* lead that uses those. What a changed pin still does not take effect on until a restart is
|
||||
* the fallback nudge destination the pin remains, the deprecated identity path for an operator
|
||||
* who still relies on it, and {@code pushReminders}/{@code pushBackoffMs}), {@code configReload:} (fleetd #326 — {@code
|
||||
* Fleetd.java:679-680} read it only at startup to decide whether to build a {@code
|
||||
* ConfigWatcher} at all and with what interval; the watcher that would apply a later change is
|
||||
* itself built once, so a running watcher keeps polling on its original enabled flag and
|
||||
* interval regardless of what a reload changes it to, the same shape as {@code lifecycle} —
|
||||
* not cold, because no already-open resource goes inconsistent with the new value, the watcher
|
||||
* (if any) simply keeps its old settings), adding or removing a profile (a new backend needs its own launcher,
|
||||
* which is constructed once), <em>and an existing profile's launch settings</em> —
|
||||
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl},
|
||||
* {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Fleetd.main}'s
|
||||
@@ -57,17 +71,104 @@ import java.util.function.Supplier;
|
||||
* {@code HerdrPeerLauncher} takes {@code Map.copyOf(profiles)} at construction and resolves
|
||||
* each spawn out of that copy, so those never reach a launch until the daemon restarts. A
|
||||
* reload logs these rather than pretending they applied.</li>
|
||||
* <li><strong>Cold</strong> — cannot change at all under a running daemon: {@code bind:},
|
||||
* {@code herdrSocket:}, {@code broker:} and {@code auth:}. The socket is bound, the broker
|
||||
* <li><strong>Split</strong> (fleetd #330; extended to a third key by fleetd #333) — read
|
||||
* <em>both</em> ways at different sites, so the key does not fit any class above as a whole:
|
||||
* {@code health:}, {@code coordinator:} and {@code fleet:}. Each is read off the startup
|
||||
* snapshot to build a long-lived object, and read live off {@link #get()} at a different,
|
||||
* unrelated site — so half of a reload's effect already applies while the other half waits
|
||||
* for a restart, and a bare "config reloaded" would under-claim by exactly that half.
|
||||
* <ul>
|
||||
* <li>{@code health:} — the monitor itself ({@code enabled}, {@code intervalSeconds},
|
||||
* {@code workingSuspectAfterSeconds}) is built once at {@code Fleetd.java:556-563}
|
||||
* and never rebuilt, so a changed value needs a restart to actually start, stop, or
|
||||
* retime it. The coverage string {@code fleet_profiles} reports
|
||||
* ({@code Fleetd.java:648-650}) is read live off {@link #get()} on every call, so it
|
||||
* already reflects the new value.</li>
|
||||
* <li>{@code coordinator:} — the {@code LeadMailbox} connection ({@code uri},
|
||||
* {@code uriEnv}, {@code selfId}, {@code prefetch}) is opened once at
|
||||
* {@code Fleetd.java:502} and never reopened, so a changed value needs a restart —
|
||||
* {@code selfId} in particular names this daemon's own AMQP inbox queue, and a peer
|
||||
* lead that learned the old name would not discover a new one on its own. The broker
|
||||
* URI env-var <em>name</em> that {@code MemberEnvAllowList} keeps out of a member's
|
||||
* environment is read live off {@link #get()} on every spawn
|
||||
* ({@code HerdrPeerLauncher.java:1530}), so it already applies.</li>
|
||||
* <li>{@code fleet:} (fleetd #333) — {@code fleet.leaders} is the frozen half:
|
||||
* {@code Fleetd.java:281} reads {@code cfg.fleet().leaders()} off the startup snapshot
|
||||
* for two long-lived objects built right after it and never rebuilt — the
|
||||
* {@code LeadTabScanner}'s {@code tab label → lead name} map ({@code Fleetd.java:301},
|
||||
* wired into {@code CallerResolver.withLeadsAndMembers} at {@code Fleetd.java:620/624},
|
||||
* which is how a caller's pane is recognised as a lead at all) and, when herdr answered
|
||||
* at startup, {@code LeadLauncher(...).ensureLeads()} ({@code Fleetd.java:315}), which
|
||||
* auto-launches each declared lead up to its {@code instances} count. So a lead added,
|
||||
* removed, or given a new {@code tab:} label under {@code fleet.leaders} needs a
|
||||
* restart — until then it is invisible to identity resolution, and this is exactly the
|
||||
* scenario fleetd #333 named: an operator edits a lead's {@code tab:} to match a
|
||||
* renamed pane, sees "config reloaded", and the pane keeps resolving as a worker,
|
||||
* because {@code CallerResolver} is still matching against the old label. The live
|
||||
* half is the rest of {@code fleet:} — {@code architects}/{@code developers}/
|
||||
* {@code reviewers}, {@code charters}, {@code tabLabel} — read live through the same
|
||||
* {@code CompositePeerLauncher} supplier the Hot bullet above names, so a reload that
|
||||
* only touches those already applies with nothing to report. Because the hot and frozen
|
||||
* halves of {@code fleet:} are disjoint sub-fields rather than the same fields read two
|
||||
* ways (contrast {@code coordinator.uriEnv} above), {@link #changedSplitKeys} compares
|
||||
* {@code fleet.leaders} alone, not the whole {@code Fleet} record — comparing the whole
|
||||
* record would report "split" for a {@code tabLabel}-only change that is actually fully
|
||||
* hot, over-claiming in exactly the direction this class exists to avoid under-claiming
|
||||
* in.</li>
|
||||
* </ul>
|
||||
* A split change is still accepted — {@link Outcome#applied()} stays {@code true}, the same
|
||||
* as a deferred change — because the live half genuinely took effect; refusing the whole
|
||||
* reload would leave the operator worse off than today. {@link Outcome#split()} names the
|
||||
* key and says which half is which each time, rather than trying to score "how changed" a
|
||||
* mixed key is or handle "both halves changed in one reload" as a special case.</li>
|
||||
* <li><strong>Cold</strong> — cannot change at all under a running daemon. All five of
|
||||
* {@link #COLD_KEYS}: {@code bind:}, {@code herdrSocket:}, {@code memberHerdrSocket:},
|
||||
* {@code broker:} and {@code auth:}. The sockets are already connected, the broker
|
||||
* connection is open, and the auth mode decides who may reach the port that is already
|
||||
* listening.</li>
|
||||
* listening. This bullet omitted {@code memberHerdrSocket:} until fleetd #333 — say "all
|
||||
* five of COLD_KEYS" rather than re-listing them, so prose and set cannot drift again.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330; recounted for fleetd #333).</strong>
|
||||
* {@code FleetConfig} has 22 top-level record components: 5 cold, 11 deferred, 3 split, 3
|
||||
* hot-excluded. Three of them are named nowhere in this file, and the reason is the same for all
|
||||
* three: {@code placement}, {@code memberCredentials} and {@code memberLoginShell} are
|
||||
* <strong>hot</strong> and correctly absent — all three are read live off {@code config.get()}
|
||||
* (placement through the {@code CompositePeerLauncher} supplier the Hot bullet names;
|
||||
* {@code memberCredentials}/{@code memberLoginShell} at spawn time, {@code Fleetd.java:198, 205, 729}
|
||||
* and {@code HerdrPeerLauncher#configuredMemberLoginShell}), so a reload takes effect on the next
|
||||
* spawn with no entry needed here.
|
||||
* {@code health} and {@code coordinator} used to be a third kind — <strong>undecided</strong>, not
|
||||
* hot — until fleetd #330 added the <strong>split</strong> class above and gave them a home. A
|
||||
* reload touching either used to report a bare "config reloaded", which under-claimed; now it names
|
||||
* the key and says which half is which. {@code fleet} was the same story in reverse: fleetd #330's
|
||||
* own fact-find named {@code health}/{@code coordinator} as "the complete set of split-shaped keys"
|
||||
* and filed {@code fleet.leaders}'s restart requirement as a documented caveat sitting in the
|
||||
* <strong>hot-excluded</strong> escape hatch instead — correctly documented, but in the one bucket
|
||||
* this file's own coverage test cannot check the truth of (see that test's javadoc). fleetd #333
|
||||
* moved it into <strong>split</strong>, where {@link #changedSplitKeys} actually reports it.
|
||||
* <p>The point of writing the count down: "not mentioned in this file" looks identical for a key
|
||||
* that is correctly hot and for a key nobody triaged. Three times now — {@code worktreeGroup} (#323),
|
||||
* {@code primary}/{@code configReload} (#326), and {@code fleet.leaders} sitting in the escape hatch
|
||||
* (#333) — the second kind hid among the first. A top-level coverage checker in the
|
||||
* {@link ConfigRefProfileCoverageTest} shape (one level up, over {@code FleetConfig} itself rather
|
||||
* than {@code FleetConfig.Profile}) proves this file's four classes exhaust the record's components
|
||||
* — see {@code ConfigRefTopLevelCoverageTest}. That test proves the record's <em>shape</em> is fully
|
||||
* triaged; it does NOT prove a {@code SPLIT_KEYS}/{@code COLD_KEYS}/{@code DEFERRED_KEYS} member has
|
||||
* any reporting code behind it at all — {@code ConfigRefTopLevelReportingCoverageTest} is what
|
||||
* fleetd #333 added for that, after measuring that a {@code SPLIT_KEYS} entry with its reporting
|
||||
* branch deleted passes both this file's own "kept in step" assert and
|
||||
* {@code ConfigRefTopLevelCoverageTest} unchanged. fleetd #337 extended it to {@code DEFERRED_KEYS}
|
||||
* after measuring the same one-way gap there directly: dropping {@code guard}'s branch out of
|
||||
* {@link #changedDeferredKeys} while {@code "guard"} stayed in the set left the whole suite green.
|
||||
*
|
||||
* <p><strong>A cold change refuses the whole reload.</strong> Not the hot half applied and the cold
|
||||
* half warned about: that would leave the running daemon in a state matching no file on disk, which
|
||||
* is the worst thing a reload can do to an operator debugging one. Refusing keeps the invariant that
|
||||
* the live config is always some version of the file, and the message names the keys that must
|
||||
* change through a restart.
|
||||
* change through a restart. A split change does <em>not</em> refuse, for a different reason than a
|
||||
* deferred change does not: its live half genuinely took effect, so refusing would throw that away
|
||||
* and leave the operator worse off than the partial-but-honest report {@link Outcome#split()} gives.
|
||||
*
|
||||
* <p>A reload that fails to parse or fails validation is also refused, and the previous config keeps
|
||||
* running. A config file being edited is normally read once mid-save; degrading a working daemon
|
||||
@@ -77,10 +178,43 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(ConfigRef.class);
|
||||
|
||||
/** Keys that cannot change under a running daemon — see the class doc. */
|
||||
private static final Set<String> COLD_KEYS =
|
||||
/**
|
||||
* Keys that cannot change under a running daemon — see the class doc.
|
||||
*
|
||||
* <p>Package-private (not {@code private}) so {@code ConfigRefTopLevelCoverageTest} can fold it
|
||||
* into the top-level triage it checks, the same way it reads {@link #SPLIT_KEYS}.
|
||||
*/
|
||||
static final Set<String> COLD_KEYS =
|
||||
Set.of("bind", "herdrSocket", "memberHerdrSocket", "broker", "auth");
|
||||
|
||||
/**
|
||||
* Keys read BOTH off the startup snapshot and live off {@link #get()} at different sites, so
|
||||
* neither the hot, deferred nor cold class fits them as a whole — see the class doc's Split
|
||||
* bullet (fleetd #330). A changed split key is accepted ({@link Outcome#applied()} stays
|
||||
* {@code true}) and reported by name, with a message naming which half is live and which needs
|
||||
* a restart.
|
||||
*/
|
||||
static final Set<String> SPLIT_KEYS = Set.of("health", "coordinator", "fleet");
|
||||
|
||||
/**
|
||||
* Top-level keys {@link #changedDeferredKeys} compares — see the class doc's Deferred bullet.
|
||||
* Promoted here from a test-side copy in {@code ConfigRefTopLevelCoverageTest} by fleetd #337,
|
||||
* the same reason {@link #COLD_KEYS} and {@link #SPLIT_KEYS} live here rather than in a test: a
|
||||
* second, hand-maintained copy of this set is exactly the kind of thing that silently drifts
|
||||
* from the method it is supposed to describe. {@code spawnReadyTimeoutMs} and
|
||||
* {@code spawnReadyPollMs} are compared together in one branch and reported under the combined
|
||||
* label {@code "spawnReady*"}; {@code profiles} is compared twice over (added/removed names,
|
||||
* then an existing profile's launch settings) — see {@link #changedDeferredKeys}.
|
||||
*
|
||||
* <p>Package-private (not {@code private}) so {@code ConfigRefTopLevelCoverageTest} and
|
||||
* {@code ConfigRefTopLevelReportingCoverageTest} can both read it, the same way they already
|
||||
* read {@link #COLD_KEYS} and {@link #SPLIT_KEYS}.
|
||||
*/
|
||||
static final Set<String> DEFERRED_KEYS = Set.of(
|
||||
"guard", "worktreeRoot", "worktreeGroup", "primary", "configReload",
|
||||
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
|
||||
"quarantineCooldownSeconds", "profiles");
|
||||
|
||||
private final Path path;
|
||||
private final AtomicReference<FleetConfig> current;
|
||||
|
||||
@@ -108,25 +242,37 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
/**
|
||||
* What a reload attempt did.
|
||||
*
|
||||
* <p>{@code split} is a separate field from {@code deferred} rather than a differently-worded
|
||||
* entry inside it, because the two carry different guarantees for any caller that branches on
|
||||
* them rather than just printing {@link #summary()}: every {@code deferred} entry means "this
|
||||
* key's whole change waits for a restart", while every {@code split} entry means "part of this
|
||||
* key's change already applied, and the message says which part" — collapsing them would force
|
||||
* a caller to re-parse the message to tell those apart. See the class doc's Split bullet
|
||||
* (fleetd #330) for why the key needs this at all.
|
||||
*
|
||||
* @param applied true when the new config is now live
|
||||
* @param coldKeys cold keys whose value changed, which is why an unapplied reload was refused
|
||||
* @param deferred keys that changed and were accepted, but whose effect waits for a restart
|
||||
* @param deferred keys that changed and were accepted, but whose effect waits entirely on a
|
||||
* restart
|
||||
* @param split split keys that changed and were accepted, each named with which half of it
|
||||
* is already live and which half waits for a restart
|
||||
* @param error the parse or validation failure that refused the reload, else {@code null}
|
||||
*/
|
||||
public record Outcome(boolean applied, List<String> coldKeys, List<String> deferred,
|
||||
String error) {
|
||||
List<String> split, String error) {
|
||||
|
||||
public Outcome {
|
||||
coldKeys = List.copyOf(coldKeys);
|
||||
deferred = List.copyOf(deferred);
|
||||
split = List.copyOf(split);
|
||||
}
|
||||
|
||||
static Outcome refusedCold(List<String> keys) {
|
||||
return new Outcome(false, keys, List.of(), null);
|
||||
return new Outcome(false, keys, List.of(), List.of(), null);
|
||||
}
|
||||
|
||||
static Outcome failed(String error) {
|
||||
return new Outcome(false, List.of(), List.of(), error);
|
||||
return new Outcome(false, List.of(), List.of(), List.of(), error);
|
||||
}
|
||||
|
||||
/** A one-line summary for the operator — the reason, not just the verdict. */
|
||||
@@ -138,11 +284,18 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
return "config reload refused — these keys cannot change under a running daemon: "
|
||||
+ String.join(", ", coldKeys) + ". Restart fleetd to apply them.";
|
||||
}
|
||||
if (!deferred.isEmpty()) {
|
||||
return "config reloaded; these changes need a restart to take effect: "
|
||||
+ String.join(", ", deferred);
|
||||
if (deferred.isEmpty() && split.isEmpty()) {
|
||||
return "config reloaded";
|
||||
}
|
||||
return "config reloaded";
|
||||
StringBuilder out = new StringBuilder("config reloaded");
|
||||
if (!deferred.isEmpty()) {
|
||||
out.append("; these changes need a restart to take effect: ")
|
||||
.append(String.join(", ", deferred));
|
||||
}
|
||||
if (!split.isEmpty()) {
|
||||
out.append("; partially live — ").append(String.join(" | ", split));
|
||||
}
|
||||
return out.toString();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -183,14 +336,21 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
}
|
||||
|
||||
List<String> deferred = changedDeferredKeys(old, fresh);
|
||||
List<String> split = changedSplitKeys(old, fresh);
|
||||
current.set(fresh);
|
||||
Outcome out = new Outcome(true, List.of(), deferred, null);
|
||||
Outcome out = new Outcome(true, List.of(), deferred, split, null);
|
||||
log.info(out.summary());
|
||||
return out;
|
||||
}
|
||||
|
||||
/** Cold keys whose value differs between the running config and the candidate. */
|
||||
private static List<String> changedColdKeys(FleetConfig old, FleetConfig fresh) {
|
||||
/**
|
||||
* Cold keys whose value differs between the running config and the candidate.
|
||||
*
|
||||
* <p>Package-private (not {@code private}) so {@code ConfigRefTopLevelReportingCoverageTest}
|
||||
* can call it directly with a reflection-built {@code FleetConfig} pair, the same reason
|
||||
* {@link #sameLaunchSettings} is package-private — see that test's class doc (fleetd #333).
|
||||
*/
|
||||
static List<String> changedColdKeys(FleetConfig old, FleetConfig fresh) {
|
||||
List<String> changed = new ArrayList<>();
|
||||
if (!Objects.equals(old.bind(), fresh.bind())) {
|
||||
changed.add("bind");
|
||||
@@ -212,8 +372,16 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
return changed;
|
||||
}
|
||||
|
||||
/** Changed keys that were accepted but whose effect waits for a restart. */
|
||||
private static List<String> changedDeferredKeys(FleetConfig old, FleetConfig fresh) {
|
||||
/**
|
||||
* Changed keys that were accepted but whose effect waits for a restart.
|
||||
*
|
||||
* <p>Package-private (not {@code private}) so {@code ConfigRefTopLevelReportingCoverageTest}
|
||||
* can call it directly with a reflection-built {@code FleetConfig} pair, the same reason
|
||||
* {@link #changedColdKeys} and {@link #changedSplitKeys} already are (fleetd #333, extended to
|
||||
* this method by fleetd #337 — membership in {@link #DEFERRED_KEYS} proved nothing about this
|
||||
* method on its own until then; see that test's class doc).
|
||||
*/
|
||||
static List<String> changedDeferredKeys(FleetConfig old, FleetConfig fresh) {
|
||||
List<String> changed = new ArrayList<>();
|
||||
if (!Objects.equals(old.lifecycle(), fresh.lifecycle())) {
|
||||
changed.add("lifecycle");
|
||||
@@ -234,6 +402,22 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
if (!Objects.equals(old.worktreeGroup(), fresh.worktreeGroup())) {
|
||||
changed.add("worktreeGroup");
|
||||
}
|
||||
// fleetd #326: Fleetd.java:506, 519, 520 read cfg.primary() only off the startup snapshot
|
||||
// (PrimaryRegistry's pinned terminal, ReplyPushLoop's reminder cap and backoff) — neither is
|
||||
// rebuilt on reload, so a changed value needs a restart. Note what it does NOT mean:
|
||||
// primary.terminal is deprecated (CB-532), identity comes from leaders:/leadScan:, so a lead
|
||||
// using those is unaffected by this pin either way. See the class doc for the exact scope.
|
||||
if (!Objects.equals(old.primary(), fresh.primary())) {
|
||||
changed.add("primary");
|
||||
}
|
||||
// fleetd #326: Fleetd.java:679-680 read cfg.configReload() only at startup to decide whether
|
||||
// to build a ConfigWatcher at all and with what interval — the watcher that would apply a
|
||||
// later change is itself built once, so a running watcher keeps its original enabled flag and
|
||||
// interval regardless of what a reload changes it to. Not cold: no already-open resource goes
|
||||
// inconsistent with the new value, a watcher (if any) simply keeps polling on the old settings.
|
||||
if (!Objects.equals(old.configReload(), fresh.configReload())) {
|
||||
changed.add("configReload");
|
||||
}
|
||||
if (!Objects.equals(old.spawnReadyTimeoutMs(), fresh.spawnReadyTimeoutMs())
|
||||
|| !Objects.equals(old.spawnReadyPollMs(), fresh.spawnReadyPollMs())) {
|
||||
changed.add("spawnReady*");
|
||||
@@ -277,6 +461,73 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
return changed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Split keys whose value differs between the running config and the candidate — see the class
|
||||
* doc's Split bullet (fleetd #330, extended for {@code fleet:} by fleetd #333). Unlike
|
||||
* {@link #changedDeferredKeys}, this does not try to tell which sub-field moved for {@code
|
||||
* health:} or {@code coordinator:}: any change to either gets the same fixed message, because
|
||||
* the message already names both halves every time, so there is no "which half changed"
|
||||
* question left for the caller to answer. {@code fleet:} is different on purpose — see below.
|
||||
*
|
||||
* <p>Package-private (not {@code private}) so {@code ConfigRefTopLevelReportingCoverageTest}
|
||||
* can call it directly with a reflection-built {@code FleetConfig} pair, the same reason
|
||||
* {@link #sameLaunchSettings} is package-private — see that test's class doc (fleetd #333). That
|
||||
* test exists because membership in {@link #SPLIT_KEYS} proves nothing about this method on its
|
||||
* own: fleetd #333 measured that dropping the {@code coordinator} branch out of this method
|
||||
* while leaving {@code "coordinator"} in {@code SPLIT_KEYS} left the whole suite green except a
|
||||
* hand-written {@code ConfigRefTest} case — neither {@code ConfigRefTopLevelCoverageTest} (it
|
||||
* only reads the set) nor the "kept in step" assert below (it only checks the reported keys are
|
||||
* a SUBSET of {@code SPLIT_KEYS}, never that every {@code SPLIT_KEYS} member has a branch here)
|
||||
* would have caught it.
|
||||
*/
|
||||
static List<String> changedSplitKeys(FleetConfig old, FleetConfig fresh) {
|
||||
List<String> changed = new ArrayList<>();
|
||||
if (!Objects.equals(old.health(), fresh.health())) {
|
||||
changed.add("health: the monitor itself (enabled, interval, workingSuspectAfter) is "
|
||||
+ "frozen at startup and needs a restart; the coverage status fleet_profiles "
|
||||
+ "reports is read live and already applied");
|
||||
}
|
||||
if (!Objects.equals(old.coordinator(), fresh.coordinator())) {
|
||||
changed.add("coordinator: the LeadMailbox connection (uri, uriEnv, selfId, prefetch) is "
|
||||
+ "opened once and needs a restart; the broker URI env-var name kept out of a "
|
||||
+ "member's environment is read live on every spawn and already applied");
|
||||
}
|
||||
// fleetd #333: unlike health/coordinator above, most of `fleet:` (architects, developers,
|
||||
// reviewers, charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
|
||||
// ThroughGet and aCharterChangeIsHotAndReachesTheLiveConfig prove it reaches the live config
|
||||
// with no restart note. Only fleet.leaders is frozen (Fleetd.java:281 reads
|
||||
// cfg.fleet().leaders() off the startup snapshot to build both the LeadTabScanner's
|
||||
// tab-label-to-name map, wired into CallerResolver.withLeadsAndMembers at Fleetd.java:620/624,
|
||||
// and — when herdr answered — LeadLauncher(...).ensureLeads() at Fleetd.java:315, which
|
||||
// auto-launches each lead up to its `instances` count; neither is rebuilt on reload). So this
|
||||
// compares fleet.leaders alone, not the whole Fleet record: comparing the whole record would
|
||||
// report "split" for a tabLabel-only or charters-only change that is actually fully hot,
|
||||
// which is the over-claim mirror of the under-claim bug this class exists to prevent.
|
||||
if (!Objects.equals(leadersOf(old), leadersOf(fresh))) {
|
||||
changed.add("fleet: fleet.leaders (each lead's tab, workspace, cwd, profile and "
|
||||
+ "instances count) is read once at startup to build the LeadTabScanner's "
|
||||
+ "identity map and to auto-launch leads, and neither is rebuilt on reload, so a "
|
||||
+ "lead added, removed, or given a new tab: label needs a restart — until then it "
|
||||
+ "stays unrecognised, and a caller from its new tab resolves as a worker, not a "
|
||||
+ "lead; the rest of fleet: (architects, developers, reviewers, charters, "
|
||||
+ "tabLabel) is read live through the supplier on CompositePeerLauncher and "
|
||||
+ "already applied");
|
||||
}
|
||||
// Kept in step with SPLIT_KEYS the same way changedColdKeys is kept in step with COLD_KEYS —
|
||||
// every message here must be traceable to one of the split keys the class doc documents.
|
||||
// NOTE what this does NOT prove, per the javadoc above: it does not catch a SPLIT_KEYS
|
||||
// member with no branch above at all, only a branch whose message is mis-worded relative to
|
||||
// the set. ConfigRefTopLevelReportingCoverageTest is what proves the former.
|
||||
assert changed.stream().allMatch(m -> SPLIT_KEYS.stream().anyMatch(k -> m.startsWith(k + ":")))
|
||||
: "a split entry was reported that does not start with a SPLIT_KEYS name: " + changed;
|
||||
return changed;
|
||||
}
|
||||
|
||||
/** {@code cfg.fleet().leaders()}, defensively, in case a caller hands in a non-defaulted config. */
|
||||
private static Map<String, FleetConfig.Leader> leadersOf(FleetConfig cfg) {
|
||||
return cfg.fleet() == null ? Map.of() : cfg.fleet().leaders();
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link FleetConfig.Profile} record components deliberately left out of
|
||||
* {@link #sameLaunchSettings} because they are read <em>live</em>, not baked in at spawn — see
|
||||
|
||||
@@ -1,10 +1,12 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
/**
|
||||
* Notified when {@link CompletionResolver} actually delivers a typed backend-error classification
|
||||
* to a waiting send (fleetd#201 / #227) — never on a race that lost. {@link CompletionResolver}
|
||||
* calls this only after {@code Rendezvous.resolveFailure} returns {@code true} for that exact
|
||||
* waiter, mirroring the win-only race rule {@link ExhaustionSink} already uses.
|
||||
* Notified when {@link CompletionResolver} has a backend-error match at the start of a pane line,
|
||||
* or both a match and its too-fast crash signature, for a waiting send (fleetd#201 / #227). A text
|
||||
* match inside ordinary pane prose can be a member's report about an error, so it fails the send
|
||||
* without notifying this sink.
|
||||
* {@link CompletionResolver} calls this only after {@code Rendezvous.resolveFailure} returns
|
||||
* {@code true} for that exact waiter, mirroring the win-only race rule {@link ExhaustionSink} uses.
|
||||
*
|
||||
* <p>The public send result is unchanged by this classification — it is still a failed send
|
||||
* ({@code Rendezvous.Kind#FAILED}); this sink is the internal seam a later stage (fleetd#201 Unit
|
||||
|
||||
@@ -387,9 +387,10 @@ public final class CompletionResolver implements TurnListener {
|
||||
// fleetd#164 (part 2) / fleetd#201: a scrape that read cleanly and produced content still
|
||||
// isn't a real reply when that content is the backend's own rejection (e.g. an HTTP 400
|
||||
// before the worker did any work). Classify it as a failure naming the member, rather than
|
||||
// handing the caller a scrape that reads like a completed answer, and — only on the
|
||||
// resolution that actually wins the race, mirroring the exhaustion sink above — notify the
|
||||
// typed backend-error sink so a later stage can act on repeated failures.
|
||||
// handing the caller a scrape that reads like a completed answer. A text match alone is not
|
||||
// enough to notify the typed backend-error sink: this assistant block can be a member's
|
||||
// normal prose about an error. A line that starts with the error match is stronger evidence;
|
||||
// the too-fast path below also has its crash signature before it records a credential failure.
|
||||
String backendError = firstMatchingLine(assistantBlock, backendErrorPatternOrFallback(target));
|
||||
if (backendError != null) {
|
||||
// Carry the whole scrape, not just the matched line. The pattern is a heuristic: a member
|
||||
@@ -401,9 +402,9 @@ public final class CompletionResolver implements TurnListener {
|
||||
if (rendezvous.resolveFailure(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("failing send to {} via turn-stall fallback: {}", target, reason);
|
||||
// fleetd#201 Unit 1: only on the resolution that actually won the race — a late
|
||||
// duplicate must never double-count one backend failure.
|
||||
backendErrorSink.onBackendError(target, backendError, reason);
|
||||
if (startsWithBackendError(backendError, backendErrorPatternOrFallback(target))) {
|
||||
backendErrorSink.onBackendError(target, backendError, reason);
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
@@ -492,8 +493,9 @@ public final class CompletionResolver implements TurnListener {
|
||||
if (rendezvous.resolveFailure(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("failing send to {} via turn-stall fallback from the raw scrape: {}", target, reason);
|
||||
// fleetd#201 Unit 1: only on the resolution that actually won the race.
|
||||
backendErrorSink.onBackendError(target, backendError, reason);
|
||||
if (startsWithBackendError(backendError, backendErrorPatternOrFallback(target))) {
|
||||
backendErrorSink.onBackendError(target, backendError, reason);
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
@@ -540,10 +542,10 @@ public final class CompletionResolver implements TurnListener {
|
||||
* {@link #MIN_TURN_NANOS} — a crash signature (e.g. a backend HTTP 400 before the worker did
|
||||
* anything) that a bare {@code BUSY -> DONE} transition cannot be told apart from a genuinely
|
||||
* fast completion. Runs the same backend-error classification the normal and raw-scrape paths
|
||||
* apply, against whatever is on screen right now: a match is a typed failure that notifies
|
||||
* {@link #backendErrorSink} (only on the resolution that wins the race); a non-match stays the
|
||||
* original generic too-fast failure, naming the member and both timings, with whatever the pane
|
||||
* shows appended so the caller sees the cause, not just "it failed".
|
||||
* apply, against whatever is on screen right now: a match together with the too-fast crash
|
||||
* signature notifies {@link #backendErrorSink} (only on the resolution that wins the race). A
|
||||
* non-match stays the original generic too-fast failure, naming the member and both timings,
|
||||
* with whatever the pane shows appended so the caller sees the cause, not just "it failed".
|
||||
*/
|
||||
private void failTooFast(String target, InFlight turn, CompletableFuture<Rendezvous.Resolution> waiter,
|
||||
long elapsedNanos) {
|
||||
@@ -611,6 +613,28 @@ public final class CompletionResolver implements TurnListener {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* True when the error pattern begins the matched pane line, rather than appearing in prose.
|
||||
*
|
||||
* <p>Leading terminal chrome is skipped first — box-drawing characters, bullets, gutter bars and
|
||||
* spaces. #339 introduced this check with a bare {@code lookingAt}, and that rejected a genuine
|
||||
* error line rendered as {@code "| 503 Service Unavailable: ..."}: the send still failed, but the
|
||||
* credential outage was never recorded. That is the false negative #339's own invariant 3 called
|
||||
* worse than the false positive it set out to fix — measured with a throwaway probe on the
|
||||
* raw-scrape path, which is exactly the path whose comment says to expect leading chrome.
|
||||
*
|
||||
* <p>Skipping only a leading run of non-letter, non-digit characters keeps the fix's intent. A
|
||||
* member's prose ({@code "I checked the retry path. An API Error: makes it back off."}) still
|
||||
* does not match, because there the pattern sits after words, not after chrome.
|
||||
*/
|
||||
private static boolean startsWithBackendError(String line, Pattern pattern) {
|
||||
int i = 0;
|
||||
while (i < line.length() && !Character.isLetterOrDigit(line.charAt(i))) {
|
||||
i++;
|
||||
}
|
||||
return pattern.matcher(line.substring(i)).lookingAt();
|
||||
}
|
||||
|
||||
/**
|
||||
* Coverage summary for the CB-578 stage A exhausted-pattern classification, logged at startup
|
||||
* the way {@link dev.ltms.fleet.health.FleetHealthMonitor#coverage} is — so an operator can
|
||||
|
||||
@@ -145,8 +145,58 @@ public final class Injector {
|
||||
return router != null ? router.agentsFor(target) : agents;
|
||||
}
|
||||
|
||||
/** The result of trying to remove an undelivered message from the injector. */
|
||||
public enum Cancellation {
|
||||
CANCELLED,
|
||||
DELIVERED,
|
||||
NOT_DELIVERED
|
||||
}
|
||||
|
||||
/**
|
||||
* An identity handle for one queued delivery. It is the only value accepted by
|
||||
* {@link #cancel(Delivery)}, so a caller cannot cancel a different message with the same target
|
||||
* or text.
|
||||
*/
|
||||
public static final class Delivery {
|
||||
private final Pending pending;
|
||||
|
||||
private Delivery(Pending pending) {
|
||||
this.pending = pending;
|
||||
}
|
||||
|
||||
public CompletableFuture<Void> completion() {
|
||||
return pending.delivered;
|
||||
}
|
||||
}
|
||||
|
||||
/** A pending message and the future that completes when it has been delivered. */
|
||||
private record Pending(String text, TurnToken token, CompletableFuture<Void> delivered) {
|
||||
private static final class Pending {
|
||||
enum State { QUEUED, DELIVERED, NOT_DELIVERED, CANCELLED }
|
||||
|
||||
final String target;
|
||||
final String text;
|
||||
final TurnToken token;
|
||||
final CompletableFuture<Void> delivered;
|
||||
volatile State state = State.QUEUED; // written under the owning Target monitor
|
||||
|
||||
Pending(String target, String text, TurnToken token, CompletableFuture<Void> delivered) {
|
||||
this.target = target;
|
||||
this.text = text;
|
||||
this.token = token;
|
||||
this.delivered = delivered;
|
||||
}
|
||||
|
||||
String text() {
|
||||
return text;
|
||||
}
|
||||
|
||||
TurnToken token() {
|
||||
return token;
|
||||
}
|
||||
|
||||
CompletableFuture<Void> delivered() {
|
||||
return delivered;
|
||||
}
|
||||
}
|
||||
|
||||
/** Per-worker delivery state, guarded by its own monitor (single writer per worker). */
|
||||
@@ -177,15 +227,47 @@ public final class Injector {
|
||||
* <p>Uses an atomic map update so a concurrent {@link #drop} cannot slip between "find the
|
||||
* target" and "queue the message" and orphan it in a target it just removed.
|
||||
*/
|
||||
public CompletableFuture<Void> enqueue(String target, String text, TurnToken token) {
|
||||
public Delivery enqueue(String target, String text, TurnToken token) {
|
||||
CompletableFuture<Void> delivered = new CompletableFuture<>();
|
||||
Pending p = new Pending(text, token, delivered);
|
||||
Pending p = new Pending(target, text, token, delivered);
|
||||
targets.compute(target, (_, existing) -> {
|
||||
Target t = (existing != null) ? existing : new Target();
|
||||
t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop
|
||||
return t;
|
||||
});
|
||||
return delivered;
|
||||
return new Delivery(p);
|
||||
}
|
||||
|
||||
/**
|
||||
* Cancel this exact queued delivery. The target monitor serializes this operation with
|
||||
* {@link #onStatus}: if delivery wins that race, this returns {@link Cancellation#DELIVERED}
|
||||
* rather than claiming the message remained queued.
|
||||
*/
|
||||
public Cancellation cancel(Delivery delivery) {
|
||||
Pending p = delivery.pending;
|
||||
Target t = targets.get(p.target);
|
||||
if (t == null) {
|
||||
return cancellationOf(p);
|
||||
}
|
||||
synchronized (t) {
|
||||
if (p.state != Pending.State.QUEUED || !t.queue.remove(p)) {
|
||||
return cancellationOf(p);
|
||||
}
|
||||
p.state = Pending.State.CANCELLED;
|
||||
if (isQuiescent(t)) {
|
||||
targets.remove(p.target, t);
|
||||
}
|
||||
return Cancellation.CANCELLED;
|
||||
}
|
||||
}
|
||||
|
||||
private static Cancellation cancellationOf(Pending p) {
|
||||
return p.state == Pending.State.DELIVERED ? Cancellation.DELIVERED : Cancellation.NOT_DELIVERED;
|
||||
}
|
||||
|
||||
private static boolean isQuiescent(Target t) {
|
||||
return t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion
|
||||
&& !t.postTurnPending && !t.awaitingPostTurnPickup && !t.postTurnObserved;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -274,6 +356,7 @@ public final class Injector {
|
||||
try {
|
||||
agentsFor(target).send(target, p.text());
|
||||
t.queue.poll();
|
||||
p.state = Pending.State.DELIVERED;
|
||||
t.awaitingPickup = true;
|
||||
t.awaitingCompletion = true;
|
||||
t.turnObserved = false;
|
||||
@@ -283,6 +366,7 @@ public final class Injector {
|
||||
// Delivery failed at herdr; drop the poisoned message and surface it
|
||||
// rather than blocking the queue behind it.
|
||||
t.queue.poll();
|
||||
p.state = Pending.State.NOT_DELIVERED;
|
||||
sent = p;
|
||||
sendError = e;
|
||||
}
|
||||
@@ -293,6 +377,9 @@ public final class Injector {
|
||||
// fail every queued message and release the target (CB-114) instead of
|
||||
// polling it indefinitely with the caller's future never completing.
|
||||
notReady = new ArrayList<>(t.queue);
|
||||
for (Pending pending : notReady) {
|
||||
pending.state = Pending.State.NOT_DELIVERED;
|
||||
}
|
||||
log.warn("readiness grace for {} expired after {} polls ({}s): target never "
|
||||
+ "became deliverable, so failing {} queued message(s) that never "
|
||||
+ "reached its pane",
|
||||
@@ -340,8 +427,7 @@ public final class Injector {
|
||||
|
||||
// Reclaim the entry once the worker is fully quiescent (nothing queued, no pickup or
|
||||
// completion awaited), so the map cannot grow without bound across short-lived workers.
|
||||
if (t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion
|
||||
&& !t.postTurnPending && !t.awaitingPostTurnPickup && !t.postTurnObserved) {
|
||||
if (isQuiescent(t)) {
|
||||
targets.remove(target, t);
|
||||
}
|
||||
}
|
||||
@@ -432,6 +518,9 @@ public final class Injector {
|
||||
boolean hadDeliveredTurn;
|
||||
synchronized (t) {
|
||||
pending = new ArrayList<>(t.queue);
|
||||
for (Pending p : pending) {
|
||||
p.state = Pending.State.NOT_DELIVERED;
|
||||
}
|
||||
t.queue.clear();
|
||||
hadDeliveredTurn = t.awaitingCompletion;
|
||||
t.awaitingCompletion = false;
|
||||
|
||||
@@ -1665,30 +1665,59 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
|
||||
/**
|
||||
* Guards the {@code effectiveAllowed == null} branch of {@link #logCredentialGap} — the
|
||||
* genuinely-unprotected report (deny-by-default, and the allow-list non-zsh fallback) — to one
|
||||
* WARN per launcher instance, not one per spawn.
|
||||
* genuinely-unprotected report (deny-by-default, and the allow-list non-zsh fallback) — AND
|
||||
* the {@code effectiveAllowed != null} / {@code keptByDerivedList} branch, the allow-list case
|
||||
* where a name is in the gap but the derived allow-list keeps it anyway. Both branches log the
|
||||
* same severity (WARN) about the same fact — a name genuinely reaching a member pane
|
||||
* unprotected — so they share this one guard, keyed per NAME rather than per launcher instance:
|
||||
* each credential-shaped name that is ever reported unprotected gets exactly one WARN, however
|
||||
* many spawns see it and whichever of the two branches first reports it.
|
||||
*
|
||||
* <p>fleetd #341: {@code memberCredentials} is a live, re-read-per-spawn supplier, so the
|
||||
* policy — and so the gap's actual member names — can change between two spawns on the same
|
||||
* launcher. Before this fix the guard was a single {@code AtomicBoolean} tripped by either
|
||||
* branch: spawn 1 could warn about name A and trip the flag, and a later spawn's gap containing
|
||||
* a different name B would never be reported, even though B is just as unprotected as A was.
|
||||
* {@code AtomicBoolean} could not express "once per distinct name" at all — only "once, ever,
|
||||
* for whichever name got there first" — so this is a {@code Set<String>} guard instead, the
|
||||
* same shape {@link OpenCodeLauncher#modelCheckSkippedWarned} already uses for its own
|
||||
* once-per-distinct-thing WARN. {@link #add}'s return value (true only the first time a name is
|
||||
* added) is what turns "log the whole gap" into "log only the names never warned about before".
|
||||
*
|
||||
* <p>Bounded by construction: every name added here first passed {@link
|
||||
* #CREDENTIAL_SHAPED_NAME}'s filter over {@link #hostEnvNames}, i.e. it is an actual
|
||||
* environment variable name from the daemon's own process — a small, OS-bounded set (the host
|
||||
* environment has, in practice, tens to a few hundred entries), not an attacker- or
|
||||
* request-controlled input. So this set cannot grow past "however many distinct credential-
|
||||
* shaped names this host's environment has ever held across this launcher's lifetime," which is
|
||||
* effectively fixed for the life of one daemon process — no separate cap is needed.
|
||||
*
|
||||
* <p>CB-633 follow-up (#192): kept SEPARATE from {@link #allowListGapLogged} on purpose.
|
||||
* {@code memberCredentials} is a live, re-read-per-spawn supplier, so the policy can change
|
||||
* between two spawns on the same launcher. A single shared flag would let a harmless allow-list
|
||||
* INFO on spawn 1 permanently suppress the real deny-by-default WARN a later spawn deserves —
|
||||
* the report that matters most getting hidden by the report that doesn't. Two flags mean each
|
||||
* report kind fires exactly once, independent of what the other kind already logged.
|
||||
* report kind (WARN vs. INFO) fires independently of what the other kind already logged; within
|
||||
* the WARN kind itself, the set above further separates by name, for the same reason.
|
||||
*/
|
||||
private final AtomicBoolean unprotectedGapLogged = new AtomicBoolean();
|
||||
private final Set<String> unprotectedGapNamesWarned = ConcurrentHashMap.newKeySet();
|
||||
|
||||
/**
|
||||
* Guards the {@code effectiveAllowed != null} branch of {@link #logCredentialGap} — the
|
||||
* allow-list-scrub-covered report — to one INFO per launcher instance. See {@link
|
||||
* #unprotectedGapLogged}'s javadoc for why this is a separate flag rather than a shared one.
|
||||
* #unprotectedGapNamesWarned}'s javadoc for why this is a separate flag rather than a shared
|
||||
* one; unlike that guard it stays a per-instance {@code AtomicBoolean}, not a per-name set —
|
||||
* fleetd #341 fixed the WARN-vs-WARN suppression, not this INFO's own one-shot shape, which was
|
||||
* not reported as broken and is out of that ticket's scope.
|
||||
*/
|
||||
private final AtomicBoolean allowListGapLogged = new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: guards {@link #warnUnknownMemberEnvironment} to one WARN per launcher
|
||||
* instance, not one per spawn — the same one-per-instance shape as {@link #unprotectedGapLogged}
|
||||
* and {@link #allowListGapLogged}, kept as its own flag for the same reason those two are split:
|
||||
* this mode is orthogonal to which of the other two branches would otherwise have fired.
|
||||
* instance, not one per spawn — the same one-shot shape {@link #unprotectedGapNamesWarned} and
|
||||
* {@link #allowListGapLogged} guard their own branches with, kept as its own flag for the same
|
||||
* reason those two are split: this mode is orthogonal to which of the other two branches would
|
||||
* otherwise have fired.
|
||||
*/
|
||||
private final AtomicBoolean unknownMemberEnvironmentWarned = new AtomicBoolean();
|
||||
|
||||
@@ -1816,14 +1845,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
List<String> blankedByScrub = gap.stream()
|
||||
.filter(name -> !MemberEnvAllowList.keeps(effectiveAllowed, name))
|
||||
.toList();
|
||||
if (!keptByDerivedList.isEmpty() && unprotectedGapLogged.compareAndSet(false, true)) {
|
||||
// fleetd #341: filter to names this guard has never warned about before — not just
|
||||
// "isEmpty" on the whole branch — so a name this spawn's gap shares with an EARLIER
|
||||
// spawn's (already-warned) gap does not re-print, while a name unique to THIS gap still
|
||||
// does, whichever of the two WARN branches reported it first.
|
||||
List<String> newlyUnprotected = keptByDerivedList.stream()
|
||||
.filter(unprotectedGapNamesWarned::add)
|
||||
.toList();
|
||||
if (!newlyUnprotected.isEmpty()) {
|
||||
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||
+ "known: nor allow: — the derived allow-list keeps them anyway (a profile's "
|
||||
+ "gitTokenEnv/gitHostEnv/tokenEnv/env: names one, or this spawn injects it), "
|
||||
+ "so every member pane inherits them UNBLOCKED — {}. Add each to "
|
||||
+ "memberCredentials.known (or .allow if a member legitimately needs it), or "
|
||||
+ "remove it from whatever profile setting derives it in.",
|
||||
keptByDerivedList.size(), keptByDerivedList);
|
||||
newlyUnprotected.size(), newlyUnprotected);
|
||||
}
|
||||
if (!blankedByScrub.isEmpty() && allowListGapLogged.compareAndSet(false, true)) {
|
||||
log.info("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||
@@ -1836,12 +1872,18 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
|
||||
/** The deny-by-default (and allow-list non-zsh fallback) WARN — unchanged byte-for-byte by #192. */
|
||||
private void warnGapUnprotected(List<String> gap) {
|
||||
if (unprotectedGapLogged.compareAndSet(false, true)) {
|
||||
// fleetd #341: same "only the names never warned before" filter as the sibling branch in
|
||||
// logCredentialGap above — see unprotectedGapNamesWarned's javadoc. Both branches share
|
||||
// this one guard because both report the exact same fact (a name reaching a member pane
|
||||
// unprotected) at the exact same severity; keying it by name is what lets a later spawn's
|
||||
// DIFFERENT name still get its own WARN after an earlier spawn's already fired.
|
||||
List<String> newlyUnprotected = gap.stream().filter(unprotectedGapNamesWarned::add).toList();
|
||||
if (!newlyUnprotected.isEmpty()) {
|
||||
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||
+ "known: nor allow: — every member pane inherits them UNBLOCKED — {}. "
|
||||
+ "Add each to memberCredentials.known (blocked by default) or .allow "
|
||||
+ "(if a member legitimately needs it).",
|
||||
gap.size(), gap);
|
||||
newlyUnprotected.size(), newlyUnprotected);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -22,6 +22,7 @@ import java.util.concurrent.ConcurrentSkipListMap;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
/**
|
||||
* AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same
|
||||
@@ -93,8 +94,23 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
private final Channel channel;
|
||||
/** All channel operations (publish/declare/ack/cancel) serialize on this — a Channel is not thread-safe. */
|
||||
private final Object channelLock = new Object();
|
||||
/** target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself. */
|
||||
/**
|
||||
* target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself.
|
||||
*
|
||||
* <p><strong>CB-318 tombstone.</strong> The value {@link #RELEASED} is a reserved sentinel: it
|
||||
* marks a target whose {@link #release} has already run, so {@link #deliverCallback} can tell a
|
||||
* delivery landing after release() apart from a fresh target it has never seen. See both methods'
|
||||
* javadoc for why a plain {@code held.remove(target)} is not enough.
|
||||
*/
|
||||
private final ConcurrentHashMap<String, LinkedHashMap<String, Held>> held = new ConcurrentHashMap<>();
|
||||
|
||||
/**
|
||||
* CB-318 sentinel stored in {@link #held} for a target whose {@link #release} has already run.
|
||||
* Never mutated — every read site compares it by reference ({@code ==}) before touching it as a
|
||||
* map, because it is a single object shared across every released target and calling a mutator on
|
||||
* it would corrupt state for all of them.
|
||||
*/
|
||||
private static final LinkedHashMap<String, Held> RELEASED = new LinkedHashMap<>();
|
||||
/** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
|
||||
private final ConcurrentHashMap<String, String> consumerTags = new ConcurrentHashMap<>();
|
||||
|
||||
@@ -201,6 +217,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
|
||||
String tag = channel.basicConsume(queue, false, deliverCallback(target), _ -> { });
|
||||
consumerTags.put(target, tag);
|
||||
// CB-318: drop a stale RELEASED tombstone from a prior ownership of this same target
|
||||
// string, so a delivery under this fresh consumer is held normally instead of being
|
||||
// nacked forever by deliverCallback's RELEASED check. Safe to do here, still under
|
||||
// channelLock: no delivery for the consumer tag just registered above can reach
|
||||
// deliverCallback before this basicConsume call returns.
|
||||
held.remove(target, RELEASED);
|
||||
log.debug("AMQP inbox owns queue {} for target {}", queue, target);
|
||||
} catch (IOException e) {
|
||||
throw new IllegalStateException("cannot own queue " + queue, e);
|
||||
@@ -243,6 +265,32 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
* later connection drop even though this release did not manage to requeue it immediately. A
|
||||
* failed {@code basicCancel} still throws, unchanged from before this fix — that failure means
|
||||
* the consumer may still be attached, so best-effort requeue is not attempted underneath it.
|
||||
*
|
||||
* <p><strong>CB-318: {@code held.remove(target)} alone leaves a second window open.</strong> The
|
||||
* bullet above already explains why cancelling first does not save a tag from going stale — but
|
||||
* that only accounts for a delivery landing before this method starts touching {@link #held}.
|
||||
* {@code basicCancel} stops <em>new</em> dispatches; it does not flush one already handed to the
|
||||
* consumer work pool. So a delivery can still land on that pool's thread and reach
|
||||
* {@link #deliverCallback} at any point during, or after, this method's body — and a plain
|
||||
* {@code held.remove(target)} does nothing to stop it: {@code deliverCallback}'s
|
||||
* {@code computeIfAbsent} finds the key gone and happily creates a brand-new map under it, which
|
||||
* this method — already past its {@code remove} — never looks at again. That entry then sits
|
||||
* delivered-but-unacked on {@link #channel} until the whole inbox closes: never requeued, never
|
||||
* redelivered, and {@link #peek} is never called again for a target nothing owns any more.
|
||||
*
|
||||
* <p>The fix is {@link #held}{@code .compute(target, ...)} instead of {@code remove}: it takes
|
||||
* whatever was held (to nack, same as before) and, in the same atomic step, leaves the
|
||||
* {@link #RELEASED} tombstone behind instead of an absent key. {@code computeIfAbsent} and
|
||||
* {@code compute} calls for the same key are mutually exclusive in {@link ConcurrentHashMap} —
|
||||
* whichever of this call and a concurrent {@code deliverCallback} runs first is fully visible to
|
||||
* the other, with no gap between them. So a delivery that loses the race sees a real map here and
|
||||
* gets nacked by the loop below, same as always; a delivery that wins the race (runs first) is
|
||||
* itself nacked by that same loop, once it settles into {@code held}. A delivery that arrives once
|
||||
* this method has stored {@link #RELEASED} finds it via {@code computeIfAbsent} and refuses itself
|
||||
* — see {@link #deliverCallback}. Either way nothing is silently retained forever, satisfying the
|
||||
* ticket's invariant against dropping a message. This closes the window rather than merely
|
||||
* narrowing it — correctness does not depend on how much time elapses between the swap and this
|
||||
* method returning.
|
||||
*/
|
||||
@Override
|
||||
public void release(String target) {
|
||||
@@ -255,8 +303,13 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
throw new IllegalStateException("cannot cancel consumer for " + target, e);
|
||||
}
|
||||
}
|
||||
var perTarget = held.remove(target);
|
||||
if (perTarget != null) {
|
||||
AtomicReference<LinkedHashMap<String, Held>> previouslyHeld = new AtomicReference<>();
|
||||
held.compute(target, (_, v) -> {
|
||||
previouslyHeld.set(v);
|
||||
return RELEASED;
|
||||
});
|
||||
var perTarget = previouslyHeld.get();
|
||||
if (perTarget != null && perTarget != RELEASED) {
|
||||
synchronized (perTarget) {
|
||||
for (Held h : perTarget.values()) {
|
||||
try {
|
||||
@@ -326,7 +379,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
@Override
|
||||
public List<InboxMessage> peek(String target) {
|
||||
var perTarget = held.get(target);
|
||||
if (perTarget == null) {
|
||||
if (perTarget == null || perTarget == RELEASED) {
|
||||
return List.of();
|
||||
}
|
||||
synchronized (perTarget) {
|
||||
@@ -337,7 +390,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
@Override
|
||||
public void ack(String target, String msgId) {
|
||||
var perTarget = held.get(target);
|
||||
if (perTarget == null) {
|
||||
if (perTarget == null || perTarget == RELEASED) {
|
||||
return;
|
||||
}
|
||||
Held h;
|
||||
@@ -370,6 +423,20 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
}
|
||||
String content = new String(delivery.getBody(), StandardCharsets.UTF_8);
|
||||
var perTarget = held.computeIfAbsent(target, _ -> new LinkedHashMap<>());
|
||||
if (perTarget == RELEASED) {
|
||||
// CB-318: release() already ran for this target and left the RELEASED tombstone in
|
||||
// held (see release()'s javadoc) — computeIfAbsent() is guaranteed to see it rather
|
||||
// than recreate a fresh map, because ConcurrentHashMap serializes compute/
|
||||
// computeIfAbsent calls for the same key against each other. Refuse the delivery
|
||||
// instead of holding it somewhere release() will never look at again: requeue it, the
|
||||
// same way release() nacks its own held entries, so a later owner (or a connection
|
||||
// drop) can still recover it. This does not need channelLock across a broker round
|
||||
// trip — basicNack, like the duplicate-ack case just below, does not wait for one.
|
||||
synchronized (channelLock) {
|
||||
channel.basicNack(tag, false, true);
|
||||
}
|
||||
return;
|
||||
}
|
||||
boolean duplicate;
|
||||
synchronized (perTarget) {
|
||||
if (perTarget.containsKey(msgId)) {
|
||||
|
||||
@@ -466,8 +466,18 @@ public final class MessageService {
|
||||
if (candidates.size() == 1) {
|
||||
Task orphan = candidates.get(0);
|
||||
if (orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
|
||||
if (orphan.turnId != null) {
|
||||
asyncTasksByTurn.remove(orphan.turnId, orphan);
|
||||
// fleetd #329 (F3): read orphan.turnId once. It used to be read twice, under no
|
||||
// lock — once for this null check, once as the removal key — the same double-read
|
||||
// shape fleetd #324 fixed in finishAsyncTask. clearAsyncQuestion's unlocked
|
||||
// forgetTurn=true path (ask()'s timeout cleanup) can null the field between the two
|
||||
// reads; capturing it once removes the torn read here too.
|
||||
String turnId = orphan.turnId;
|
||||
if (turnId != null) {
|
||||
if (replyOrphanTurnIdRaceHookForTest != null) {
|
||||
// Test-only (fleetd #329, F3): see the field's own javadoc.
|
||||
replyOrphanTurnIdRaceHookForTest.run();
|
||||
}
|
||||
asyncTasksByTurn.remove(turnId, orphan);
|
||||
}
|
||||
count(FleetMetrics.REPLIES, "path", "async-recovered");
|
||||
return true; // the ticket itself took it — no inbox stranding at all
|
||||
@@ -853,16 +863,26 @@ public final class MessageService {
|
||||
if (onAccepted != null) {
|
||||
onAccepted.run();
|
||||
}
|
||||
CompletableFuture<Void> delivered = injector.enqueue(target, content, token);
|
||||
Injector.Delivery delivery = injector.enqueue(target, content, token);
|
||||
try {
|
||||
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
||||
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
|
||||
} catch (TimeoutException e) {
|
||||
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
|
||||
boolean wasDelivered = delivery.completion().isDone()
|
||||
&& !delivery.completion().isCompletedExceptionally();
|
||||
if (!wasDelivered) {
|
||||
if (timeoutCancellationRaceHookForTest != null) {
|
||||
// Test-only (fleetd #345): see the field's own javadoc.
|
||||
timeoutCancellationRaceHookForTest.run();
|
||||
}
|
||||
// The target monitor makes cancellation atomic with onStatus picking this
|
||||
// Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
|
||||
wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
|
||||
}
|
||||
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
|
||||
if (!wasDelivered) {
|
||||
// CB-640: still sitting in the injector's queue, waiting for the member to
|
||||
// go idle — record the fact for fleet health (see queuedDeliveries).
|
||||
// CB-640: record that delivery did not happen for fleet health (see
|
||||
// queuedDeliveries). The exact Pending was cancelled, so it cannot arrive later.
|
||||
queuedDeliveries.put(target, Boolean.TRUE);
|
||||
}
|
||||
return recorded(new Reply(
|
||||
@@ -1002,13 +1022,49 @@ public final class MessageService {
|
||||
// Measured when #282 was merged: this guard is DEFENCE IN DEPTH, not the thing
|
||||
// that makes the chained ask work. ask() calls markAsyncQuestion (:860) before
|
||||
// resolveQuestion (:861), so by the time this thread wakes, the task has already
|
||||
// moved to the new turnId and finishAsyncTask(oldTurnId, ...) finds nothing. Removing
|
||||
// this guard alone leaves the test green. Keep it anyway: it mirrors sendAsync's
|
||||
// sibling guard, and that sibling's own comment (:1017) warns the two orderings are
|
||||
// not something to rely on. Do NOT delete it as dead code without re-checking that
|
||||
// ordering, and do not treat it as the sole protection either.
|
||||
if (result.outcome() != Outcome.QUESTION) {
|
||||
finishAsyncTask(turnId, result);
|
||||
// moved to the new turnId — this outcome check, not a Task lookup, is what tells the
|
||||
// two cases apart (see fleetd #329 below). Removing this guard alone leaves the test
|
||||
// green. Keep it anyway: it mirrors sendAsync's sibling guard, and that sibling's own
|
||||
// comment (:1017) warns the two orderings are not something to rely on. Do NOT delete
|
||||
// it as dead code without re-checking that ordering, and do not treat it as the sole
|
||||
// protection either.
|
||||
//
|
||||
// fleetd #329 (F1): complete the SAME Task object this method already looked up at
|
||||
// :991, rather than re-resolving it from turnId a second time. The old
|
||||
// finishAsyncTask(turnId, result) did its own asyncTasksByTurn.get(turnId) here, and
|
||||
// that second lookup races ask()'s own timeout path: ask()'s ticket.answer().get(...)
|
||||
// can time out (or lose that exact race) at essentially the same instant this method's
|
||||
// rendezvous.answerAsk(turnId, ...) above already succeeded, running
|
||||
// markAskTimedOut + clearAsyncQuestion(turnId, true) with no lock at all — which
|
||||
// forgets turnId (removes it from asyncTasksByTurn, nulls Task.turnId) before this
|
||||
// thread ever gets here. The worker's real reply then arrived, this method's own wait
|
||||
// woke up with it, and the by-turnId lookup found nothing: the async ticket's future
|
||||
// was never completed, so fleet_poll{ticket} stayed PENDING forever even though
|
||||
// answer() itself correctly returned REPLIED. Reusing the reference captured at :991
|
||||
// — before that race window opens — sidesteps the second lookup entirely: it is the
|
||||
// identical Task whatever asyncTasksByTurn or Task.turnId say by the time we reach
|
||||
// this point, so finishAsyncTask(Task, Reply) can still complete its future and detach
|
||||
// it using whatever turnId it now reads. The #282 chained-ask case is unaffected
|
||||
// because it is still gated purely by result.outcome() == QUESTION above, which does
|
||||
// not depend on this lookup — widening what "task" means here cannot complete a ticket
|
||||
// the chained ask deliberately left open.
|
||||
//
|
||||
// A null task is NOT only "this was never an async ticket". That reading was in this
|
||||
// comment when #329 merged and it is wrong. A genuine async ticket also lands here
|
||||
// with task == null, because ask()'s timeout path runs clearAsyncQuestion(turnId,
|
||||
// true) — which drops the asyncTasksByTurn entry — in its catch block, while
|
||||
// rendezvous.closeAsk(turnId) runs later, in its finally. Between those two the ask
|
||||
// is still answerable but the map entry is already gone, so the lookup at :991
|
||||
// returns null and this ticket is never completed. Measured on 2026-09-04: a probe
|
||||
// firing only that first half before answer() runs printed
|
||||
// "answer=REPLIED phase=PENDING reply=null" — the same stranded ticket #329 set out
|
||||
// to fix, one step earlier in the same race. The probe used forgetTurnForTest, which
|
||||
// omits ask()'s markAskTimedOut; that cannot change the outcome, because askTimedOut
|
||||
// is read only by askAnsweredAsyncTasks, and reply() never reaches it while this
|
||||
// method's own waiter is live. So #329 narrows this window rather than closing it.
|
||||
// Open as fleetd #334 — do not read this guard as complete.
|
||||
if (result.outcome() != Outcome.QUESTION && task != null) {
|
||||
finishAsyncTask(task, result);
|
||||
}
|
||||
return result;
|
||||
} catch (TimeoutException e) {
|
||||
@@ -1059,7 +1115,7 @@ public final class MessageService {
|
||||
// this fires exactly once, from whichever path completes it: finishAsyncTask(task, result)
|
||||
// below on any non-QUESTION outcome of send() — a worker's fleet_reply, the CB-106
|
||||
// completion fallback, a CB-109 wedge, TIMED_OUT, BUSY, or BACKEND_EXHAUSTED — the same
|
||||
// finishAsyncTask reached via answer()'s finishAsyncTask(turnId, result) once a QUESTION
|
||||
// finishAsyncTask reached via answer()'s own finishAsyncTask(task, result) once a QUESTION
|
||||
// is resolved, completeExceptionally(t) just below when send() itself throws, or a CB-516
|
||||
// abandon() on teardown. Without this, MessageService.reply's rendezvous fast path (the
|
||||
// one an async ticket always takes) never told the push loop anything happened — see the
|
||||
@@ -1079,7 +1135,22 @@ public final class MessageService {
|
||||
finishAsyncTask(task, result);
|
||||
}
|
||||
} catch (Throwable t) {
|
||||
task.future.completeExceptionally(t);
|
||||
// fleetd #329 (F2): finishAsyncTask above already completes task.future — on its very
|
||||
// first line — before doing anything else, so anything that throws afterward (inside
|
||||
// finishAsyncTask's own cleanup, or from a future addition to this try block) lands
|
||||
// here with the future already resolved. completeExceptionally on an already-completed
|
||||
// future is a silent no-op: it returns false and does nothing, so without the check
|
||||
// below the exception simply vanished — no log, no metric, nothing. Measured (see the
|
||||
// ticket): temporarily reintroducing the fleetd #324 NPE reproduced 19 real exceptions
|
||||
// on the ordinary path, with 74/74 tests staying green and not one log line produced.
|
||||
// Do not stop completing the future first — finishAsyncTask completing before it
|
||||
// cleans up is what makes a late failure harmless to the ticket's own result — only
|
||||
// add the missing visibility for the case where that step, or whatever ran after it,
|
||||
// has already lost the race to report through the future.
|
||||
if (!task.future.completeExceptionally(t)) {
|
||||
log.error("async send {} -> {} threw after its ticket was already resolved",
|
||||
ticket, target, t);
|
||||
}
|
||||
}
|
||||
});
|
||||
pruneTerminalTickets();
|
||||
@@ -1230,20 +1301,124 @@ public final class MessageService {
|
||||
}
|
||||
}
|
||||
|
||||
/** Complete and detach an async ticket after its worker's actual terminal reply. */
|
||||
/**
|
||||
* Complete and detach an async ticket after its worker's actual terminal reply.
|
||||
*
|
||||
* <p><strong>fleetd #324.</strong> {@code task.turnId} is read into {@code turnId} exactly once.
|
||||
* It used to be read twice — once for the null check, once as the removal key — and {@code
|
||||
* volatile} makes each of those reads individually fresh but does not make the pair atomic.
|
||||
* {@link #answer} calls this while holding {@code sessionLocks} for the target; {@link #ask}'s
|
||||
* own timeout path calls {@link #clearAsyncQuestion} (which nulls {@link Task#turnId}) under no
|
||||
* lock at all. When that unlocked null-out landed between the two reads here, the second read saw
|
||||
* {@code null} and {@code asyncTasksByTurn.remove(null, task)} threw {@code NullPointerException}
|
||||
* on the lead's own {@code answer()} call — even though {@code task.future.complete(result)} on
|
||||
* the line above had already run, so the answer was in fact delivered. Capturing the field once
|
||||
* removes the torn read; see the ticket for why the wider asymmetry between the locked and
|
||||
* unlocked sides is not fixed by this alone.
|
||||
*/
|
||||
private void finishAsyncTask(Task task, Reply result) {
|
||||
task.future.complete(result);
|
||||
if (task.turnId != null) {
|
||||
asyncTasksByTurn.remove(task.turnId, task);
|
||||
if (afterFinishAsyncTaskCompleteHookForTest != null) {
|
||||
// Test-only (fleetd #329, F2): see the field's own javadoc.
|
||||
afterFinishAsyncTaskCompleteHookForTest.run();
|
||||
}
|
||||
String turnId = task.turnId;
|
||||
if (turnId != null) {
|
||||
if (finishAsyncTaskRaceHook != null) {
|
||||
// Test-only (fleetd #324): see the field's own javadoc.
|
||||
finishAsyncTaskRaceHook.run();
|
||||
}
|
||||
asyncTasksByTurn.remove(turnId, task);
|
||||
}
|
||||
}
|
||||
|
||||
/** Complete the async ticket correlated to a specific answered turn. */
|
||||
private void finishAsyncTask(String turnId, Reply result) {
|
||||
Task task = asyncTasksByTurn.get(turnId);
|
||||
if (task != null) {
|
||||
finishAsyncTask(task, result);
|
||||
}
|
||||
/**
|
||||
* Null in production; test seam for fleetd #324 — invoked from {@link #finishAsyncTask(Task,
|
||||
* Reply)} right after {@code task.turnId}'s null-check passes and before the (now-local) value is
|
||||
* used for the removal. A test installs this to force, deterministically, the exact interleaving
|
||||
* that a real race between this method and {@link #ask}'s unlocked timeout cleanup can otherwise
|
||||
* only produce by chance: firing it here reproduces "the field went null between the check and the
|
||||
* use" against the pre-fix code, and demonstrates the fix tolerates it (the captured local is used
|
||||
* unconditionally, so a hook that nulls the field afterward cannot affect this call).
|
||||
*/
|
||||
private volatile Runnable finishAsyncTaskRaceHook;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #324): install {@link #finishAsyncTaskRaceHook}. Package-private so the test,
|
||||
* in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setFinishAsyncTaskRaceHookForTest(Runnable hook) {
|
||||
this.finishAsyncTaskRaceHook = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #329 (F2) — invoked from {@link #finishAsyncTask(Task,
|
||||
* Reply)} unconditionally, immediately after {@code task.future.complete(result)} runs (before
|
||||
* {@link Task#turnId} is even read, so it fires regardless of whether this task was ever asked).
|
||||
* A test installs this to force an exception into exactly the shape fleetd #329 identified:
|
||||
* something throws inside {@link #sendAsync}'s executor task after the async ticket's future is
|
||||
* already resolved, so the surrounding {@code catch (Throwable t)} can only report failure through
|
||||
* {@code completeExceptionally} — a silent no-op on an already-completed future. Engineering a
|
||||
* real exception to land in that exact post-completion window is what the ticket itself had to do
|
||||
* by temporarily deleting a production guard (fleetd #324's {@code turnId != null} check); this
|
||||
* hook drives the identical shape deterministically instead.
|
||||
*/
|
||||
private volatile Runnable afterFinishAsyncTaskCompleteHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #329, F2): install {@link #afterFinishAsyncTaskCompleteHookForTest}.
|
||||
* Package-private so the test, in the same package, can reach it without widening any production
|
||||
* API.
|
||||
*/
|
||||
void setAfterFinishAsyncTaskCompleteHookForTest(Runnable hook) {
|
||||
this.afterFinishAsyncTaskCompleteHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #345 — invoked in {@link #send}'s timeout path after
|
||||
* {@link Injector.Delivery#completion()} reports incomplete and before {@link Injector#cancel}
|
||||
* takes the target monitor. A test installs this to make {@code onStatus} pick the exact queued
|
||||
* delivery up in that window, so {@code cancel} returns {@link Injector.Cancellation#DELIVERED}.
|
||||
* This deterministically covers the caller's need to use that result rather than relying on a
|
||||
* timing-sensitive real race.
|
||||
*/
|
||||
private volatile Runnable timeoutCancellationRaceHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #345): install {@link #timeoutCancellationRaceHookForTest}. Package-private
|
||||
* so the test, in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setTimeoutCancellationRaceHookForTest(Runnable hook) {
|
||||
this.timeoutCancellationRaceHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #329 (F3) — invoked from {@link #reply} right after
|
||||
* the single local read of {@code orphan.turnId} passes its null-check and before that (now-local)
|
||||
* value is used as the {@code asyncTasksByTurn} removal key. Mirrors {@link
|
||||
* #finishAsyncTaskRaceHook} exactly, for the structurally identical double-read fleetd #324 fixed
|
||||
* in {@link #finishAsyncTask}: a test installs this to force, deterministically, {@code
|
||||
* clearAsyncQuestion}'s unlocked {@code forgetTurn=true} path nulling {@link Task#turnId} in that
|
||||
* exact window, and to confirm the single-read fix tolerates it (the captured local is used
|
||||
* unconditionally, so a hook that nulls the field afterward cannot affect this call).
|
||||
*/
|
||||
private volatile Runnable replyOrphanTurnIdRaceHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #329, F3): install {@link #replyOrphanTurnIdRaceHookForTest}. Package-private
|
||||
* so the test, in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setReplyOrphanTurnIdRaceHookForTest(Runnable hook) {
|
||||
this.replyOrphanTurnIdRaceHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #324): run the exact production cleanup {@link #ask}'s own timeout path runs
|
||||
* unlocked — {@link #clearAsyncQuestion(String, boolean)} with {@code forgetTurn=true} — so a test
|
||||
* can reproduce that specific mutation instead of hand-rolling an approximation of it.
|
||||
*/
|
||||
void forgetTurnForTest(String turnId) {
|
||||
clearAsyncQuestion(turnId, true);
|
||||
}
|
||||
|
||||
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
|
||||
|
||||
@@ -386,6 +386,31 @@ public final class SessionManager implements TurnListener {
|
||||
// from the registry with no pane stop is an orphaned pane — a live terminal burning a fleet
|
||||
// slot that no longer appears in the roster and can never be reclaimed.
|
||||
launcher.stop(paneId);
|
||||
if (removed != null && !preserveWorktree && removed.worktree() != null) {
|
||||
// fleetd #316: the `dirty` read above ran while the worker could still write to this
|
||||
// worktree, so a stale `false` must not be trusted to authorise the --force removal
|
||||
// below. Re-read the worktree's state one more time, right here — immediately before
|
||||
// the one step that would destroy it, and only on the path that is actually about to
|
||||
// do that (invariant 4: no second unconditional `git status` on a release that already
|
||||
// decided to preserve). By now `launcher.stop` has returned, so this read reflects
|
||||
// whatever the worker managed to write up to and including its teardown, not whatever
|
||||
// it had written at release-start time.
|
||||
if (dirtyImmediatelyBeforeRemoval(removed)) {
|
||||
preserveWorktree = true;
|
||||
// The pre-stop snapshot above never ran for this session (the pre-stop read said
|
||||
// clean), so this is the only chance to get the newly-discovered work into
|
||||
// refs/wip/* rather than leaving the on-disk preserve as the sole copy. Best-effort,
|
||||
// like every other snapshot attempt — trySnapshot logs and swallows its own failure.
|
||||
String lateSnapshotRef = trySnapshot(removed, cause);
|
||||
log.warn("release {} preserves worktree {} for pane={} terminal={}: it reported "
|
||||
+ "clean before the pane stopped but dirty immediately before removal — the "
|
||||
+ "worker wrote to it during teardown, and --force removing it now would "
|
||||
+ "have destroyed that work{}",
|
||||
cause, removed.worktree(), paneId, removed.terminalId(),
|
||||
lateSnapshotRef == null ? "" : " (snapshotted to refs/wip/" + removed.branch()
|
||||
+ " commit=" + lateSnapshotRef + ")");
|
||||
}
|
||||
}
|
||||
if (removed != null && !preserveWorktree && removed.worktree() != null) {
|
||||
// fleetd #283: this is the one cleanup step in this method that used to be bare. By the
|
||||
// time it runs, the registry entry, the retained handle, and the pane are all already
|
||||
@@ -404,6 +429,24 @@ public final class SessionManager implements TurnListener {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #316: the read that actually authorises {@code worktrees.remove}, taken with the
|
||||
* worker's pane already stopped. Fails toward preserving (returns {@code true}) on any
|
||||
* exception — the same rule the pre-stop check applies (CB-581): once we can no longer tell
|
||||
* whether the worktree is dirty, preserving costs disk while deleting on a guess can destroy
|
||||
* work that has no other copy.
|
||||
*/
|
||||
private boolean dirtyImmediatelyBeforeRemoval(MemberSession removed) {
|
||||
try {
|
||||
return worktrees.hasUncommitted(removed.worktree());
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("release could not re-check worktree {} for pane={} terminal={} immediately "
|
||||
+ "before removal; preserving it rather than risk destroying unsaved work: {}",
|
||||
removed.worktree(), removed.paneId(), removed.terminalId(), e.toString());
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Best-effort snapshot of a dirty worktree into {@code refs/wip/<branch>} (CB-578 stage C). A
|
||||
* failure here must never escalate: the caller has already decided to preserve the worktree
|
||||
|
||||
@@ -153,6 +153,21 @@ class ConfigRefProfileCoverageTest {
|
||||
"ConfigRef.LAUNCH_SETTINGS_EXCLUDED names a component that does not exist on "
|
||||
+ "FleetConfig.Profile — check for a typo: " + excluded);
|
||||
|
||||
// The exclusion set is this mechanism's own escape hatch, so it has to be pinned too.
|
||||
// Found by mutation while verifying fleetd #323: moving autoCompactWindow and
|
||||
// ideOpenCommand OUT of the comparison and INTO the exclusion set left the whole suite
|
||||
// green — the loop below simply skips them, and the denominator assertion still balances.
|
||||
// That is exactly the lazy move a failing coverage test invites, and it silently restores
|
||||
// the #323 bug. Only ideProjectDir and worktreeGroup were saved by a behavioural test in
|
||||
// ConfigRefTest; the other two had none. So: growing this set now requires editing this
|
||||
// line as well, which is a visible, deliberate diff rather than a quiet one.
|
||||
assertEquals(Set.of("weight", "maxLoad", "credentialId"), excluded,
|
||||
"ConfigRef.LAUNCH_SETTINGS_EXCLUDED changed. A component belongs in it ONLY if it "
|
||||
+ "is read live off the config supplier, not baked into a launcher at "
|
||||
+ "startup. If you are adding one to silence this test, that is fleetd #323 "
|
||||
+ "happening again: compare it in sameLaunchSettings instead. If it really "
|
||||
+ "is read live, name where it is read and update this assertion.");
|
||||
|
||||
FleetConfig.Profile base = profileOf(BASE);
|
||||
List<String> uncovered = new java.util.ArrayList<>();
|
||||
int compared = 0;
|
||||
|
||||
@@ -506,6 +506,292 @@ class ConfigRefTest {
|
||||
assertEquals("devgroup2", ref.get().worktreeGroup());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #326: {@code primary} is read only off the startup snapshot — {@code Fleetd.java:506,
|
||||
* 519, 520} feed {@code PrimaryRegistry} and {@code ReplyPushLoop} at construction and neither is
|
||||
* rebuilt on reload — but it was missing from {@link ConfigRef#changedDeferredKeys}, so a reload
|
||||
* that only changed the pinned primary terminal reported a bare "config reloaded" while a lead
|
||||
* whose tab no longer matched stayed demoted to worker.
|
||||
*/
|
||||
@Test
|
||||
void changingPrimaryIsReportedAsDeferred(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("""
|
||||
primary:
|
||||
terminal: term-a
|
||||
"""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("""
|
||||
primary:
|
||||
terminal: term-b
|
||||
"""));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertEquals(java.util.List.of("primary"), out.deferred());
|
||||
assertTrue(out.summary().contains("need") && out.summary().contains("restart"), out.summary());
|
||||
// The snapshot still carries the new value — a restart is what makes it take effect.
|
||||
assertEquals("term-b", ref.get().primary().terminal());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #326: {@code configReload} itself is read only at startup ({@code Fleetd.java:679-680})
|
||||
* to decide whether to build a {@code ConfigWatcher} at all, and with what interval — the watcher
|
||||
* that would apply a later change is itself built once, so it is deferred rather than cold (see
|
||||
* {@link ConfigRef}'s class doc: cold means an already-open resource would go inconsistent with
|
||||
* the new value, and there is no such resource here — a running watcher just keeps polling on its
|
||||
* original enabled/interval until a restart, exactly like {@code lifecycle} or {@code guard}).
|
||||
* Before this fix, turning reload off (or changing its interval) through a reload reported a bare
|
||||
* "config reloaded" — the obvious joke the issue names.
|
||||
*/
|
||||
@Test
|
||||
void changingConfigReloadIsReportedAsDeferred(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("""
|
||||
configReload:
|
||||
enabled: true
|
||||
intervalSeconds: 10
|
||||
"""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("""
|
||||
configReload:
|
||||
enabled: false
|
||||
intervalSeconds: 30
|
||||
"""));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertEquals(java.util.List.of("configReload"), out.deferred());
|
||||
assertTrue(out.summary().contains("need") && out.summary().contains("restart"), out.summary());
|
||||
// The snapshot still carries the new value — a restart is what makes it take effect.
|
||||
assertFalse(ref.get().configReload().isEnabled());
|
||||
assertEquals(30, ref.get().configReload().intervalSeconds());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #330: {@code health:} is read both ways — {@code Fleetd.java:556-563} builds the
|
||||
* monitor off the startup snapshot and never rebuilds it, but {@code Fleetd.java:648-650} reads
|
||||
* {@code config.get().health()} live on every {@code fleet_profiles} call. A changed value is
|
||||
* neither purely hot nor purely deferred, so it gets its own {@code split} report naming both
|
||||
* halves rather than a bare "config reloaded" (which would hide the frozen half) or a plain
|
||||
* {@code deferred} entry (which would hide that the coverage string already applied).
|
||||
*/
|
||||
@Test
|
||||
void changingHealthIsReportedAsSplit(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("""
|
||||
health:
|
||||
enabled: true
|
||||
intervalSeconds: 30
|
||||
"""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("""
|
||||
health:
|
||||
enabled: true
|
||||
intervalSeconds: 90
|
||||
"""));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
|
||||
assertEquals(1, out.split().size(), out.split().toString());
|
||||
assertTrue(out.split().getFirst().startsWith("health:"), out.split().toString());
|
||||
assertTrue(out.split().getFirst().contains("restart"), out.split().toString());
|
||||
assertTrue(out.split().getFirst().contains("live"), out.split().toString());
|
||||
assertTrue(out.summary().contains("partially live"), out.summary());
|
||||
// The snapshot still carries the new value — the monitor itself is what waits for a restart.
|
||||
assertEquals(90, ref.get().health().intervalSeconds());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #330: {@code coordinator:} is the other split key — {@code Fleetd.java:502} opens the
|
||||
* {@code LeadMailbox} off the startup snapshot and never reopens it, but
|
||||
* {@code HerdrPeerLauncher.java:1530} reads {@code config.get().coordinator()} live on every
|
||||
* spawn to keep the broker URI env-var name out of a member's environment.
|
||||
*/
|
||||
@Test
|
||||
void changingCoordinatorIsReportedAsSplit(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("""
|
||||
coordinator:
|
||||
selfId: mac-a
|
||||
"""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("""
|
||||
coordinator:
|
||||
selfId: mac-b
|
||||
"""));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
|
||||
assertEquals(1, out.split().size(), out.split().toString());
|
||||
assertTrue(out.split().getFirst().startsWith("coordinator:"), out.split().toString());
|
||||
assertTrue(out.split().getFirst().contains("restart"), out.split().toString());
|
||||
assertTrue(out.split().getFirst().contains("live"), out.split().toString());
|
||||
// The snapshot still carries the new value — the LeadMailbox connection is what waits for a
|
||||
// restart; selfId names this daemon's own inbox queue and a peer cannot discover a rename.
|
||||
assertEquals("mac-b", ref.get().coordinator().selfId());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #333: {@code fleet:} is split too — {@code fleet.leaders} is read only at
|
||||
* {@code Fleetd.java:281} to build the {@code LeadTabScanner}'s identity map (fed into
|
||||
* {@code CallerResolver}) and to auto-launch leads via {@code LeadLauncher.ensureLeads()}, and
|
||||
* neither is rebuilt on reload, while the rest of {@code fleet:} (role pools, charters,
|
||||
* tabLabel) is read live through the {@code CompositePeerLauncher} supplier. Before this fix
|
||||
* {@code fleet} sat in the coverage checker's hot-excluded escape hatch, so a reload that
|
||||
* changed only {@code fleet.leaders} reported a bare "config reloaded" — exactly the
|
||||
* under-claim the split class exists to prevent for {@code health}/{@code coordinator}.
|
||||
*/
|
||||
@Test
|
||||
void changingFleetLeadersIsReportedAsSplit(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("""
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
tab: "lead: opus-a"
|
||||
profile: sonnet
|
||||
"""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("""
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
tab: "lead: opus-b"
|
||||
profile: sonnet
|
||||
"""));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
|
||||
assertEquals(1, out.split().size(), out.split().toString());
|
||||
assertTrue(out.split().getFirst().startsWith("fleet:"), out.split().toString());
|
||||
assertTrue(out.split().getFirst().contains("restart"), out.split().toString());
|
||||
assertTrue(out.split().getFirst().contains("live"), out.split().toString());
|
||||
assertTrue(out.summary().contains("partially live"), out.summary());
|
||||
// The snapshot still carries the new value — the LeadTabScanner's identity map and
|
||||
// LeadLauncher's auto-launch are what wait for a restart; a reload rebuilds neither.
|
||||
assertEquals("lead: opus-b", ref.get().fleet().leaders().get("opus").tab());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #333: {@code fleet.leaders} is the ONLY frozen part of {@code fleet:}. A reload that
|
||||
* changes {@code tabLabel} (or charters, or a role pool) without touching {@code fleet.leaders}
|
||||
* must stay fully hot with nothing reported — proving {@link ConfigRef#changedSplitKeys}
|
||||
* compares {@code fleet.leaders} specifically rather than the whole {@code Fleet} record, which
|
||||
* would over-claim "needs a restart" for a change that is genuinely all live (the mirror
|
||||
* mistake of the under-claim this class exists to prevent).
|
||||
*/
|
||||
@Test
|
||||
void changingFleetTabLabelWithoutLeadersStaysFullyHot(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("""
|
||||
fleet:
|
||||
tabLabel: "{role}: {profile} #{n}"
|
||||
"""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("""
|
||||
fleet:
|
||||
tabLabel: "[{profile}] {role}"
|
||||
"""));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertTrue(out.deferred().isEmpty(), out.deferred().toString());
|
||||
assertTrue(out.split().isEmpty(), out.split().toString());
|
||||
assertEquals("config reloaded", out.summary());
|
||||
assertEquals("[{profile}] {role}", ref.get().fleet().tabLabel());
|
||||
}
|
||||
|
||||
/**
|
||||
* A split change must not refuse the reload (invariant 2 of fleetd #330) and
|
||||
* {@code Outcome.applied()} must stay {@code true} (invariant 3) — unlike a cold change, the
|
||||
* live half of a split key genuinely took effect, so refusing would throw that away.
|
||||
*/
|
||||
@Test
|
||||
void aSplitChangeDoesNotRefuseTheReload(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("coordinator:\n selfId: mac-a\n"));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("coordinator:\n selfId: mac-b\n"));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertTrue(out.error() == null);
|
||||
assertTrue(out.coldKeys().isEmpty());
|
||||
}
|
||||
|
||||
/**
|
||||
* Both split keys can change in one reload — the report names both, and a caller reading
|
||||
* {@code split} does not have to guess which half of which key already applied.
|
||||
*/
|
||||
@Test
|
||||
void changingBothSplitKeysReportsBoth(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("""
|
||||
health:
|
||||
enabled: true
|
||||
coordinator:
|
||||
selfId: mac-a
|
||||
"""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("""
|
||||
health:
|
||||
enabled: false
|
||||
coordinator:
|
||||
selfId: mac-b
|
||||
"""));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertEquals(2, out.split().size(), out.split().toString());
|
||||
assertTrue(out.split().stream().anyMatch(s -> s.startsWith("health:")), out.split().toString());
|
||||
assertTrue(out.split().stream().anyMatch(s -> s.startsWith("coordinator:")), out.split().toString());
|
||||
}
|
||||
|
||||
/**
|
||||
* A split key and a deferred key changing in the same reload must both show up, each in its own
|
||||
* list — proving the two fields do not step on each other and {@link ConfigRef.Outcome#summary()}
|
||||
* reports both halves of the message.
|
||||
*/
|
||||
@Test
|
||||
void aSplitChangeAndADeferredChangeCoexist(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(f, yaml("""
|
||||
coordinator:
|
||||
selfId: mac-a
|
||||
lifecycle:
|
||||
drainTimeoutSeconds: 30
|
||||
"""));
|
||||
ConfigRef ref = refFor(f);
|
||||
|
||||
Files.writeString(f, yaml("""
|
||||
coordinator:
|
||||
selfId: mac-b
|
||||
lifecycle:
|
||||
drainTimeoutSeconds: 60
|
||||
"""));
|
||||
ConfigRef.Outcome out = ref.reload();
|
||||
|
||||
assertTrue(out.applied());
|
||||
assertEquals(java.util.List.of("lifecycle"), out.deferred());
|
||||
assertEquals(1, out.split().size(), out.split().toString());
|
||||
assertTrue(out.split().getFirst().startsWith("coordinator:"), out.split().toString());
|
||||
assertTrue(out.summary().contains("need a restart") || out.summary().contains("needs a restart"),
|
||||
out.summary());
|
||||
assertTrue(out.summary().contains("partially live"), out.summary());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFixedRefHasNoFileAndRefusesToReload() {
|
||||
FleetConfig cfg = new FleetConfig(null, null, null, null, null, null,
|
||||
|
||||
@@ -0,0 +1,171 @@
|
||||
package dev.ltms.fleet.config;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.RecordComponent;
|
||||
import java.util.LinkedHashSet;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #330: a new top-level {@link FleetConfig} record component must be triaged into a reload
|
||||
* class before it ships, or it repeats fleetd #323 ({@code worktreeGroup} missing from {@code
|
||||
* ConfigRef.changedDeferredKeys}) and fleetd #326 ({@code primary}/{@code configReload} missing the
|
||||
* same way) — a key silently absent from {@link ConfigRef}'s reload machinery, so a reload changing
|
||||
* only that key reports a bare "config reloaded" for a change the running daemon never picked up.
|
||||
*
|
||||
* <p>This is the top-level counterpart of {@link ConfigRefProfileCoverageTest}: instead of
|
||||
* enumerating {@link FleetConfig.Profile}'s record components, it enumerates {@link FleetConfig}'s
|
||||
* own — {@code bind}, {@code health}, {@code memberCredentials}, and so on — and requires each to
|
||||
* fall into exactly one of four homes: {@link ConfigRef#COLD_KEYS}, {@link #DEFERRED_TOP_LEVEL_KEYS}
|
||||
* (compared in {@code ConfigRef.changedDeferredKeys}), {@link ConfigRef#SPLIT_KEYS}, or
|
||||
* {@link #HOT_EXCLUDED_TOP_LEVEL_KEYS} (the escape hatch: read live off the config supplier, so no
|
||||
* reload bookkeeping is needed for it at all).
|
||||
*
|
||||
* <h2>What this checker can and cannot prove</h2>
|
||||
* It proves the record's <em>shape</em> is fully triaged: every one of {@code FleetConfig}'s
|
||||
* components sits in exactly one of the four sets, none sits in two, and the escape hatch
|
||||
* ({@link #HOT_EXCLUDED_TOP_LEVEL_KEYS}) cannot silently grow without a visible diff to this file.
|
||||
* That is what "a new component cannot be added without someone triaging it" means in practice.
|
||||
*
|
||||
* <p>It CANNOT prove that any of the citations are <em>true</em>. "Compared in {@code
|
||||
* changedDeferredKeys}" and "read live off {@code config.get()}" are facts about {@code
|
||||
* ConfigRef.java}, {@code Fleetd.java} and {@code HerdrPeerLauncher.java} that a reflection-only
|
||||
* test over {@code FleetConfig}'s shape has no way to inspect — this test would pass identically
|
||||
* whether or not the cited line still does what the comment next to it says. Trust the citation
|
||||
* because a person read the source (the exact call sites are named next to each set below), not
|
||||
* because this test is green. {@link ConfigRefTest} is what behaviourally proves the deferred and
|
||||
* split keys it covers actually get reported; {@link ConfigRefProfileCoverageTest} does the same,
|
||||
* behaviourally, for {@code FleetConfig.Profile}'s own fields.
|
||||
*/
|
||||
class ConfigRefTopLevelCoverageTest {
|
||||
|
||||
private static final RecordComponent[] COMPONENTS = FleetConfig.class.getRecordComponents();
|
||||
|
||||
/**
|
||||
* Top-level components whose change {@code ConfigRef.changedDeferredKeys} reads and reports on.
|
||||
* Fleetd #337 promoted this out of a hand-maintained copy here into {@link ConfigRef#DEFERRED_KEYS}
|
||||
* itself, the same reason {@link ConfigRef#COLD_KEYS} and {@link ConfigRef#SPLIT_KEYS} are
|
||||
* production constants rather than test-side copies: two lists that are supposed to describe the
|
||||
* same method are exactly the shape that silently drifts apart. {@code spawnReadyTimeoutMs}/
|
||||
* {@code spawnReadyPollMs} are compared together and reported under one combined label
|
||||
* ({@code "spawnReady*"}); {@code profiles} is compared twice over — once for added/removed
|
||||
* profile names, once for an existing profile's launch settings — and that second comparison
|
||||
* excludes {@code weight}/{@code maxLoad}/{@code credentialId} as hot sub-fields, which is what
|
||||
* {@link ConfigRefProfileCoverageTest} exists to keep honest at the sub-field level. {@code
|
||||
* profiles} itself still belongs here, not in the hot-exclusion set below: most of a profile's
|
||||
* fields are NOT read live, so citing "read live off the config supplier" for the whole
|
||||
* top-level key would be false.
|
||||
*/
|
||||
private static final Set<String> DEFERRED_TOP_LEVEL_KEYS = ConfigRef.DEFERRED_KEYS;
|
||||
|
||||
/**
|
||||
* The escape hatch: top-level components with no reload bookkeeping at all, because every read
|
||||
* of them goes live through {@link ConfigRef#get()} rather than off a startup snapshot. A
|
||||
* component belongs here ONLY if that is true — never because adding it here makes this test
|
||||
* pass. fleetd #323 is the cautionary tale for exactly this pattern: the identical hatch on
|
||||
* {@code ConfigRef.LAUNCH_SETTINGS_EXCLUDED} let two profile fields be silently re-broken with
|
||||
* the whole suite green, and it was only caught by mutating the checker itself (see this class's
|
||||
* own mutation test below, and {@code ConfigRefProfileCoverageTest}'s equivalent).
|
||||
*
|
||||
* <ul>
|
||||
* <li>{@code placement} — read live by the placement policy on every spawn (class doc, Hot
|
||||
* bullet; {@code ConfigRefTest.aConsumerHoldingTheRefSeesTheNewValue} proves it
|
||||
* behaviourally).</li>
|
||||
* <li>{@code memberCredentials} — read live at {@code Fleetd.java:198, 205, 729}.</li>
|
||||
* <li>{@code memberLoginShell} — read live at
|
||||
* {@code HerdrPeerLauncher#configuredMemberLoginShell}.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>{@code fleet} used to sit here too, on the strength of most of it (role pools, charters,
|
||||
* tabLabel) being read the same live way — but {@code fleet.leaders} inside the same key is
|
||||
* read only at startup and genuinely needs a restart, which is a real gap this escape hatch
|
||||
* cannot represent: it excuses a whole top-level key from reload bookkeeping, and {@code fleet}
|
||||
* needed exactly half of it excused. fleetd #333 moved it into {@link ConfigRef#SPLIT_KEYS}
|
||||
* instead, where {@code changedSplitKeys} reports the frozen half by name. That is the
|
||||
* cautionary tale this test's own javadoc already told: this checker proves the record's shape
|
||||
* is triaged, never that a bucket a key sits in is the right one — a person has to read the
|
||||
* source, which is exactly how fleetd #333 found {@code fleet} sitting in the wrong bucket
|
||||
* while this test stayed green throughout.</p>
|
||||
*/
|
||||
private static final Set<String> HOT_EXCLUDED_TOP_LEVEL_KEYS =
|
||||
Set.of("placement", "memberCredentials", "memberLoginShell");
|
||||
|
||||
@Test
|
||||
void everyTopLevelComponentIsAccountedForInExactlyOneClass() {
|
||||
Set<String> allNames = new TreeSet<>();
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
allNames.add(rc.getName());
|
||||
}
|
||||
|
||||
Set<String> cold = ConfigRef.COLD_KEYS;
|
||||
Set<String> split = ConfigRef.SPLIT_KEYS;
|
||||
Set<String> deferred = DEFERRED_TOP_LEVEL_KEYS;
|
||||
Set<String> hot = HOT_EXCLUDED_TOP_LEVEL_KEYS;
|
||||
|
||||
// Typo guard on each set — the same check ConfigRefProfileCoverageTest runs on
|
||||
// LAUNCH_SETTINGS_EXCLUDED. A name that does not exist on FleetConfig is a silent no-op.
|
||||
assertTrue(allNames.containsAll(cold),
|
||||
"ConfigRef.COLD_KEYS names a component that does not exist on FleetConfig: " + cold);
|
||||
assertTrue(allNames.containsAll(split),
|
||||
"ConfigRef.SPLIT_KEYS names a component that does not exist on FleetConfig: " + split);
|
||||
assertTrue(allNames.containsAll(deferred),
|
||||
"DEFERRED_TOP_LEVEL_KEYS names a component that does not exist on FleetConfig: " + deferred);
|
||||
assertTrue(allNames.containsAll(hot),
|
||||
"HOT_EXCLUDED_TOP_LEVEL_KEYS names a component that does not exist on FleetConfig: " + hot);
|
||||
|
||||
// The escape hatch is pinned. Growing it requires editing this line — a visible, deliberate
|
||||
// diff, not a quiet one. See the field javadoc above for what "belongs here" actually means.
|
||||
assertEquals(Set.of("placement", "memberCredentials", "memberLoginShell"), hot,
|
||||
"HOT_EXCLUDED_TOP_LEVEL_KEYS changed. A component belongs here ONLY if it is read "
|
||||
+ "live off the config supplier, never because adding it makes this test "
|
||||
+ "pass. If you are adding one to silence this test, that is fleetd #323 "
|
||||
+ "happening again: account for it in ConfigRef.changedDeferredKeys (or "
|
||||
+ "COLD_KEYS/SPLIT_KEYS) instead. If it really is read live, name where and "
|
||||
+ "update this assertion and the field javadoc together.");
|
||||
|
||||
// No component may sit in two buckets at once — the denominator check below could not catch
|
||||
// that on its own (two buckets double-booking one key still sums to the right total if
|
||||
// another key is simultaneously missing), so check every pair directly and name the culprit.
|
||||
record Bucket(String name, Set<String> keys) {}
|
||||
List<Bucket> buckets = List.of(
|
||||
new Bucket("COLD_KEYS", cold), new Bucket("SPLIT_KEYS", split),
|
||||
new Bucket("DEFERRED_TOP_LEVEL_KEYS", deferred), new Bucket("HOT_EXCLUDED_TOP_LEVEL_KEYS", hot));
|
||||
for (int i = 0; i < buckets.size(); i++) {
|
||||
for (int j = i + 1; j < buckets.size(); j++) {
|
||||
Set<String> overlap = new LinkedHashSet<>(buckets.get(i).keys());
|
||||
overlap.retainAll(buckets.get(j).keys());
|
||||
assertEquals(Set.of(), overlap, "a component is in both " + buckets.get(i).name()
|
||||
+ " and " + buckets.get(j).name() + ": " + overlap);
|
||||
}
|
||||
}
|
||||
|
||||
Set<String> union = new TreeSet<>();
|
||||
union.addAll(cold);
|
||||
union.addAll(split);
|
||||
union.addAll(deferred);
|
||||
union.addAll(hot);
|
||||
|
||||
System.out.printf(
|
||||
"FleetConfig top-level coverage — %d components total: %d cold %s, %d deferred %s, "
|
||||
+ "%d split %s, %d hot-excluded %s%n",
|
||||
allNames.size(), cold.size(), cold, deferred.size(), deferred, split.size(), split,
|
||||
hot.size(), hot);
|
||||
|
||||
Set<String> missing = new TreeSet<>(allNames);
|
||||
missing.removeAll(union);
|
||||
assertEquals(Set.of(), missing,
|
||||
"these FleetConfig components are in none of COLD_KEYS, DEFERRED_TOP_LEVEL_KEYS, "
|
||||
+ "SPLIT_KEYS or HOT_EXCLUDED_TOP_LEVEL_KEYS — triage each one into whichever "
|
||||
+ "actually describes it: " + missing);
|
||||
assertEquals(allNames.size(), cold.size() + deferred.size() + split.size() + hot.size(),
|
||||
"counts don't sum to the component total even though every component was found in "
|
||||
+ "the union — " + allNames.size() + " components, " + cold.size()
|
||||
+ " cold + " + deferred.size() + " deferred + " + split.size() + " split + "
|
||||
+ hot.size() + " hot-excluded");
|
||||
}
|
||||
}
|
||||
+284
@@ -0,0 +1,284 @@
|
||||
package dev.ltms.fleet.config;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.lang.reflect.Constructor;
|
||||
import java.lang.reflect.RecordComponent;
|
||||
import java.util.ArrayList;
|
||||
import java.util.Arrays;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.TreeSet;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
/**
|
||||
* fleetd #333, finding F2: {@link ConfigRefTopLevelCoverageTest} proves every {@link FleetConfig}
|
||||
* top-level component sits in exactly one of {@link ConfigRef#COLD_KEYS}, {@link
|
||||
* ConfigRef#DEFERRED_KEYS}, {@link ConfigRef#SPLIT_KEYS} or the hot-excluded set. It does
|
||||
* <strong>not</strong> prove that a key's membership in one of the first three sets corresponds to
|
||||
* any actual comparison in {@link ConfigRef}: a key can sit in a set with no branch in {@code
|
||||
* changedColdKeys}/{@code changedSplitKeys}/{@code changedDeferredKeys} checking it, and both
|
||||
* {@link ConfigRefTopLevelCoverageTest} and the "kept in step" {@code assert} inside the first two
|
||||
* of those methods stay green, because neither one reads the method body — the coverage test only
|
||||
* reads set membership, and the assert only checks that reported entries are a SUBSET of the set,
|
||||
* never that every set member produced a reported entry. {@code changedDeferredKeys} does not even
|
||||
* have a "kept in step" assert of its own.
|
||||
*
|
||||
* <p>Measured directly, live, while fixing fleetd #333: dropping the {@code coordinator} branch out
|
||||
* of {@code ConfigRef.changedSplitKeys} while leaving {@code "coordinator"} in
|
||||
* {@link ConfigRef#SPLIT_KEYS} left {@link ConfigRefTopLevelCoverageTest} and the in-method assert
|
||||
* both green — only a hand-written behavioural case in {@link ConfigRefTest} caught it, because it
|
||||
* happened to name that exact key. This is the {@link ConfigRefProfileCoverageTest} mechanism one
|
||||
* level up, generalised over every {@code COLD_KEYS}/{@code SPLIT_KEYS}/{@code DEFERRED_KEYS}
|
||||
* member rather than one hand-picked field: enumerate {@link FleetConfig}'s own record components by
|
||||
* reflection, build "a config where only {@code <key>} differs" for each key, and call the real
|
||||
* {@link ConfigRef#changedColdKeys}/{@link ConfigRef#changedSplitKeys}/
|
||||
* {@link ConfigRef#changedDeferredKeys} methods (all package-private for exactly this, the same
|
||||
* reason {@link ConfigRef#sameLaunchSettings} already is) to prove each one is actually reported —
|
||||
* not assumed from a set literal.
|
||||
*
|
||||
* <h2>fleetd #337 — DEFERRED_KEYS was the gap left open here</h2>
|
||||
* This class originally covered {@code COLD_KEYS} and {@code SPLIT_KEYS} only — {@code
|
||||
* DEFERRED_TOP_LEVEL_KEYS} (now {@link ConfigRef#DEFERRED_KEYS}) carried the identical one-way risk
|
||||
* in principle, unexercised. fleetd #337 measured the real consequence rather than assuming it from
|
||||
* the shape of the gap: dropping {@code guard}'s comparison out of {@code changedDeferredKeys} while
|
||||
* {@code "guard"} stayed in the set left all 1355 tests green — the same failure mode {@code
|
||||
* coordinator} demonstrated for {@code SPLIT_KEYS} in fleetd #333, now confirmed for {@code
|
||||
* DEFERRED_KEYS} too. Re-deriving the full list by mutation (drop each key's branch in turn, run the
|
||||
* suite, restore) found six of the eleven {@code DEFERRED_KEYS} members with no behavioural test in
|
||||
* {@link ConfigRefTest} naming them: {@code guard}, {@code leadHeartbeat}, {@code worktreeRoot},
|
||||
* {@code spawnReadyTimeoutMs}, {@code spawnReadyPollMs} and {@code quarantineCooldownSeconds}. That
|
||||
* list corrects fleetd #333's own guess at it in two ways the mutation proved and a reading did not:
|
||||
* {@code lifecycle} is NOT on it — {@code ConfigRefTest.aDeferredChangeIsAppliedAndReported} already
|
||||
* names it, and dropping its branch fails that test — and {@code worktreeRoot} IS on it, which #333
|
||||
* never named at all. The other five {@code DEFERRED_KEYS} members ({@code lifecycle}, {@code
|
||||
* worktreeGroup}, {@code primary}, {@code configReload}, {@code profiles}) already had a hand-written
|
||||
* case each. {@link #everyDeferredKeyIsActuallyReportedByChangedDeferredKeys} below now covers all
|
||||
* eleven the reflective way, so the six with no hand-written test are no longer silently unpinned —
|
||||
* every {@code DEFERRED_KEYS} component turned out to be a scalar or a simple record, so, unlike
|
||||
* {@code fleet.leaders} in fleetd #333, none needed an exclusion set: {@link #BASE}/{@link #ALT} give
|
||||
* every top-level component (not only {@code COLD_KEYS}/{@code SPLIT_KEYS}) a real, distinct value.
|
||||
*/
|
||||
class ConfigRefTopLevelReportingCoverageTest {
|
||||
|
||||
private static final RecordComponent[] COMPONENTS = FleetConfig.class.getRecordComponents();
|
||||
|
||||
/**
|
||||
* One valid value per top-level {@link FleetConfig} component — "the a value". fleetd #337 gave
|
||||
* every {@code DEFERRED_KEYS} component a real value here too (previously left {@code null} on
|
||||
* both sides, which meant {@code mutate(key)} produced no actual difference for any of them);
|
||||
* only {@code placement}, {@code memberCredentials} and {@code memberLoginShell} — the
|
||||
* hot-excluded set, never compared by any {@code changed*Keys} method — stay {@code null}.
|
||||
* {@link FleetConfig}'s compact constructor only normalizes {@code profiles}, so every other
|
||||
* field accepts whatever is put here unmutated.
|
||||
*/
|
||||
private static final Map<String, Object> BASE = baseValues();
|
||||
|
||||
/** The same shape, each value distinct from {@link #BASE} — "the b value". */
|
||||
private static final Map<String, Object> ALT = altValues();
|
||||
|
||||
private static Map<String, Object> baseValues() {
|
||||
Map<String, Object> v = new LinkedHashMap<>();
|
||||
v.put("bind", new FleetConfig.Bind("127.0.0.1", 8765));
|
||||
v.put("herdrSocket", "~/.config/herdr/a.sock");
|
||||
v.put("memberHerdrSocket", "~/.config/herdr/member-a.sock");
|
||||
v.put("profiles", Map.of());
|
||||
v.put("guard", new FleetConfig.Guard(List.of("host-a")));
|
||||
v.put("worktreeRoot", "/wt/a");
|
||||
v.put("lifecycle", new FleetConfig.Lifecycle(300, 5, 30, false));
|
||||
v.put("spawnReadyTimeoutMs", 5000);
|
||||
v.put("spawnReadyPollMs", 100);
|
||||
v.put("broker", new FleetConfig.Broker("amqp://a", null, 1));
|
||||
v.put("primary", new FleetConfig.Primary("term-a", 1, 1000));
|
||||
v.put("fleet", new FleetConfig.Fleet(
|
||||
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-a", 1, null, 10,
|
||||
"claude", null, null, null)),
|
||||
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
|
||||
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(300, 60_000L, 3));
|
||||
v.put("health", new FleetConfig.Health(true, 30, 600, null, null));
|
||||
v.put("placement", null);
|
||||
v.put("auth", new FleetConfig.Auth("loopback-trust", null));
|
||||
v.put("configReload", new FleetConfig.ConfigReload(true, 10));
|
||||
v.put("quarantineCooldownSeconds", 1800);
|
||||
v.put("memberCredentials", null);
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-a", null, "self-a", 1));
|
||||
v.put("worktreeGroup", "group-a");
|
||||
v.put("memberLoginShell", null);
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
private static Map<String, Object> altValues() {
|
||||
Map<String, Object> v = new LinkedHashMap<>();
|
||||
v.put("bind", new FleetConfig.Bind("127.0.0.2", 8766));
|
||||
v.put("herdrSocket", "~/.config/herdr/b.sock");
|
||||
v.put("memberHerdrSocket", "~/.config/herdr/member-b.sock");
|
||||
// A single added profile — enough to trip the "added/removed" comparison in
|
||||
// ConfigRef.changedDeferredKeys, which is all this mechanism needs to prove "profiles" has
|
||||
// a branch behind it; the launch-settings comparison already has its own hand-written cases
|
||||
// in ConfigRefTest (changingAProfilesLaunchSettingsIsReportedAsDeferred and siblings).
|
||||
v.put("profiles", Map.of("sonnet", minimalProfile("sonnet")));
|
||||
v.put("guard", new FleetConfig.Guard(List.of("host-b")));
|
||||
v.put("worktreeRoot", "/wt/b");
|
||||
v.put("lifecycle", new FleetConfig.Lifecycle(600, 10, 60, true));
|
||||
v.put("spawnReadyTimeoutMs", 10_000);
|
||||
v.put("spawnReadyPollMs", 200);
|
||||
v.put("broker", new FleetConfig.Broker("amqp://b", null, 2));
|
||||
v.put("primary", new FleetConfig.Primary("term-b", 2, 2000));
|
||||
// Differs from BASE.fleet only in fleet.leaders (a different tab for "opus") — the frozen
|
||||
// sub-field ConfigRef.changedSplitKeys actually compares. A Fleet that instead differed only
|
||||
// in tabLabel would correctly NOT be reported (see
|
||||
// ConfigRefTest.changingFleetTabLabelWithoutLeadersStaysFullyHot) and would wrongly fail this
|
||||
// test — that is by design, not a gap: this map exists to prove fleet.leaders is covered.
|
||||
v.put("fleet", new FleetConfig.Fleet(
|
||||
Map.of("opus", new FleetConfig.Leader("sonnet", "lead: opus-b", 1, null, 10,
|
||||
"claude", null, null, null)),
|
||||
Map.of(), Map.of(), Map.of(), Map.of(), "{role}: {profile} #{n}"));
|
||||
v.put("leadHeartbeat", new FleetConfig.LeadHeartbeat(600, 120_000L, 5));
|
||||
v.put("health", new FleetConfig.Health(false, 90, 900, null, null));
|
||||
v.put("placement", null);
|
||||
v.put("auth", new FleetConfig.Auth("token", "TOKEN_ENV"));
|
||||
v.put("configReload", new FleetConfig.ConfigReload(false, 20));
|
||||
v.put("quarantineCooldownSeconds", 3600);
|
||||
v.put("memberCredentials", null);
|
||||
v.put("coordinator", new FleetConfig.Coordinator("amqp://coord-b", null, "self-b", 2));
|
||||
v.put("worktreeGroup", "group-b");
|
||||
v.put("memberLoginShell", null);
|
||||
assertNamesMatchComponents(v);
|
||||
return v;
|
||||
}
|
||||
|
||||
/** A minimal, otherwise-null {@link FleetConfig.Profile} — just enough to name one in a map. */
|
||||
private static FleetConfig.Profile minimalProfile(String name) {
|
||||
return new FleetConfig.Profile(name, null, null, null, null, null, null, null, null, null,
|
||||
null, null, null, null, null, null, null, null, null, null, null, null, null, null,
|
||||
null, null);
|
||||
}
|
||||
|
||||
private static void assertNamesMatchComponents(Map<String, Object> values) {
|
||||
Set<String> names = new TreeSet<>();
|
||||
for (RecordComponent rc : COMPONENTS) {
|
||||
names.add(rc.getName());
|
||||
}
|
||||
assertEquals(names, new TreeSet<>(values.keySet()),
|
||||
"this test's value map has drifted from FleetConfig's actual top-level components — "
|
||||
+ "update BASE/ALT alongside the record");
|
||||
}
|
||||
|
||||
private static FleetConfig configOf(Map<String, Object> values) throws ReflectiveOperationException {
|
||||
Class<?>[] types = Arrays.stream(COMPONENTS).map(RecordComponent::getType).toArray(Class<?>[]::new);
|
||||
Object[] args = Arrays.stream(COMPONENTS).map(rc -> values.get(rc.getName())).toArray();
|
||||
Constructor<FleetConfig> ctor = FleetConfig.class.getDeclaredConstructor(types);
|
||||
return ctor.newInstance(args);
|
||||
}
|
||||
|
||||
/** {@code BASE} with exactly one named top-level component swapped for its {@code ALT} value. */
|
||||
private static FleetConfig mutate(String componentName) throws ReflectiveOperationException {
|
||||
Map<String, Object> values = new LinkedHashMap<>(BASE);
|
||||
values.put(componentName, ALT.get(componentName));
|
||||
return configOf(values);
|
||||
}
|
||||
|
||||
/**
|
||||
* The mechanism fleetd #333 F2 asked for, applied to {@link ConfigRef#COLD_KEYS}: mutate each
|
||||
* cold key in isolation and prove {@code changedColdKeys} actually names it, not just that
|
||||
* {@code COLD_KEYS} claims it does.
|
||||
*/
|
||||
@Test
|
||||
void everyColdKeyIsActuallyReportedByChangedColdKeys() throws ReflectiveOperationException {
|
||||
FleetConfig base = configOf(BASE);
|
||||
List<String> uncovered = new ArrayList<>();
|
||||
for (String key : new TreeSet<>(ConfigRef.COLD_KEYS)) {
|
||||
FleetConfig mutated = mutate(key);
|
||||
if (!ConfigRef.changedColdKeys(base, mutated).contains(key)) {
|
||||
uncovered.add(key);
|
||||
}
|
||||
}
|
||||
System.out.printf(
|
||||
"ConfigRef.changedColdKeys reporting coverage — %d COLD_KEYS, %d verified%n",
|
||||
ConfigRef.COLD_KEYS.size(), ConfigRef.COLD_KEYS.size() - uncovered.size());
|
||||
assertEquals(List.of(), uncovered,
|
||||
"these keys are in ConfigRef.COLD_KEYS but mutating them alone produces no matching "
|
||||
+ "entry from changedColdKeys — a set entry with no comparison behind it: "
|
||||
+ uncovered);
|
||||
}
|
||||
|
||||
/**
|
||||
* The same mechanism applied to {@link ConfigRef#SPLIT_KEYS}: mutate each split key in
|
||||
* isolation and prove {@code changedSplitKeys} actually names it (message starting with
|
||||
* {@code "<key>:"}), not just that {@code SPLIT_KEYS} claims it does. This is the exact check
|
||||
* that would have failed fleetd #333's own reproduction — dropping the {@code coordinator}
|
||||
* branch from {@code changedSplitKeys} while {@code "coordinator"} stayed in {@code SPLIT_KEYS}
|
||||
* — see the mutation proof in the fleetd #333 PR description; the earlier checkers here (the
|
||||
* shape test and the in-method assert) do not.
|
||||
*/
|
||||
@Test
|
||||
void everySplitKeyIsActuallyReportedByChangedSplitKeys() throws ReflectiveOperationException {
|
||||
FleetConfig base = configOf(BASE);
|
||||
List<String> uncovered = new ArrayList<>();
|
||||
for (String key : new TreeSet<>(ConfigRef.SPLIT_KEYS)) {
|
||||
FleetConfig mutated = mutate(key);
|
||||
List<String> split = ConfigRef.changedSplitKeys(base, mutated);
|
||||
if (split.stream().noneMatch(s -> s.startsWith(key + ":"))) {
|
||||
uncovered.add(key);
|
||||
}
|
||||
}
|
||||
System.out.printf(
|
||||
"ConfigRef.changedSplitKeys reporting coverage — %d SPLIT_KEYS, %d verified%n",
|
||||
ConfigRef.SPLIT_KEYS.size(), ConfigRef.SPLIT_KEYS.size() - uncovered.size());
|
||||
assertEquals(List.of(), uncovered,
|
||||
"these keys are in ConfigRef.SPLIT_KEYS but mutating them alone produces no matching "
|
||||
+ "entry from changedSplitKeys — a set entry with no comparison behind it, "
|
||||
+ "exactly the fleetd #333 F2 shape: " + uncovered);
|
||||
}
|
||||
|
||||
/**
|
||||
* For most {@link ConfigRef#DEFERRED_KEYS} members, {@code changedDeferredKeys} reports the key
|
||||
* name verbatim — the default this map assumes. Two entries don't: {@code spawnReadyTimeoutMs}
|
||||
* and {@code spawnReadyPollMs} are compared together in one branch and reported under the
|
||||
* combined label {@code "spawnReady*"} (see {@link ConfigRef#changedDeferredKeys}). {@code
|
||||
* profiles} keeps the default: mutating it here only exercises the added/removed comparison
|
||||
* (see {@link #altValues}), which reports {@code "profiles (added/removed: …)"} — starts with
|
||||
* {@code "profiles"}, same as the default would expect.
|
||||
*/
|
||||
private static final Map<String, String> DEFERRED_REPORT_PREFIX = Map.of(
|
||||
"spawnReadyTimeoutMs", "spawnReady*",
|
||||
"spawnReadyPollMs", "spawnReady*");
|
||||
|
||||
/**
|
||||
* fleetd #337: the same mechanism applied to {@link ConfigRef#DEFERRED_KEYS}, closing the gap
|
||||
* this class's own javadoc left open since fleetd #333. Mutate each deferred key in isolation
|
||||
* and prove {@code changedDeferredKeys} actually names it (message starting with the key's
|
||||
* expected report prefix — see {@link #DEFERRED_REPORT_PREFIX}), not just that {@code
|
||||
* DEFERRED_KEYS} claims it does. This is the exact check that fails for {@code guard} the way
|
||||
* {@code coordinator} failed {@link #everySplitKeyIsActuallyReportedByChangedSplitKeys} in
|
||||
* fleetd #333 — verified live: dropping {@code guard}'s branch from {@code changedDeferredKeys}
|
||||
* while {@code "guard"} stayed in {@code DEFERRED_KEYS} left the whole 1355-test suite green,
|
||||
* and this test is what now catches it (it fails naming {@code guard} with that mutation in
|
||||
* place).
|
||||
*/
|
||||
@Test
|
||||
void everyDeferredKeyIsActuallyReportedByChangedDeferredKeys() throws ReflectiveOperationException {
|
||||
FleetConfig base = configOf(BASE);
|
||||
List<String> uncovered = new ArrayList<>();
|
||||
for (String key : new TreeSet<>(ConfigRef.DEFERRED_KEYS)) {
|
||||
FleetConfig mutated = mutate(key);
|
||||
List<String> deferred = ConfigRef.changedDeferredKeys(base, mutated);
|
||||
String prefix = DEFERRED_REPORT_PREFIX.getOrDefault(key, key);
|
||||
if (deferred.stream().noneMatch(s -> s.startsWith(prefix))) {
|
||||
uncovered.add(key);
|
||||
}
|
||||
}
|
||||
System.out.printf(
|
||||
"ConfigRef.changedDeferredKeys reporting coverage — %d DEFERRED_KEYS, %d verified%n",
|
||||
ConfigRef.DEFERRED_KEYS.size(), ConfigRef.DEFERRED_KEYS.size() - uncovered.size());
|
||||
assertEquals(List.of(), uncovered,
|
||||
"these keys are in ConfigRef.DEFERRED_KEYS but mutating them alone produces no "
|
||||
+ "matching entry from changedDeferredKeys — a set entry with no comparison "
|
||||
+ "behind it, exactly the fleetd #333 F2 shape, confirmed here for "
|
||||
+ "DEFERRED_KEYS by fleetd #337: " + uncovered);
|
||||
}
|
||||
}
|
||||
@@ -804,7 +804,28 @@ class CompletionResolverTest {
|
||||
// --- fleetd#201 Unit 1: target-keyed backend-error pattern + typed sink ----------------------
|
||||
|
||||
@Test
|
||||
void aConfiguredBackendErrorPatternClassifiesAMatchAsAFailureAndNotifiesTheSinkOnce() {
|
||||
void aNormalMemberReportMentioningTheFallbackErrorPatternFailsButDoesNotNotifyTheSink() {
|
||||
String block = "⏺ I checked the retry path. An API Error: makes it back off.\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
java.util.List<String> notified = new java.util.ArrayList<>();
|
||||
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), BackendErrorPatternLookup.legacy(), sink);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"a scrape mentioning the pattern must still fail the send");
|
||||
assertTrue(waiter.getNow(null).text().contains("I checked the retry path. An API Error: makes it back off."),
|
||||
"the failure must keep the whole pane tail");
|
||||
assertTrue(notified.isEmpty(),
|
||||
"a normal report mentioning the fallback pattern must not record a credential failure");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aConfiguredBackendErrorPatternAtTheStartOfALineClassifiesAMatchAndNotifiesTheSinkOnce() {
|
||||
String block = "⏺ 503 Service Unavailable: upstream credential rejected\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
@@ -924,7 +945,7 @@ class CompletionResolverTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void aConfiguredPatternAlsoClassifiesTheRawScrapeFallbackAndNotifiesTheSink() {
|
||||
void aConfiguredPatternAtTheStartOfALineAlsoClassifiesTheRawScrapeFallbackAndNotifiesTheSink() {
|
||||
// No ⏺ marker and leading TUI chrome ⇒ lastAssistantBlock() yields "", so classification must
|
||||
// fall back to the raw scrape (fleetd#211) — and it must use the configured pattern too.
|
||||
String block = """
|
||||
@@ -949,6 +970,40 @@ class CompletionResolverTest {
|
||||
assertTrue(notified.get(0).contains("503 Service Unavailable"), notified.get(0));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #339 follow-up: a genuine backend error rendered behind terminal chrome must still
|
||||
* record the credential outage. #339 added a start-of-line check to stop a member's own prose
|
||||
* being counted as an outage, and a bare {@code lookingAt} also rejected this — the send failed
|
||||
* but the sink never fired. #339's invariant 3 named that direction as the worse one: a real
|
||||
* outage going unrecorded leaves the fleet spawning into a dead credential.
|
||||
*
|
||||
* <p>The raw-scrape path is where this matters, because its own comment says to expect leading
|
||||
* TUI chrome there.
|
||||
*/
|
||||
@Test
|
||||
void aRealErrorBehindTerminalChromeStillNotifiesTheSink() {
|
||||
String block = """
|
||||
╭──────────────────────────────────────╮
|
||||
│ 503 Service Unavailable: upstream credential rejected
|
||||
""";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable");
|
||||
java.util.List<String> notified = new java.util.ArrayList<>();
|
||||
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), patterns, sink);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"a real backend error must still fail the send");
|
||||
assertEquals(1, notified.size(),
|
||||
"a real error line behind box chrome is still a real outage — it must reach the sink, "
|
||||
+ "or the fleet keeps spawning into a dead credential");
|
||||
}
|
||||
|
||||
// --- fleetd#201 Unit 1: classification inside the fleetd#164 MIN_TURN_NANOS floor -------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -48,7 +48,7 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void deliversWhenIdle() {
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "hello", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "hello", TestTurnTokens.inert(T)).completion();
|
||||
assertFalse(f.isDone(), "not delivered until an injectable status arrives");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
assertTrue(f.isDone());
|
||||
@@ -170,6 +170,32 @@ class InjectorTest {
|
||||
assertEquals(List.of("a", "b", "c"), sent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void cancellingTheMiddleDeliveryKeepsTheFollowingDeliveryReachable() {
|
||||
Injector.Delivery first = injector.enqueue(T, "same text", TestTurnTokens.inert(T));
|
||||
Injector.Delivery cancelled = injector.enqueue(T, "same text", TestTurnTokens.inert(T));
|
||||
injector.enqueue(T, "after cancelled", TestTurnTokens.inert(T));
|
||||
|
||||
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(cancelled));
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertEquals(List.of("same text", "after cancelled"), sent(),
|
||||
"cancellation must match the exact Delivery and preserve the remaining FIFO queue");
|
||||
assertTrue(first.completion().isDone());
|
||||
}
|
||||
|
||||
@Test
|
||||
void cancellationReportsDeliveredWhenPickupWonTheRace() {
|
||||
Injector.Delivery delivery = injector.enqueue(T, "already sent", TestTurnTokens.inert(T));
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertEquals(Injector.Cancellation.DELIVERED, injector.cancel(delivery),
|
||||
"a cancellation after pickup must not claim that the text stayed queued");
|
||||
assertEquals(List.of("already sent"), sent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void activeWhileQueuedOrInFlightThenQuietAfterTurnCompletes() {
|
||||
assertTrue(injector.activeTargets().isEmpty());
|
||||
@@ -378,7 +404,7 @@ class InjectorTest {
|
||||
void sendFailureDropsMessageAndFailsItsFuture() {
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T)).completion();
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
assertTrue(f.isCompletedExceptionally());
|
||||
@@ -387,7 +413,7 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void dropFailsPendingWaiters() {
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "orphan", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "orphan", TestTurnTokens.inert(T)).completion();
|
||||
injector.drop(T, new HerdrException("worker gone", "pane_not_found", null));
|
||||
assertTrue(f.isCompletedExceptionally(), "queued waiters unblock when the worker vanishes");
|
||||
}
|
||||
@@ -396,8 +422,8 @@ class InjectorTest {
|
||||
void dropPassesTheRealCauseForQueuedAndDeliveredWork() {
|
||||
Captor cap = new Captor();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap);
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> queued = inj.enqueue(T, "queued", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T)).completion();
|
||||
CompletableFuture<Void> queued = inj.enqueue(T, "queued", TestTurnTokens.inert(T)).completion();
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // deliver the first message
|
||||
inj.onStatus(T, AgentStatus.WORKING); // its turn is now in flight; one remains queued
|
||||
@@ -449,7 +475,7 @@ class InjectorTest {
|
||||
Captor cap = new Captor();
|
||||
List<String> forgotten = new ArrayList<>();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap, _ -> false, forgotten::add);
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "task", TestTurnTokens.inert(T)).completion();
|
||||
|
||||
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
@@ -563,7 +589,7 @@ class InjectorTest {
|
||||
StatusPoller poller = new StatusPoller(new AgentControl(idle), inj, 10);
|
||||
poller.start();
|
||||
try {
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "via-poller", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "via-poller", TestTurnTokens.inert(T)).completion();
|
||||
delivered.get(2, TimeUnit.SECONDS); // completes when the poller drives the send
|
||||
} finally {
|
||||
poller.stop();
|
||||
@@ -582,7 +608,7 @@ class InjectorTest {
|
||||
void deliveredFutureCarriesSendFailure() {
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T)).completion();
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
ExecutionException ex = assertThrows(ExecutionException.class, f::get);
|
||||
assertInstanceOf(HerdrException.class, ex.getCause());
|
||||
|
||||
@@ -41,7 +41,7 @@ class StatusPollerRoutingTest {
|
||||
poller.start();
|
||||
try {
|
||||
CompletableFuture<Void> delivered =
|
||||
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET));
|
||||
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)).completion();
|
||||
// Must resolve quickly: refining against the WRONG daemon (member) never classifies
|
||||
// out of UNKNOWN, so this would time out under the bug.
|
||||
delivered.get(2, TimeUnit.SECONDS);
|
||||
@@ -65,7 +65,7 @@ class StatusPollerRoutingTest {
|
||||
poller.start();
|
||||
try {
|
||||
CompletableFuture<Void> delivered =
|
||||
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET));
|
||||
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)).completion();
|
||||
assertThrows(TimeoutException.class, () -> delivered.get(500, TimeUnit.MILLISECONDS),
|
||||
"a lead target must never be refined from the member daemon's pane content");
|
||||
} finally {
|
||||
|
||||
@@ -23,6 +23,7 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
@@ -370,6 +371,161 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
"expected the pre-existing 'scrub blanks them' INFO unchanged, got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #341: {@code unprotectedGapLogged} guarded TWO WARN branches that name DIFFERENT env
|
||||
* var names — the allow-list branch ({@code keptByDerivedList}, below) and the deny-by-default
|
||||
* / non-zsh-fallback branch ({@link HerdrPeerLauncher#warnGapUnprotected}). {@code
|
||||
* memberCredentials} is a live, re-read-per-spawn supplier, so the policy can change between
|
||||
* two spawns on the same launcher instance — a config reload needs no restart. Spawn 1 runs
|
||||
* under {@code deny-by-default} with a gap of {@code SPAWN_ONE_UNCOVERED_TOKEN}, which trips
|
||||
* the (before this fix) SHARED one-shot flag. The policy is then reloaded to {@code
|
||||
* allow-list}; spawn 2's gap is {@code FLEETD_WORKER_TOKEN} instead — the test profile's own
|
||||
* {@code tokenEnv}, which the derived allow-list keeps even though it is on neither {@code
|
||||
* known:} nor {@code allow:}, so it is genuinely unprotected and deserves its own WARN. Before
|
||||
* this fix that WARN never fires, because the shared flag was already {@code true} — the
|
||||
* operator is never told {@code FLEETD_WORKER_TOKEN} reaches every member pane unblocked. Real
|
||||
* path: two real {@link HerdrPeerLauncher#spawn} calls on ONE launcher instance, with mutable
|
||||
* {@code memberCredentials}/host-env suppliers standing in for a live config reload between
|
||||
* spawns.
|
||||
*/
|
||||
@Test
|
||||
void aDifferentUnprotectedGapOnALaterSpawnIsNotSuppressedByAnEarlierSpawnsWarn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AtomicReference<FleetConfig.MemberCredentials> credsState = new AtomicReference<>(
|
||||
new FleetConfig.MemberCredentials(null, List.of(), List.of(), null)); // deny-by-default
|
||||
AtomicReference<Set<String>> hostEnvState =
|
||||
new AtomicReference<>(Set.of("SPAWN_ONE_UNCOVERED_TOKEN"));
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, credsState::get, "/bin/zsh", hostEnvState::get);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.WARN);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
// Spawn 1: deny-by-default, gap = {SPAWN_ONE_UNCOVERED_TOKEN} — the effectiveAllowed ==
|
||||
// null branch, via warnGapUnprotected.
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
// Live policy reload to allow-list, with a DIFFERENT gap name.
|
||||
credsState.set(new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null));
|
||||
hostEnvState.set(Set.of("FLEETD_WORKER_TOKEN"));
|
||||
// Spawn 2: allow-list, gap = {FLEETD_WORKER_TOKEN} — kept by the derived allow-list
|
||||
// (the profile's own tokenEnv), so it is the effectiveAllowed != null / keptByDerivedList
|
||||
// branch, at the SAME log line HerdrPeerLauncher:1819 guards with the shared flag.
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
List<String> messages = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
assertTrue(messages.stream().anyMatch(
|
||||
m -> m.contains("UNBLOCKED") && m.contains("SPAWN_ONE_UNCOVERED_TOKEN")),
|
||||
"spawn 1's deny-by-default gap must still warn — got: " + messages);
|
||||
assertTrue(messages.stream().anyMatch(
|
||||
m -> m.contains("UNBLOCKED") && m.contains("FLEETD_WORKER_TOKEN")),
|
||||
"spawn 2's gap names a DIFFERENT env var than spawn 1 (FLEETD_WORKER_TOKEN, not "
|
||||
+ "SPAWN_ONE_UNCOVERED_TOKEN) — it must still be warned about even though a "
|
||||
+ "flag already fired once for spawn 1's unrelated name. Before fleetd #341's "
|
||||
+ "fix this WARN never fires because unprotectedGapLogged was already true. "
|
||||
+ "Got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #341 follow-up: the OTHER half of the guard's contract. The set exists to report every
|
||||
* distinct name, but it must still report each one only ONCE — the noise control is the reason
|
||||
* a guard is here at all, and the ticket named it as invariant 1. Two spawns, same policy, same
|
||||
* gap name: exactly one WARN mentioning it.
|
||||
*
|
||||
* <p>Measured before this test existed: replacing {@code .filter(unprotectedGapNamesWarned::add)}
|
||||
* with a filter that adds and always returns {@code true} — so every name is logged on every
|
||||
* spawn — left all 1358 tests green. The fix was correct and nothing held it there. That is the
|
||||
* "a test on the seam does not prove the caller" shape: the {@code Set} behaves, and nothing
|
||||
* proved this class used it as a guard rather than as a record.
|
||||
*/
|
||||
@Test
|
||||
void theSameUnprotectedNameIsWarnedAboutOnlyOnceAcrossSpawns() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AtomicReference<FleetConfig.MemberCredentials> credsState = new AtomicReference<>(
|
||||
new FleetConfig.MemberCredentials(null, List.of(), List.of(), null)); // deny-by-default
|
||||
AtomicReference<Set<String>> hostEnvState =
|
||||
new AtomicReference<>(Set.of("REPEATED_UNCOVERED_TOKEN"));
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, credsState::get, "/bin/zsh", hostEnvState::get);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.WARN);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
// Same policy, same gap, second spawn. Nothing new to tell the operator.
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
List<String> messages = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
long mentioning = messages.stream()
|
||||
.filter(m -> m.contains("REPEATED_UNCOVERED_TOKEN"))
|
||||
.count();
|
||||
assertEquals(1, mentioning,
|
||||
"one unchanged unprotected name across two spawns must produce exactly one WARN — "
|
||||
+ "the set is a guard, not just a record. Got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #341 follow-up: the reverse policy order. The original defect was found going
|
||||
* deny-by-default then allow-list, and a guard that is fixed in one direction is not
|
||||
* necessarily fixed in the other — "ask which states still OPEN the gate". Here spawn 1 runs
|
||||
* under {@code allow-list} (the {@code keptByDerivedList} branch) and spawn 2 under
|
||||
* {@code deny-by-default} ({@link HerdrPeerLauncher#warnGapUnprotected}), with a different name
|
||||
* each time. Both must be reported.
|
||||
*/
|
||||
@Test
|
||||
void anAllowListWarnDoesNotSuppressALaterDenyByDefaultWarnForADifferentName() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AtomicReference<FleetConfig.MemberCredentials> credsState = new AtomicReference<>(
|
||||
new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null));
|
||||
AtomicReference<Set<String>> hostEnvState =
|
||||
new AtomicReference<>(Set.of("FLEETD_WORKER_TOKEN"));
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, credsState::get, "/bin/zsh", hostEnvState::get);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.WARN);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
// Spawn 1: allow-list, gap kept by the derived list — the keptByDerivedList WARN.
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
// Live reload the OTHER way: back to deny-by-default, with a different name.
|
||||
credsState.set(new FleetConfig.MemberCredentials(null, List.of(), List.of(), null));
|
||||
hostEnvState.set(Set.of("LATER_UNCOVERED_TOKEN"));
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
List<String> messages = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
assertTrue(messages.stream().anyMatch(
|
||||
m -> m.contains("UNBLOCKED") && m.contains("FLEETD_WORKER_TOKEN")),
|
||||
"spawn 1's allow-list gap must warn — got: " + messages);
|
||||
assertTrue(messages.stream().anyMatch(
|
||||
m -> m.contains("UNBLOCKED") && m.contains("LATER_UNCOVERED_TOKEN")),
|
||||
"spawn 2's deny-by-default gap names a different variable and must still be warned "
|
||||
+ "about, even though an allow-list WARN already fired. Got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: with {@code memberHerdrSocket:} configured, member panes run under a
|
||||
* different OS user — {@link HerdrPeerLauncher#hostEnvNames} describes fleetd's own process, not
|
||||
|
||||
@@ -0,0 +1,264 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import com.rabbitmq.client.AMQP;
|
||||
import com.rabbitmq.client.Channel;
|
||||
import com.rabbitmq.client.Connection;
|
||||
import com.rabbitmq.client.DeliverCallback;
|
||||
import com.rabbitmq.client.Delivery;
|
||||
import com.rabbitmq.client.Envelope;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.Timeout;
|
||||
|
||||
import java.lang.reflect.InvocationHandler;
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.time.Duration;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-318: a delivery landing on the consumer work-pool thread <em>after</em>
|
||||
* {@link AmqpReplyInbox#release} has already swapped the target's {@code held} entry for its
|
||||
* tombstone, but <em>before</em> {@code release()} itself returns, must be nacked-with-requeue —
|
||||
* never silently retained in a fresh map {@code release()} has already stopped looking at.
|
||||
*
|
||||
* <p><strong>This forces the actual interleaving, not a sequence of calls.</strong> {@code release()}
|
||||
* runs on its own thread and is made to block <em>inside</em> its nack loop, via a fake
|
||||
* {@link Channel} whose {@code basicNack} blocks on its first invocation. That block is only
|
||||
* reachable after {@code release()}'s {@code held.compute(...)} has already swapped in the
|
||||
* {@code RELEASED} tombstone (the compute call happens strictly before the loop that calls
|
||||
* {@code basicNack}), so observing it is direct, ordering-guaranteed proof that the tombstone is in
|
||||
* place and {@code release()} has not yet returned — still holding {@code channelLock} — when a
|
||||
* second thread fires {@code own()}'s captured {@link DeliverCallback} for a brand-new message on the
|
||||
* same target. No mocking library is on the classpath, so the fake broker is a {@link Proxy}, the
|
||||
* same pattern {@code AmqpReplyInboxRecoveryRaceTest} already uses.
|
||||
*
|
||||
* <p><strong>What this does and does not prove.</strong> It proves that a delivery whose
|
||||
* {@code computeIfAbsent} call is ordered strictly after {@code release()}'s tombstone swap — while
|
||||
* {@code release()} is still running — is nacked-with-requeue rather than silently parked forever.
|
||||
* It does not drive a real broker: {@code basicNack} here is a recorded call on a fake channel, not a
|
||||
* verified requeue-and-redeliver. That half of the contract (a nacked-with-requeue delivery really
|
||||
* does come back to a later owner) is already covered against a real broker by
|
||||
* {@code AmqpReplyInboxContractTest.releaseCancelsConsumerAndRequeuesHeldDeliveryForRecovery}, which
|
||||
* this test does not duplicate.
|
||||
*/
|
||||
class AmqpReplyInboxReleaseRaceTest {
|
||||
|
||||
@Test
|
||||
@Timeout(15)
|
||||
void deliveryArrivingWhileReleaseIsStillRunningIsNackedNotStranded() throws Exception {
|
||||
String target = "worker-release-race";
|
||||
|
||||
List<long[]> nacks = new CopyOnWriteArrayList<>(); // {deliveryTag, requeue(1/0)}
|
||||
List<Long> acks = new CopyOnWriteArrayList<>();
|
||||
AtomicReference<DeliverCallback> deliverCallback = new AtomicReference<>();
|
||||
CountDownLatch nackStarted = new CountDownLatch(1);
|
||||
CountDownLatch releaseMayFinishNack = new CountDownLatch(1);
|
||||
AtomicInteger nackCallCount = new AtomicInteger();
|
||||
|
||||
Channel consumeChannel = fakeConsumeChannel(deliverCallback, nacks, acks, nackCallCount,
|
||||
nackStarted, releaseMayFinishNack);
|
||||
Channel publishChannel = fakeInertChannel();
|
||||
Connection connection = fakeConnection(consumeChannel, publishChannel);
|
||||
|
||||
AmqpReplyInbox inbox = new AmqpReplyInbox(connection, AmqpReplyInbox.DEFAULT_PREFETCH);
|
||||
inbox.own(target);
|
||||
assertTrue(deliverCallback.get() != null, "own() must have registered a DeliverCallback");
|
||||
|
||||
// Seed one already-held delivery (m0) so release()'s nack loop has something to iterate, and
|
||||
// therefore somewhere to block, before it can return.
|
||||
deliverCallback.get().handle("ctag", delivery(1L, "m0", "first"));
|
||||
|
||||
AtomicReference<Throwable> releaseError = new AtomicReference<>();
|
||||
Thread releaseThread = new Thread(() -> {
|
||||
try {
|
||||
inbox.release(target);
|
||||
} catch (Throwable t) {
|
||||
releaseError.set(t);
|
||||
}
|
||||
}, "release-under-test");
|
||||
releaseThread.start();
|
||||
|
||||
// This latch only fires from inside the fake channel's basicNack — i.e. from inside
|
||||
// release()'s nack loop, which release()'s code only reaches AFTER held.compute(...) has
|
||||
// already swapped in RELEASED. Waiting for it is direct proof the swap has happened and
|
||||
// release() has not yet returned (it is stuck mid-loop, still holding channelLock).
|
||||
assertTrue(nackStarted.await(10, TimeUnit.SECONDS),
|
||||
"release() never reached its nack loop — it may not have started");
|
||||
|
||||
// The exact interleaving CB-318 describes: a delivery for a NEW message on the same target
|
||||
// lands on the "consumer work-pool thread" (this second thread) while release() is still
|
||||
// running. With the pre-fix code (a bare held.remove(target)) this created a brand-new map
|
||||
// under computeIfAbsent that release() — already past its remove — never looks at again.
|
||||
AtomicReference<Throwable> deliveryError = new AtomicReference<>();
|
||||
Thread deliveryThread = new Thread(() -> {
|
||||
try {
|
||||
deliverCallback.get().handle("ctag", delivery(2L, "m1", "second"));
|
||||
} catch (Throwable t) {
|
||||
deliveryError.set(t);
|
||||
}
|
||||
}, "concurrent-delivery");
|
||||
deliveryThread.start();
|
||||
|
||||
// Head start for the delivery thread to reach (and, on the fixed code, block on)
|
||||
// channelLock — release() still holds it at this point, so a correct fix cannot have
|
||||
// resolved m1's nack yet. Purely in-memory work (computeIfAbsent, a reference compare)
|
||||
// separates deliveryThread.start() from that block point, so 300ms is a large margin, not a
|
||||
// tight timing assumption.
|
||||
Thread.sleep(300);
|
||||
assertEquals(1, nackCallCount.get(),
|
||||
"the concurrent delivery must not resolve its nack before release() gives up "
|
||||
+ "channelLock — if this is 2 already, the interleaving below is not being "
|
||||
+ "tested, only a sequential call");
|
||||
|
||||
releaseMayFinishNack.countDown(); // let release() finish nacking m0 and return
|
||||
assertTrue(releaseThread.join(Duration.ofSeconds(10)), "release() did not finish");
|
||||
assertTrue(deliveryThread.join(Duration.ofSeconds(10)), "the concurrent delivery did not finish");
|
||||
|
||||
assertNull(releaseError.get(), "release() threw: " + releaseError.get());
|
||||
assertNull(deliveryError.get(), "the concurrent delivery threw: " + deliveryError.get());
|
||||
|
||||
assertEquals(2, nacks.size(),
|
||||
"both the pre-held m0 and the concurrently-arriving m1 must be nacked, got: "
|
||||
+ nacks.stream().map(n -> "[tag=" + n[0] + " requeue=" + n[1] + "]").toList());
|
||||
assertTrue(nacks.stream().allMatch(n -> n[1] == 1L),
|
||||
"invariant 1 (never drop): every nack must set requeue=true");
|
||||
assertTrue(nacks.stream().anyMatch(n -> n[0] == 1L), "m0's delivery tag must be nacked");
|
||||
assertTrue(nacks.stream().anyMatch(n -> n[0] == 2L),
|
||||
"m1 — delivered while release() was still running, after the tombstone swap — must be "
|
||||
+ "nacked, not silently retained in a map release() will never look at again");
|
||||
assertTrue(acks.isEmpty(), "invariant 1 (never drop): a held reply must never be basicAck'd");
|
||||
}
|
||||
|
||||
private static Delivery delivery(long tag, String msgId, String body) {
|
||||
Envelope envelope = new Envelope(tag, false, "", "irrelevant");
|
||||
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder().messageId(msgId).build();
|
||||
return new Delivery(envelope, props, body.getBytes(StandardCharsets.UTF_8));
|
||||
}
|
||||
|
||||
/** A {@link Proxy}-backed consume {@link Channel}: blocks the FIRST {@code basicNack} call on
|
||||
* {@code releaseMayFinishNack}, after signalling {@code nackStarted} — everything else records
|
||||
* the call and returns a harmless default, matching the style already used by
|
||||
* {@code AmqpReplyInboxRecoveryRaceTest}. */
|
||||
private static Channel fakeConsumeChannel(AtomicReference<DeliverCallback> deliverCallback,
|
||||
List<long[]> nacks, List<Long> acks,
|
||||
AtomicInteger nackCallCount,
|
||||
CountDownLatch nackStarted,
|
||||
CountDownLatch releaseMayFinishNack) {
|
||||
InvocationHandler handler = (proxy, method, args) -> {
|
||||
String name = method.getName();
|
||||
if (name.equals("basicConsume")) {
|
||||
deliverCallback.set((DeliverCallback) args[2]);
|
||||
return "ctag";
|
||||
}
|
||||
if (name.equals("basicNack")) {
|
||||
long tag = (long) args[0];
|
||||
boolean requeue = (boolean) args[2];
|
||||
if (nackCallCount.incrementAndGet() == 1) {
|
||||
nackStarted.countDown();
|
||||
if (!releaseMayFinishNack.await(10, TimeUnit.SECONDS)) {
|
||||
throw new IllegalStateException("test never released the nack latch");
|
||||
}
|
||||
}
|
||||
nacks.add(new long[] {tag, requeue ? 1L : 0L});
|
||||
return null;
|
||||
}
|
||||
if (name.equals("basicAck")) {
|
||||
acks.add((long) args[0]);
|
||||
return null;
|
||||
}
|
||||
if (name.equals("equals")) {
|
||||
return proxy == args[0];
|
||||
}
|
||||
if (name.equals("hashCode")) {
|
||||
return System.identityHashCode(proxy);
|
||||
}
|
||||
if (name.equals("toString")) {
|
||||
return "FakeConsumeChannel";
|
||||
}
|
||||
return defaultValue(method.getReturnType());
|
||||
};
|
||||
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
|
||||
new Class<?>[] {Channel.class}, handler);
|
||||
}
|
||||
|
||||
/** A {@link Proxy}-backed {@link Channel} that answers every call with a harmless default — used
|
||||
* as the publish channel, which this test never actually publishes on. */
|
||||
private static Channel fakeInertChannel() {
|
||||
InvocationHandler handler = (proxy, method, args) -> {
|
||||
String name = method.getName();
|
||||
if (name.equals("equals")) {
|
||||
return proxy == args[0];
|
||||
}
|
||||
if (name.equals("hashCode")) {
|
||||
return System.identityHashCode(proxy);
|
||||
}
|
||||
if (name.equals("toString")) {
|
||||
return "FakeInertChannel";
|
||||
}
|
||||
return defaultValue(method.getReturnType());
|
||||
};
|
||||
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
|
||||
new Class<?>[] {Channel.class}, handler);
|
||||
}
|
||||
|
||||
/** A {@link Proxy}-backed {@link Connection} handing out {@code first} then {@code second} from
|
||||
* successive {@code createChannel()} calls, matching {@link AmqpReplyInbox}'s constructor. */
|
||||
private static Connection fakeConnection(Channel first, Channel second) {
|
||||
AtomicInteger calls = new AtomicInteger();
|
||||
InvocationHandler handler = (proxy, method, args) -> {
|
||||
String name = method.getName();
|
||||
if (name.equals("createChannel") && (args == null || args.length == 0)) {
|
||||
return calls.getAndIncrement() == 0 ? first : second;
|
||||
}
|
||||
if (name.equals("equals")) {
|
||||
return proxy == args[0];
|
||||
}
|
||||
if (name.equals("hashCode")) {
|
||||
return System.identityHashCode(proxy);
|
||||
}
|
||||
if (name.equals("toString")) {
|
||||
return "FakeConnection";
|
||||
}
|
||||
return defaultValue(method.getReturnType());
|
||||
};
|
||||
return (Connection) Proxy.newProxyInstance(AmqpReplyInboxReleaseRaceTest.class.getClassLoader(),
|
||||
new Class<?>[] {Connection.class}, handler);
|
||||
}
|
||||
|
||||
private static Object defaultValue(Class<?> type) {
|
||||
if (!type.isPrimitive() || type == void.class) {
|
||||
return null;
|
||||
}
|
||||
if (type == boolean.class) {
|
||||
return Boolean.FALSE;
|
||||
}
|
||||
if (type == long.class) {
|
||||
return 0L;
|
||||
}
|
||||
if (type == short.class) {
|
||||
return (short) 0;
|
||||
}
|
||||
if (type == byte.class) {
|
||||
return (byte) 0;
|
||||
}
|
||||
if (type == char.class) {
|
||||
return (char) 0;
|
||||
}
|
||||
if (type == double.class) {
|
||||
return 0.0d;
|
||||
}
|
||||
if (type == float.class) {
|
||||
return 0.0f;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,9 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
@@ -11,6 +15,7 @@ import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
@@ -253,7 +258,7 @@ class MessageServiceTest {
|
||||
injector.onStatus(T, AgentStatus.IDLE); // first delivery
|
||||
injector.onStatus(T, AgentStatus.WORKING); // first turn in flight
|
||||
|
||||
CompletableFuture<Void> queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T)).completion();
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(T);
|
||||
injector.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
|
||||
|
||||
@@ -405,6 +410,29 @@ class MessageServiceTest {
|
||||
"a delivered send whose worker never replies times out as still working");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #345. This forces the injector to pick up the exact pending delivery after {@code send}
|
||||
* first observes its completion as incomplete, but before {@code cancel} takes the target monitor.
|
||||
* The timeout must use {@link Injector.Cancellation#DELIVERED} from {@code cancel} and report
|
||||
* {@link MessageService.Outcome#TIMED_OUT_WORKING}, because the text landed.
|
||||
*
|
||||
* <p>What this does not prove: that this precise interleaving happens by itself under production
|
||||
* timing. The test forces it through a test-only hook; it proves the timeout caller handles the
|
||||
* injector result when the interleaving occurs.
|
||||
*/
|
||||
@Test
|
||||
void sendTimeoutUsesCancellationDeliveredWhenPickupWinsTheRace() {
|
||||
messages.setTimeoutCancellationRaceHookForTest(() -> injector.onStatus(T, AgentStatus.IDLE));
|
||||
try {
|
||||
MessageService.Reply reply = messages.send(T, "race delivery", 50);
|
||||
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, reply.outcome(),
|
||||
"cancel reporting DELIVERED means the worker received the timed-out message");
|
||||
} finally {
|
||||
messages.setTimeoutCancellationRaceHookForTest(null);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
@@ -940,6 +968,54 @@ class MessageServiceTest {
|
||||
assertEquals("PR opened: https://example/pulls/42", view.reply());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #324: {@code answer()} holds {@code sessionLocks} for the target and, once the worker's
|
||||
* real terminal reply arrives, calls {@code finishAsyncTask}, which used to read the volatile
|
||||
* {@code task.turnId} twice — once to check it is non-null, once as the key for
|
||||
* {@code asyncTasksByTurn.remove}. {@code ask()}'s own timeout path mutates the same field with no
|
||||
* lock at all. This test does not wait for a real race to land in that narrow window between the
|
||||
* two reads — instead it drives the exact sequence the ticket describes (worker asks, primary
|
||||
* answers, worker's real reply arrives) and, via a package-private test hook wired to fire at
|
||||
* precisely that point, runs the identical production cleanup {@code ask()}'s timeout catch block
|
||||
* runs ({@code clearAsyncQuestion(turnId, true)}) so the field goes {@code null} between the two
|
||||
* reads deterministically rather than by chance.
|
||||
*
|
||||
* <p>What this proves: given that exact interleaving, {@code answer()} must not throw and the
|
||||
* ticket must still resolve to the worker's real reply. What it does not prove: that the
|
||||
* interleaving itself is reachable in production — that is established by reading the code (see
|
||||
* the ticket), not by this test, since forcing it via a hook is not the same as two independent
|
||||
* threads racing on their own schedules.
|
||||
*/
|
||||
@Test
|
||||
void finishAsyncTaskSurvivesTurnIdGoingNullBetweenItsTwoReads() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
String turnId = asking.turnId();
|
||||
|
||||
// Fire ask()'s own unlocked timeout cleanup at the moment finishAsyncTask has already checked
|
||||
// task.turnId is non-null but has not yet used it — the exact torn-read window fleetd #324
|
||||
// describes.
|
||||
messages.setFinishAsyncTaskRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting(); // answer() opened its own forward waiter for the resumed worker turn
|
||||
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
|
||||
"the lead's own answer() call must not throw because ask()'s timeout cleanup raced it");
|
||||
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("PR opened: https://example/pulls/42", done.reply(),
|
||||
"the ticket must still resolve to the worker's real reply despite the forced race");
|
||||
}
|
||||
|
||||
@Test
|
||||
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
@@ -993,6 +1069,105 @@ class MessageServiceTest {
|
||||
"the reply completed its own ticket directly and never touched the inbox");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #329 (F1). {@code answer()} completes the async ticket by looking {@code turnId} up in
|
||||
* {@code asyncTasksByTurn} a SECOND time (the first is at :991, purely to re-register the
|
||||
* {@code asyncTasksByWaiter} entry for #282's chained-ask case). That second lookup races
|
||||
* {@code ask()}'s own unlocked timeout cleanup ({@code clearAsyncQuestion(turnId, true)}, run from
|
||||
* {@code markAskTimedOut} + forgetting): {@code ask()}'s {@code ticket.answer().get(timeoutMillis)}
|
||||
* can time out at essentially the same instant {@code answer()}'s {@code rendezvous.answerAsk}
|
||||
* call above already succeeded and unblocked the worker. fleetd #324's single-read fix does not
|
||||
* help here — that fixed a torn read of one already-held {@link MessageService} internal
|
||||
* {@code Task}; this is a second, independent map lookup by a caller that no longer holds the
|
||||
* {@code Task} it already found once.
|
||||
*
|
||||
* <p>This test does not wait for that real race to land on its own schedule — it drives the exact
|
||||
* sequence the ticket describes (worker asks, primary answers, worker's real reply arrives) and
|
||||
* fires the identical production cleanup {@code ask()}'s timeout path runs
|
||||
* ({@code clearAsyncQuestion(turnId, true)}, via the {@code forgetTurnForTest} seam fleetd #324
|
||||
* already merged) at the point between "the worker's own {@code ask()} call has unblocked" and
|
||||
* "the worker's real {@code fleet_reply} arrives" — the exact window fleetd #329 names.
|
||||
*
|
||||
* <p>What this proves: given that exact interleaving, the async ticket must still resolve to the
|
||||
* worker's real reply, not stay stuck {@link MessageService.Phase#PENDING} forever. What it does
|
||||
* not prove: that the interleaving itself is reachable in production on its own timing — that is
|
||||
* established by reading the code (see the ticket's "path in"), not by this test, for the same
|
||||
* reason fleetd #324's own race test says so.
|
||||
*/
|
||||
@Test
|
||||
void aReplyRacingAsksTimeoutCleanupStillCompletesTheAsyncTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
String turnId = asking.turnId();
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(turnId, "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer(),
|
||||
"the worker's own ask() call must have already unblocked with the primary's answer "
|
||||
+ "before we force the race below");
|
||||
|
||||
// ask()'s own unlocked timeout cleanup can forget this exact turnId at essentially the same
|
||||
// instant answer() has already unblocked the worker and is now waiting on the resumed turn's
|
||||
// real reply — reproduce that interleaving directly instead of trying to win a real race.
|
||||
messages.forgetTurnForTest(turnId);
|
||||
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome(),
|
||||
"the primary's own answer() call must still see the worker's real reply");
|
||||
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("PR opened: https://example/pulls/42", done.reply(),
|
||||
"fleet_poll{ticket} must return the worker's real reply, not stay PENDING forever "
|
||||
+ "just because ask()'s timeout cleanup forgot this turnId first");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #329 (F3). {@link MessageService#reply} reads the same orphan {@code Task}'s {@code
|
||||
* turnId} twice on this path (once to check it is non-null, once as the {@code
|
||||
* asyncTasksByTurn.remove} key) — the exact double-read shape fleetd #324 fixed in {@code
|
||||
* finishAsyncTask}. This test drives the same real sequence as {@code
|
||||
* aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket} (worker asks, primary answers, the
|
||||
* primary's own bounded wait for the resumed turn expires) to reach a task with {@code turnId}
|
||||
* genuinely stamped and no live rendezvous waiter open — the state {@code askAnsweredAsyncTasks}
|
||||
* matches here — then fires the identical production cleanup {@code ask()}'s own timeout path
|
||||
* runs ({@code clearAsyncQuestion(turnId, true)}, via the {@code forgetTurnForTest} seam fleetd
|
||||
* #324 merged) at the point between the check and the removal use, via a dedicated test hook
|
||||
* mirroring {@code finishAsyncTaskRaceHook}.
|
||||
*/
|
||||
@Test
|
||||
void replyToAnOrphanedTaskSurvivesTurnIdGoingNullBetweenItsTwoReads() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
String turnId = asking.turnId();
|
||||
|
||||
MessageService.Reply answerReply = messages.answer(turnId, "config.yaml", 150);
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, answerReply.outcome(),
|
||||
"the primary's own bounded wait must give up first, leaving turnId stamped with no "
|
||||
+ "live waiter — the state reply()'s F3 code path matches");
|
||||
|
||||
messages.setReplyOrphanTurnIdRaceHookForTest(() -> messages.forgetTurnForTest(turnId));
|
||||
try {
|
||||
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
|
||||
|
||||
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("PR opened: https://example/pulls/42", done.reply(),
|
||||
"the ticket must still resolve to the worker's real reply despite the forced race");
|
||||
} finally {
|
||||
messages.setReplyOrphanTurnIdRaceHookForTest(null);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #307's ambiguity guard: an ask timeout frees its target ({@code hasAsyncQuestion}
|
||||
* becomes false the instant it lapses — proven above), so a second, independent delegation can
|
||||
@@ -1103,6 +1278,63 @@ class MessageServiceTest {
|
||||
assertEquals("done", done.reply());
|
||||
}
|
||||
|
||||
// --- fleetd #329 (F2): an exception after the ticket's future completes must reach a log -----
|
||||
//
|
||||
// sendAsync's own executor task ends with catch (Throwable t) { task.future.completeExceptionally(t); }
|
||||
// — but finishAsyncTask completes that same future on its first line, so anything that throws
|
||||
// afterward hits an already-completed future: completeExceptionally returns false and does
|
||||
// nothing, and (before this fix) nothing logged it either. Measured on the ticket: temporarily
|
||||
// reintroducing the fleetd #324 NPE reproduced 19 real exceptions on the ordinary path with
|
||||
// 74/74 tests staying green and zero log lines. There is no reachable production call site where
|
||||
// finishAsyncTask throws after completing the future (fleetd #324 already closed the one that
|
||||
// used to), so this test drives the shape directly via a dedicated test-only hook
|
||||
// (afterFinishAsyncTaskCompleteHookForTest) rather than trying to engineer a real exception into
|
||||
// that narrow window — the same technique fleetd #324's own race test and this ticket's F1/F3
|
||||
// tests use for their own races.
|
||||
|
||||
private static ListAppender<ILoggingEvent> attachMessageServiceLog() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(MessageService.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detachMessageServiceLog(ListAppender<ILoggingEvent> appender) {
|
||||
((Logger) LoggerFactory.getLogger(MessageService.class)).detachAppender(appender);
|
||||
}
|
||||
|
||||
@Test
|
||||
void anExceptionAfterTheTicketFutureCompletesStillReachesTheLog() throws Exception {
|
||||
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
|
||||
try {
|
||||
messages.setAfterFinishAsyncTaskCompleteHookForTest(() -> {
|
||||
throw new RuntimeException("PROBE-329-F2");
|
||||
});
|
||||
|
||||
String ticket = messages.sendAsync(T, "do the task");
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
|
||||
// finishAsyncTask's own first line already completed the future with the real result
|
||||
// BEFORE the hook threw — F2 is a visibility gap, not a correctness gap for the ticket
|
||||
// itself, so the ticket's own outcome must be unaffected by the swallowed exception.
|
||||
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertEquals("done", done.reply());
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.ERROR
|
||||
&& e.getFormattedMessage().contains(ticket)
|
||||
&& e.getThrowableProxy() != null
|
||||
&& "PROBE-329-F2".equals(e.getThrowableProxy().getMessage())),
|
||||
"an exception thrown after the ticket's future already completed must still "
|
||||
+ "reach the log, not vanish silently");
|
||||
} finally {
|
||||
messages.setAfterFinishAsyncTaskCompleteHookForTest(null);
|
||||
detachMessageServiceLog(appender);
|
||||
}
|
||||
}
|
||||
|
||||
// --- CB-582: fleet_status pendingAsk() ------------------------------------------------------
|
||||
|
||||
@Test
|
||||
@@ -1559,7 +1791,10 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome());
|
||||
|
||||
assertTrue(messages.hasQueuedDelivery(T),
|
||||
"a TIMED_OUT_QUEUED send leaves the message still queued in the injector");
|
||||
"a TIMED_OUT_QUEUED send still records the undelivered delivery for fleet health");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
assertTrue(herdr.calls.stream().noneMatch(c -> c.method().equals("agent.prompt")),
|
||||
"a TIMED_OUT_QUEUED send must be cancelled, not delivered when the worker later goes idle");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -86,17 +86,50 @@ class SessionManagerTest {
|
||||
private volatile RuntimeException hasUncommittedFailure;
|
||||
private volatile RuntimeException snapshotFailure;
|
||||
private final java.util.concurrent.atomic.AtomicLong snapshotSeq = new java.util.concurrent.atomic.AtomicLong();
|
||||
/** fleetd #316: successive {@code hasUncommitted} answers, one per call, last one sticky
|
||||
* once exhausted — models a worktree whose state changes between reads. Empty (the
|
||||
* default) falls back to the plain {@link #dirty} flag, so every existing test using this
|
||||
* fake keeps returning one fixed answer. */
|
||||
private final List<Boolean> dirtySequence = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
private int failHasUncommittedOnCall = -1;
|
||||
private RuntimeException hasUncommittedCallFailure;
|
||||
private final java.util.concurrent.atomic.AtomicInteger hasUncommittedCalls =
|
||||
new java.util.concurrent.atomic.AtomicInteger();
|
||||
|
||||
RecordingWorktrees dirty(boolean dirty) {
|
||||
this.dirty = dirty;
|
||||
return this;
|
||||
}
|
||||
|
||||
/** fleetd #316: return {@code answers[0]} on the first {@code hasUncommitted} call,
|
||||
* {@code answers[1]} on the second, and so on; the last element repeats after that. */
|
||||
RecordingWorktrees dirtySequence(boolean... answers) {
|
||||
for (boolean a : answers) {
|
||||
dirtySequence.add(a);
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
int hasUncommittedCallCount() {
|
||||
return hasUncommittedCalls.get();
|
||||
}
|
||||
|
||||
RecordingWorktrees failHasUncommittedWith(RuntimeException e) {
|
||||
this.hasUncommittedFailure = e;
|
||||
return this;
|
||||
}
|
||||
|
||||
/**
|
||||
* Throw from {@code hasUncommitted} on one specific call only, counting from 0. The
|
||||
* whole-double {@link #failHasUncommittedWith} cannot express fleetd #316's fail-safe
|
||||
* case, which needs the pre-stop read to succeed and only the late read to fail.
|
||||
*/
|
||||
RecordingWorktrees failHasUncommittedOnCall(int call, RuntimeException e) {
|
||||
this.failHasUncommittedOnCall = call;
|
||||
this.hasUncommittedCallFailure = e;
|
||||
return this;
|
||||
}
|
||||
|
||||
RecordingWorktrees failRemoveFor(String worktreePath) {
|
||||
failRemoveFor.add(worktreePath);
|
||||
return this;
|
||||
@@ -127,9 +160,16 @@ class SessionManagerTest {
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
int call = hasUncommittedCalls.getAndIncrement();
|
||||
if (hasUncommittedFailure != null) {
|
||||
throw hasUncommittedFailure;
|
||||
}
|
||||
if (call == failHasUncommittedOnCall) {
|
||||
throw hasUncommittedCallFailure;
|
||||
}
|
||||
if (!dirtySequence.isEmpty()) {
|
||||
return dirtySequence.get(Math.min(call, dirtySequence.size() - 1));
|
||||
}
|
||||
return dirty;
|
||||
}
|
||||
|
||||
@@ -1207,6 +1247,104 @@ class SessionManagerTest {
|
||||
+ "dirty check threw");
|
||||
}
|
||||
|
||||
// --- fleetd #316: the dirty check must be re-taken after the worker is stopped, not trusted
|
||||
// stale from before it ------------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void releaseDoesNotRemoveAWorktreeThatBecameDirtyBetweenTheFirstCheckAndRemoval() {
|
||||
// Models the exact race #316 reports: hasUncommitted answers clean while the worker is
|
||||
// still running (call 1), the worker then writes new work, and by the time release is
|
||||
// about to force-remove the worktree a second read (call 2) would see it as dirty. Without
|
||||
// the fix this test fails: release() never re-reads and force-removes the worktree anyway.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees().dirtySequence(false, true);
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-316a", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertTrue(worktrees.removeCalls().isEmpty(),
|
||||
"a worktree that turned dirty between the pre-stop read and removal must be preserved");
|
||||
assertEquals(2, worktrees.hasUncommittedCallCount(),
|
||||
"the fix re-reads hasUncommitted exactly once more, immediately before removal");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releasePreservesAWorktreeWhoseLateRecheckCannotBeRead() {
|
||||
// fleetd #316 invariant 1, which no test pinned when the fix landed: the late re-check
|
||||
// fails toward PRESERVING. Found by mutation — flipping dirtyImmediatelyBeforeRemoval's
|
||||
// catch from `return true` to `return false` turned the guard into a cause of the very
|
||||
// data loss it was added to stop, and the whole suite stayed green. The pre-stop read
|
||||
// succeeds and says clean (call 0); the read that authorises the removal throws (call 1).
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees()
|
||||
.dirtySequence(false)
|
||||
.failHasUncommittedOnCall(1, new WorktreeException("git status exited 128"));
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-316c", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertTrue(worktrees.removeCalls().isEmpty(),
|
||||
"a worktree whose state cannot be read immediately before removal must be kept: "
|
||||
+ "preserving costs disk, deleting on a guess destroys work with no other copy");
|
||||
assertEquals(2, worktrees.hasUncommittedCallCount(),
|
||||
"the late re-check still runs — it is the throwing call, not a skipped one");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseSnapshotsWorkFoundOnlyByTheLateRecheck() {
|
||||
// #316's second half: the pre-stop dirty=false means trySnapshot never ran for this
|
||||
// session, so the late-discovered work would otherwise have no refs/wip/* copy at all —
|
||||
// only the on-disk preserve. The re-check path must snapshot it too.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees().dirtySequence(false, true);
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-316b", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(java.util.List.of(s.worktree()), worktrees.snapshotCalls(),
|
||||
"the newly-dirty worktree is snapshotted even though the pre-stop check saw it clean");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseStillRemovesAWorktreeThatStaysCleanOnTheLateRecheck() {
|
||||
// The ordinary, non-racing case: nothing else changes behaviour when the second read
|
||||
// agrees with the first.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees().dirty(false);
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-316c", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(java.util.List.of(s.worktree()), worktrees.removeCalls(),
|
||||
"a worktree that is still clean on the late recheck is removed as before");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseNeverReChecksAWorktreeAlreadyPreservedByTheFirstDirtyCheck() {
|
||||
// Invariant 4 from #316: no second unconditional git status. A release that already
|
||||
// decided to preserve (the ordinary CB-576 dirty path) must not pay for a second read.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
RecordingWorktrees worktrees = new RecordingWorktrees().dirty(true);
|
||||
SessionManager sessions = sessionManager(herdr, worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-316d", null));
|
||||
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(1, worktrees.hasUncommittedCallCount(),
|
||||
"a release that already preserves on the first read must not re-check before "
|
||||
+ "skipping the removal it was never going to do");
|
||||
assertTrue(worktrees.removeCalls().isEmpty());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #283 defect 1 changed this test's own premise, so its assertions are updated along
|
||||
* with the production fix. Before #283, the middle session's worktree-removal failure escaped
|
||||
|
||||
Reference in New Issue
Block a user