Compare commits
13 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 554395b104 | |||
| 823976c1b5 | |||
| c8388a7f92 | |||
| 02e6aef98c | |||
| 6d493bc7bb | |||
| e5cb51a90e | |||
| e545c08082 | |||
| efa0deb9b2 | |||
| fa1f49675b | |||
| 8426c3528f | |||
| 65f98ba910 | |||
| 667254df47 | |||
| 2926cd1784 |
@@ -22,7 +22,7 @@ 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>
|
||||
@@ -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,63 @@ 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>Split</strong> (fleetd #330) — read <em>both</em> ways at different sites, so the
|
||||
* key does not fit any class above as a whole: {@code health:} and {@code coordinator:}.
|
||||
* 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>
|
||||
* </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: {@code bind:},
|
||||
* {@code herdrSocket:}, {@code broker:} and {@code auth:}. The socket is bound, the broker
|
||||
* connection is open, and the auth mode decides who may reach the port that is already
|
||||
* listening.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p><strong>The denominator, measured on 2026-09-04 (fleetd #330).</strong> {@code FleetConfig} has
|
||||
* 22 top-level record components. Two of them are named nowhere in this file, and the reason is the
|
||||
* same for both: {@code memberCredentials} and {@code memberLoginShell} are <strong>hot</strong> and
|
||||
* correctly absent — both are read live off {@code config.get()} 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.
|
||||
* <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. Twice now — {@code worktreeGroup} (#323) and
|
||||
* {@code primary}/{@code configReload} (#326) — 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}.
|
||||
*
|
||||
* <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 +137,24 @@ 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");
|
||||
|
||||
private final Path path;
|
||||
private final AtomicReference<FleetConfig> current;
|
||||
|
||||
@@ -108,25 +182,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 +224,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,8 +276,9 @@ 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;
|
||||
}
|
||||
@@ -234,6 +328,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 +387,32 @@ 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). Unlike {@link #changedDeferredKeys}, this does not try to
|
||||
* tell which sub-field moved: any change to {@code health:} or {@code coordinator:} 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.
|
||||
*/
|
||||
private 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");
|
||||
}
|
||||
// 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 two split keys the class doc documents.
|
||||
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;
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link FleetConfig.Profile} record components deliberately left out of
|
||||
* {@link #sameLaunchSettings} because they are read <em>live</em>, not baked in at spawn — see
|
||||
|
||||
@@ -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)) {
|
||||
|
||||
@@ -1230,14 +1230,61 @@ 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);
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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;
|
||||
}
|
||||
|
||||
/**
|
||||
* 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);
|
||||
}
|
||||
|
||||
/** Complete the async ticket correlated to a specific answered turn. */
|
||||
private void finishAsyncTask(String turnId, Reply result) {
|
||||
Task task = asyncTasksByTurn.get(turnId);
|
||||
|
||||
@@ -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,219 @@ 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());
|
||||
}
|
||||
|
||||
/**
|
||||
* 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,167 @@
|
||||
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
|
||||
* — verified by reading that method as of fleetd #330, not derived from this test.
|
||||
* {@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 = Set.of(
|
||||
"guard", "worktreeRoot", "worktreeGroup", "primary", "configReload",
|
||||
"leadHeartbeat", "lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs",
|
||||
"quarantineCooldownSeconds", "profiles");
|
||||
|
||||
/**
|
||||
* 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 fleet} — role pools, {@code charters} and {@code tabLabel} are read live through
|
||||
* the supplier on {@code CompositePeerLauncher} (class doc, Hot bullet). The one
|
||||
* documented exception, {@code fleet.leaders}, is read only at startup and genuinely needs
|
||||
* a restart — a real gap in {@code changedDeferredKeys}, but one the class doc already
|
||||
* carries and that fleetd #330 explicitly did not re-open (its "seven readers" fact-find
|
||||
* named {@code health}/{@code coordinator} as the complete set of split-shaped keys, not
|
||||
* {@code fleet}). Reported as a caveat, not fixed here.</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>
|
||||
*/
|
||||
private static final Set<String> HOT_EXCLUDED_TOP_LEVEL_KEYS =
|
||||
Set.of("placement", "fleet", "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", "fleet", "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");
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
@@ -940,6 +940,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");
|
||||
|
||||
@@ -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