Compare commits

...

8 Commits

Author SHA1 Message Date
Dai Ha 6d493bc7bb fleetd#326: classify primary and configReload as deferred top-level keys
CI / contract (pull_request) Successful in 1m3s
CI / build (pull_request) Failing after 1m44s
ConfigRef.changedDeferredKeys only classified seven top-level FleetConfig
keys (#323 fixed the profile side). Two more keys are read only off the
startup snapshot and were missing:

- primary: Fleetd.java:506/519/520 feed PrimaryRegistry and ReplyPushLoop
  at construction; neither is rebuilt on reload.
- configReload: Fleetd.java:679-680 decide once at startup whether to
  build a ConfigWatcher at all, and with what interval; the watcher that
  would apply a later change is itself built once, so it is deferred
  (not cold — no already-open resource goes inconsistent, a running
  watcher just keeps its original settings).

health and coordinator are deliberately left unclassified: both are read
both off the startup snapshot AND live off the config supplier at a
second call site, so no single bucket is correct for either — see the
PR body for the options writeup and the coordinator.uriEnv exposure
question the issue asked to be answered.

Each fix is proven with a failing-first test in ConfigRefTest and a
revert-quote-restore mutation check (see PR body for the transcripts).
2026-09-04 14:42:24 +07:00
Dai Ha e545c08082 #323: pin the exclusion set — the coverage mechanism's own escape hatch, found by mutation
CI / contract (push) Successful in 1m29s
CI / build (push) Successful in 2m10s
2026-09-04 14:32:08 +07:00
Dai Ha efa0deb9b2 Merge #323: the reload classifier proves its own coverage instead of claiming it 2026-09-04 14:29:10 +07:00
Dai Ha fa1f49675b Merge #318: a delivery landing after release is refused, not parked in a map nobody reads
CI / contract (push) Successful in 1m0s
CI / build (push) Successful in 2m8s
2026-09-04 14:21:29 +07:00
Dai Ha 8426c3528f #316: pin the fail-toward-preserve rule on the late re-check, found by mutation
CI / contract (push) Successful in 1m22s
CI / build (push) Successful in 1m42s
2026-09-04 14:20:36 +07:00
Dai Ha 65f98ba910 Merge #316: the dirty check that authorises the worktree removal is taken after the worker stops 2026-09-04 14:17:13 +07:00
Dai Ha 667254df47 #316: re-check worktree dirtiness after the pane stops, before removing it
CI / contract (pull_request) Successful in 39s
CI / build (pull_request) Successful in 1m46s
SessionManager.releaseRemoved read hasUncommitted() once, while the worker
could still write, then used that stale boolean after launcher.stop() to
authorise `git worktree remove --force`. The same stale read also gated
trySnapshot, so a worker that wrote between the read and the stop lost its
work with neither a preserve nor a snapshot.

Add a second, best-effort hasUncommitted read immediately before the
removal, taken only on the path that is actually about to delete something
(never on a release that already decided to preserve, and never for
SHUTDOWN, which preserves unconditionally). If the tree is now dirty,
preserve it and attempt a fresh snapshot, since the original snapshot never
ran when the pre-stop read said clean. A failing re-check also preserves,
matching the existing CB-581 fail-safe rule.
2026-09-04 14:12:49 +07:00
Dai Ha 2926cd1784 #318: release() no longer strands a delivery that lands while it is running
CI / build (pull_request) Failing after 1m59s
CI / contract (pull_request) Successful in 2m14s
AmqpReplyInbox.release() used held.remove(target) then iterated the old
map. A delivery landing on the consumer work-pool thread after the
remove (basicCancel does not flush one already handed to that pool) hit
deliverCallback's computeIfAbsent, found the key gone, and created a
brand-new map release() never looks at again — delivered-but-unacked
forever, never requeued, never redelivered (#298 only closed the
"already in held when release runs" case).

Fix: release() swaps in a RELEASED tombstone via held.compute(...)
instead of held.remove(...). ConcurrentHashMap serializes compute/
computeIfAbsent calls for the same key against each other, so whichever
of release() and a concurrent deliverCallback runs first is fully
visible to the other — no gap. deliverCallback checks for the
tombstone and nacks-with-requeue instead of recreating a map; peek/ack
treat it as empty; own() clears a stale tombstone so a target is never
poisoned if its id is ever reused (the issue's own text says id reuse
doesn't happen, but the tombstone would otherwise sit in `held` forever
either way).

New test AmqpReplyInboxReleaseRaceTest forces the actual interleaving
with a latch (blocks release() inside its nack loop, which is only
reachable after the tombstone swap, then fires a concurrent delivery)
rather than a sequential call — a sequential test would not have caught
this, since #298's own contract test forces settlement before release()
runs. Mutation-tested: reverting the fix makes this test fail with
"expected: <2> but was: <1>" (m1 never nacked); restored after
confirming that failure.
2026-09-04 14:10:21 +07:00
7 changed files with 621 additions and 6 deletions
@@ -42,7 +42,17 @@ 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; a lead whose pinned terminal changed under a running daemon stays
* unresolved as primary until a restart), {@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
@@ -234,6 +244,20 @@ 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 pin needs a restart before a lead resolves as primary again.
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*");
@@ -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)) {
@@ -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,70 @@ 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());
}
@Test
void aFixedRefHasNoFileAndRefusesToReload() {
FleetConfig cfg = new FleetConfig(null, null, null, null, null, null,
@@ -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;
}
}
@@ -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