Compare commits

...

8 Commits

Author SHA1 Message Date
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
Dai Ha b9c2cf69f4 Merge #307: a worker's real reply after an ask timeout completes its ticket instead of stranding
CI / build (push) Successful in 2m5s
CI / contract (push) Successful in 2m19s
2026-09-04 13:31:43 +07:00
Dai Ha 8beae50fe7 Merge #308: refuse spawns once the shutdown drain has started, and sweep stragglers
CI / contract (push) Successful in 1m49s
CI / build (push) Successful in 3m0s
2026-09-04 13:28:09 +07:00
Dai Ha b2a58cb966 Merge #309: clean up partial worktree state when git worktree add fails 2026-09-04 13:28:04 +07:00
Dai Ha f159ca7d27 #310: log when a reap is skipped because the record changed
CI / contract (push) Successful in 1m7s
CI / build (push) Successful in 1m46s
The compare-and-release declines silently. This race is unobservable by
construction, so a reaper that quietly stops reaping is the hardest kind of
behaviour to diagnose later. One debug line names the pane and the likely
cause.
2026-09-04 13:22:55 +07:00
Dai Ha 83f2aea60f #308: refuse a spawn once the shutdown drain has started, and sweep stragglers
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Successful in 1m36s
drainAll iterated a one-shot registry snapshot with nothing to refuse a new
fleet_spawn while the drain was still running (mcp.close() only runs 8 calls
after sessions.close() in the shutdown hook). A session registered in that
window was never visited by the drain loop: its pane kept running and its
worktree was never preserved, with the in-memory registry gone at exit.

Fix, both mechanisms as the issue asked for (neither alone is complete):

- SessionManager.acquire now checks a `draining` flag, flipped true at the
  very start of drainAll before the registry snapshot is even taken, and
  throws the new ShuttingDownException (invariant 3: fail loudly, say why).
  FleetMcp.spawn and FleetApp.spawnMember surface it as a clean error/503
  rather than an uncaught RuntimeException.
- The flag alone cannot close the whole race: a caller already past the
  check can still be mid-launcher.spawn() (a real herdr round trip) when
  drainAll snapshots the registry. drainAll now re-reads the registry once
  its main pass finishes and drains whatever straggler landed there too,
  bounded by the SAME whole-drain deadline (invariant 1: timeoutNanos stays
  a budget for the whole drain, never extended for a straggler).
- ReleaseCause.SHUTDOWN still preserves worktrees for both the initial pass
  and the sweep (invariant 2, unchanged release() path).

Tests (SessionManagerTest): a guard test proving acquire() throws once
drainAll has started, and a race test using a launcher double that blocks
the second spawn() and the first stop() call to force, deterministically,
the exact interleaving where a spawn passes the guard before drainAll flips
it and only registers after the initial snapshot — proving the post-loop
sweep catches it.

Shape check (SessionManager.java only, not fixed): reapIdle has the same
shape — a decision made from a roster() snapshot, then acted on via
release(s.paneId()) with no re-check of the session's current state.
2026-09-04 13:21:53 +07:00
Dai Ha 002329adb5 #309: clean partial worktrees after add failure
CI / contract (pull_request) Successful in 1m23s
CI / build (pull_request) Successful in 2m2s
2026-09-04 13:18:57 +07:00
Dai Ha a49671ceb9 #310: prevent idle reap from stopping delivered workers
CI / contract (pull_request) Successful in 1m19s
CI / build (pull_request) Successful in 2m1s
2026-09-04 13:17:51 +07:00
9 changed files with 738 additions and 21 deletions
@@ -22,6 +22,7 @@ import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.ShuttingDownException;
import dev.ltms.fleet.session.WorktreeRequest;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerLauncher;
@@ -962,6 +963,10 @@ public final class FleetMcp {
return text(json(memberView(member)));
} catch (GuardException e) {
return error("subscription boundary: " + e.getMessage());
} catch (ShuttingDownException e) {
// fleetd #308: the daemon's shutdown drain has already started — refuse loudly rather
// than register a session drainAll will never see again.
return error("shutting down: " + e.getMessage());
} catch (PlacementException e) {
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — distinct
// from "profile does not exist" below.
@@ -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)) {
@@ -19,6 +19,7 @@ import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.ShuttingDownException;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.WorktreeRequest;
@@ -501,6 +502,11 @@ public final class FleetApp {
ctx.status(201).json(view(member));
} catch (GuardException e) {
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
} catch (ShuttingDownException e) {
// fleetd #308: the daemon's shutdown drain has already started — 503, not a bare 500,
// so this reads the same as PlacementException below: valid request, refused because
// of a transient daemon state rather than a bad argument.
ctx.status(503).json(Map.of("error", "shutting_down", "detail", e.getMessage()));
} catch (PlacementException e) {
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — a benign,
// likely-transient refusal, distinct from "profile does not exist" below. 503: the
@@ -93,6 +93,9 @@ public final class GitWorktrees implements Worktrees {
/** OS group name for {@link #shareWithGroup} (fleetd #185 stage 3); {@code null} ⇒ feature off. */
private final String group;
private final Consumer<String> afterWorktreeAdded;
/** How the initial {@code git worktree add} command runs. Package-private test seam for an
* interrupted command after Git has made worktree state. */
private final Function<String[], String> worktreeAddRunner;
/** How {@link #shareWithGroup}'s processes (git config / chgrp / chmod / find) actually run.
* Defaults to the real {@link #exec(String...)}. Package-private test seam so a unit test can
* prove "no group configured ⇒ zero processes spawned" and inspect exactly what a configured
@@ -132,7 +135,13 @@ public final class GitWorktrees implements Worktrees {
/** Test seam combining a configurable {@code group} with {@link #afterWorktreeAdded}. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded) {
this(configuredRoot, group, afterWorktreeAdded, null);
this(configuredRoot, group, afterWorktreeAdded, null, null);
}
/** Test seam for changing how {@link #shareWithGroup}'s processes run. */
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner) {
this(configuredRoot, group, afterWorktreeAdded, shareGroupRunner, null);
}
/**
@@ -141,13 +150,15 @@ public final class GitWorktrees implements Worktrees {
* exactly what commands a configured group runs, without a real second OS user/group.
*
* @param shareGroupRunner {@code null} ⇒ the real {@link #exec(String...)}.
* @param worktreeAddRunner {@code null} ⇒ the real {@link #exec(String...)}.
*/
GitWorktrees(String configuredRoot, String group, Consumer<String> afterWorktreeAdded,
Function<String[], String> shareGroupRunner) {
Function<String[], String> shareGroupRunner, Function<String[], String> worktreeAddRunner) {
this.configuredRoot = configuredRoot;
this.group = (group == null || group.isBlank()) ? null : group;
this.afterWorktreeAdded = afterWorktreeAdded == null ? _ -> {} : afterWorktreeAdded;
this.shareGroupRunner = shareGroupRunner != null ? shareGroupRunner : this::exec;
this.worktreeAddRunner = worktreeAddRunner != null ? worktreeAddRunner : this::exec;
}
@Override
@@ -166,8 +177,8 @@ public final class GitWorktrees implements Worktrees {
String wt = path.toAbsolutePath().toString();
log.info("adding worktree branch={} path={} base={}", branch, wt, base);
removeUserInfoFromHttpsOrigin(repoRoot);
exec("git", "-C", repoRoot, "worktree", "add", wt, "-b", branch, base);
try {
worktreeAddRunner.apply(new String[] {"git", "-C", repoRoot, "worktree", "add", wt, "-b", branch, base});
afterWorktreeAdded.accept(wt);
requireCredentialFreeHttpsOrigin(wt);
configureEnvironmentCredentialHelper(repoRoot, wt);
@@ -181,10 +192,11 @@ public final class GitWorktrees implements Worktrees {
}
/**
* {@code add()} has already created the worktree and its branch by the time any step from
* {@link #afterWorktreeAdded} through {@link #isolateToolSurface} can throw — including
* {@code git worktree add} may have created the worktree and its branch by the time it, or any
* later step through {@link #isolateToolSurface}, throws. This includes
* {@link #requireCredentialFreeHttpsOrigin}, an intended security refusal, not only an IO
* accident. Without this, {@code add()} never returns, so its caller
* accident. A Git-reported {@code worktree add} failure usually creates nothing, but an
* interrupted command can leave partial state. Without cleanup, {@code add()} never returns, so its caller
* ({@code SessionManager#acquireWithWorktree}) never receives a path to register or clean up:
* its local {@code path} stays null, the {@code if (path != null)} guard in its own catch block
* never runs, and the worktree directory and branch leak on disk forever with nothing tracking
@@ -204,6 +216,10 @@ public final class GitWorktrees implements Worktrees {
* used in {@code SessionManager#acquireWithWorktree}'s own catch block.
*/
private void cleanupAfterAddFailure(String repoRoot, String worktreePath, String branch, RuntimeException original) {
if (!Files.exists(Path.of(worktreePath))) {
log.debug("provisioning failed before worktree {} existed; nothing to clean up", worktreePath);
return;
}
log.warn("provisioning failed for branch={} path={}: {} — cleaning up before rethrowing",
branch, worktreePath, original.getMessage());
try {
@@ -21,6 +21,7 @@ import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;
import java.util.function.LongSupplier;
@@ -61,6 +62,8 @@ public final class SessionManager implements TurnListener {
private final LongSupplier nowNanos;
private final int contextCap;
private final boolean clearAfterTurn;
/** Null in production; test seam for the interval before an idle session's conditional release. */
private final Consumer<MemberSession> beforeIdleRelease;
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
/**
* CB-586: the repo root the fleet actually works in, remembered the first time a worktree
@@ -71,6 +74,16 @@ public final class SessionManager implements TurnListener {
*/
private volatile String fleetRepoRoot;
/**
* fleetd #308: flips true the instant {@link #drainAll} starts, before its registry snapshot
* is even taken — so a spawn already in flight sees the refusal as early as a plain flag can
* make it. This alone cannot close the race completely: a caller that read {@code false} just
* before the flip can still land in the registry after the snapshot. {@link #drainAll}'s
* post-loop sweep is what catches that straggler; the two mechanisms are deliberately paired,
* see {@link #drainAll}'s javadoc.
*/
private final AtomicBoolean draining = new AtomicBoolean(false);
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
/** CB-516: notified with a {@link ReleaseDetail} on every release; no-op until wired. */
@@ -102,13 +115,23 @@ public final class SessionManager implements TurnListener {
}
public SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
int contextCap, boolean clearAfterTurn) {
int contextCap, boolean clearAfterTurn) {
this(launcher, worktrees, nowNanos, contextCap, clearAfterTurn, null);
}
/**
* Package-private constructor for a deterministic reap/delivery race test. Production callers
* use the constructor above, whose null hook adds no callback or lock to an ordinary reap.
*/
SessionManager(PeerLauncher launcher, Worktrees worktrees, LongSupplier nowNanos,
int contextCap, boolean clearAfterTurn, Consumer<MemberSession> beforeIdleRelease) {
this.launcher = launcher;
this.worktrees = worktrees;
this.presence = new PresenceFleet(this);
this.nowNanos = nowNanos;
this.contextCap = contextCap;
this.clearAfterTurn = clearAfterTurn;
this.beforeIdleRelease = beforeIdleRelease;
}
/**
@@ -181,6 +204,14 @@ public final class SessionManager implements TurnListener {
public MemberSession acquire(String profile, MemberRole role, String requestedCwd, String callerCwd,
String ownerTerminal, WorktreeRequest wt,
String sessionName, String resumeSessionId) {
// fleetd #308: refuse before anything else runs — no slot reservation, no launcher spawn —
// so a caller learns the daemon is going down instead of getting a session drainAll will
// never see again. Checked here because every other acquire(...) overload delegates to
// this one, so this is the single point every spawn path passes through.
if (draining.get()) {
throw new ShuttingDownException("fleetd is shutting down; refusing to spawn a session "
+ "the shutdown drain would never see");
}
MemberRole memberRole = (role == null) ? MemberRole.DEV : role;
requireResumeCapability(profile, resumeSessionId);
// CB-619 / fleetd #123: an explicit profile bypasses placement (CompositePeerLauncher only
@@ -272,10 +303,32 @@ public final class SessionManager implements TurnListener {
*/
private void release(String paneId, ReleaseCause cause) {
MemberSession removed = registry.remove(paneId);
releaseRemoved(paneId, removed, handles.remove(paneId), cause);
}
/**
* Tear a session down only while {@code expected} is still its registry value. A lifecycle
* transition replaces the immutable record, so this prevents a reap based on an old READY or
* DONE record from stopping a worker that delivery has made BUSY.
*/
private boolean releaseIfCurrent(MemberSession expected, ReleaseCause cause) {
if (!registry.remove(expected.paneId(), expected)) {
// A lifecycle transition replaced the record between the caller's check and this remove.
// Log it: this race is by definition unobservable otherwise, and a reaper that silently
// declines to reap is the hardest kind of behaviour to diagnose after the fact.
log.debug("skipping reap of pane={}: its registry record changed after the idle check "
+ "(most likely a delivery made it BUSY)", expected.paneId());
return false;
}
releaseRemoved(expected.paneId(), expected, handles.remove(expected.paneId()), cause);
return true;
}
private void releaseRemoved(String paneId, MemberSession removed, PeerHandle removedHandle,
ReleaseCause cause) {
// fleetd #209: remove right alongside the registry entry so a released session's handle is
// never leaked — but keep the local reference below, so the id can still be resolved for
// the ReleaseDetail this teardown notifies with.
PeerHandle removedHandle = handles.remove(paneId);
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
String snapshotRef = null;
if (removed != null) {
@@ -846,14 +899,18 @@ public final class SessionManager implements TurnListener {
}
long idleNanos = now - s.lastActivityAtNanos();
if (idleNanos > idleTtlNanos) {
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
// CB-581: one session that fails to release must not abort the whole reaping pass —
// match drainAll's per-session try/catch so the rest of the roster still gets reaped.
try {
release(s.paneId());
reaped++;
if (beforeIdleRelease != null) {
beforeIdleRelease.accept(s);
}
if (releaseIfCurrent(s, ReleaseCause.COMPLETED)) {
log.debug("reaping idle session terminal={} pane={}: idle {}s exceeds the {}s ttl",
s.terminalId(), s.paneId(), TimeUnit.NANOSECONDS.toSeconds(idleNanos),
TimeUnit.NANOSECONDS.toSeconds(idleTtlNanos));
reaped++;
}
} catch (RuntimeException e) {
log.warn("reap failed for pane={} terminal={} worktree={}; continuing with "
+ "remaining sessions", s.paneId(), s.terminalId(), s.worktree(), e);
@@ -884,10 +941,40 @@ public final class SessionManager implements TurnListener {
* a reason to delete a worker's only copy of its uncommitted work. A session still {@code BUSY}
* when the timeout expired is abandoned mid-turn and logged loudly so an operator can find its
* kept worktree.
*
* <p>fleetd #308: {@code roster()} is a one-shot snapshot (see its javadoc), and nothing used
* to stop a new session from registering after it was taken — {@link #acquire} stayed open for
* as long as this drain waited on a {@code BUSY} session, up to the whole {@code timeoutNanos}
* budget. Two things close that window, deliberately paired because neither alone is complete:
* {@link #draining} is flipped true before the snapshot is even taken, so {@link #acquire}
* refuses (invariant 3: loudly, via {@link ShuttingDownException}) as much of the window as a
* plain flag can close; and the sweep below re-reads the registry once the initial snapshot has
* fully drained and drains whatever a straggler — a caller that read the flag as {@code false}
* a moment before it flipped — still managed to register. The sweep shares the same
* {@code deadline} rather than getting its own: {@code timeoutNanos} is a budget for the WHOLE
* drain (see above), and a straggler must not buy the drain more time than the flag it lost the
* race against would have. In the ordinary case the sweep finds nothing and costs one empty
* {@link #roster()} call.
*/
void drainAll(long timeoutNanos) {
long deadline = System.nanoTime() + timeoutNanos;
for (MemberSession s : roster()) {
draining.set(true);
drainSnapshot(roster(), deadline);
List<MemberSession> stragglers = roster();
if (!stragglers.isEmpty()) {
log.warn("drain sweep found {} session(s) registered after the drain snapshot was "
+ "taken (raced past the shutdown guard); draining them too", stragglers.size());
drainSnapshot(stragglers, deadline);
}
}
/**
* Drain exactly the sessions in {@code snapshot}, waiting out a {@code BUSY} one against the
* shared whole-drain {@code deadline} before releasing it. Shared by {@link #drainAll}'s main
* pass and its post-loop straggler sweep (fleetd #308) so both honor the same one budget.
*/
private void drainSnapshot(List<MemberSession> snapshot, long deadline) {
for (MemberSession s : snapshot) {
try {
if (s.state() == MemberSession.State.BUSY) {
while (System.nanoTime() < deadline) {
@@ -0,0 +1,18 @@
package dev.ltms.fleet.session;
/**
* Thrown by {@link SessionManager#acquire} when a spawn is requested after the daemon's shutdown
* drain has already begun (fleetd #308).
*
* <p>{@link SessionManager#drainAll} snapshots the registry once and tears down exactly what is
* in that snapshot. A session registered after the snapshot is invisible to the drain loop: its
* pane is left running and its worktree is never preserved, and nothing else ever reclaims
* either — the daemon's in-memory registry dies with the process. Refusing the spawn here,
* loudly, is what stops that session from ever being created in the first place, rather than
* silently handing the caller a session the daemon can no longer manage.
*/
public final class ShuttingDownException extends RuntimeException {
public ShuttingDownException(String message) {
super(message);
}
}
@@ -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;
}
}
@@ -403,6 +403,52 @@ class GitWorktreesTest {
assertTrue(heads.isBlank(), "the branch leaked after a post-creation step threw:\n" + heads);
}
/**
* fleetd #309. The add runner is a narrow seam for the case where Git has created state but
* the caller then kills the process. The runner first performs the real add in this throwaway
* repo, then throws the same kind of exception that {@link GitWorktrees#exec} uses for a timeout.
* This proves the failure path cleans both real Git objects without waiting for a slow checkout.
*/
@Test
void addCleansUpWhenTheWorktreeAddRunnerFailsAfterCreatingState(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
String branch = "cb-309-timeout";
WorktreeException timeout = new WorktreeException("command timed out: synthetic git worktree add");
AtomicReference<String> createdPath = new AtomicReference<>();
GitWorktrees worktrees = new GitWorktrees(tmp.resolve("wts").toString(), null, null, null, command -> {
createdPath.set(command[5]);
try {
git(repo, "worktree", "add", command[5], "-b", command[7], command[8]);
} catch (Exception e) {
throw new AssertionError("test setup could not create the worktree", e);
}
throw timeout;
});
WorktreeException thrown = assertThrows(WorktreeException.class,
() -> worktrees.add(repo.toString(), branch, "HEAD"));
assertSame(timeout, thrown, "cleanup must not replace the add failure");
assertNotNull(createdPath.get(), "the add runner must receive the worktree path");
assertFalse(Files.exists(Path.of(createdPath.get())),
"the worktree directory leaked after the add runner failed");
assertFalse(refExists(repo, "refs/heads/" + branch), "the branch leaked after the add runner failed");
}
/** An ordinary Git refusal must not delete the existing branch or log a cleanup warning. */
@Test
void addFailureBeforeCreatingAWorktreeIsQuiet(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
String branch = "already-exists";
git(repo, "branch", branch);
assertThrows(WorktreeException.class, () -> new GitWorktrees(tmp.resolve("wts").toString())
.add(repo.toString(), branch, "HEAD"));
assertTrue(refExists(repo, "refs/heads/" + branch), "the existing branch must remain");
assertTrue(capturedMessages().isEmpty(), "an ordinary Git refusal logged a warning: " + capturedMessages());
}
// ---- CB-189: broader remote-URL coverage — every remote, both fetch and push URLs, any
// non-SSH scheme. Reporting only, additive to the origin/https strip-and-refuse tests above. ----
@@ -26,7 +26,12 @@ import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.*;
@@ -178,14 +183,19 @@ class SessionManagerTest {
}
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
boolean clearAfterTurn) {
boolean clearAfterTurn) {
return sessionManager(herdr, clock, contextCap, clearAfterTurn, null);
}
private SessionManager sessionManager(FakeHerdr herdr, LongSupplier clock, int contextCap,
boolean clearAfterTurn, java.util.function.Consumer<MemberSession> hook) {
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher workers = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn);
return new SessionManager(workers, new GitWorktrees(), clock, contextCap, clearAfterTurn, hook);
}
@Test
@@ -670,6 +680,24 @@ class SessionManagerTest {
"BUSY session remains");
}
@Test
void reapIdleDoesNotReleaseSessionDeliveredAfterItsEligibilityCheck() {
long[] clock = {0};
FakeHerdr herdr = new FakeHerdr();
SessionManager[] manager = new SessionManager[1];
SessionManager sessions = sessionManager(herdr, () -> clock[0], 0, false,
session -> manager[0].onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId())));
manager[0] = sessions;
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
clock[0] = 11;
assertEquals(0, sessions.reapIdle(10), "delivery replaces the idle snapshot before release");
assertEquals(MemberSession.State.BUSY, sessions.get(session.paneId()).orElseThrow().state(),
"a just-delivered session stays registered and busy");
assertFalse(herdr.called("pane.close"), "the busy session pane is not stopped");
}
@Test
void doneSessionPastIdleTtlIsReaped() {
long[] clock = {0};
@@ -824,6 +852,186 @@ class SessionManagerTest {
.count();
}
// --- fleetd #308: a spawn accepted while the shutdown drain is running must not orphan ---
@Test
void acquireRefusesANewSpawnOnceDrainAllHasStarted() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(50)); // empty roster — returns immediately,
// but the shutdown guard it flips must stay tripped for the life of the process.
ShuttingDownException e = assertThrows(ShuttingDownException.class,
() -> sessions.acquire("ltms-local", null, "/caller", "term_primary"),
"a spawn requested after the drain has begun must be refused loudly (invariant 3), "
+ "not silently registered into a registry the drain will never revisit");
assertNotNull(e.getMessage());
assertFalse(e.getMessage().isBlank(), "the refusal must say why, not just that it failed");
assertTrue(sessions.roster().isEmpty(), "the refused spawn must never reach the registry");
}
/**
* fleetd #308: the guard above closes most of the shutdown-race window, but it cannot close
* all of it — a caller that already passed the {@code draining} check before {@code drainAll}
* flips it can still be mid-{@code launcher.spawn()} (a real herdr round trip, not
* instantaneous) when {@code drainAll} takes its registry snapshot. This test forces exactly
* that interleaving with a launcher double that blocks the second {@code spawn()} call and the
* first {@code stop()} call until released, then proves the post-loop sweep in {@code
* drainAll} still finds and tears down the straggler that lands in the registry afterward.
*/
@Test
void drainAllSweepsAStragglerThatRegisteredAfterTheInitialSnapshot() throws Exception {
FakeHerdr herdr = new FakeHerdr();
FleetConfig.Profile cfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "FLEETD_WORKER_TOKEN",
List.of("ccs", "ltms-local"), "tab", "fleetd-workers",
"worker: {profile} #{n}", null, null, null);
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
RaceLauncher race = new RaceLauncher(delegate);
SessionManager sessions = new SessionManager(race);
// Registered normally, before the drain starts — the first spawn call, never blocked.
MemberSession ready = sessions.acquire("ltms-local", "/ready", "/caller", "ownerR");
ExecutorService exec = Executors.newFixedThreadPool(2);
try {
// The straggler's acquire() reads `draining == false` (checked before this call ever
// touches the launcher) and then blocks inside its own spawn() — the second spawn call.
Future<MemberSession> straggler = exec.submit(() ->
sessions.acquire("ltms-local", "/late", "/caller", "ownerLate"));
assertTrue(race.enteredSecondSpawn.await(5, TimeUnit.SECONDS),
"the straggler must have passed the shutdown guard and reached spawn() before "
+ "drainAll ever runs");
assertEquals(1, sessions.roster().size(),
"the straggler is still inside spawn() — not registered yet");
// drainAll flips `draining`, snapshots the registry (only `ready` is in it), and starts
// releasing that snapshot — its first release() call stops `ready`'s pane, which this
// launcher double blocks on so the interleaving below is deterministic, not a timing bet.
Future<?> drain = exec.submit(() -> sessions.drainAll(TimeUnit.SECONDS.toNanos(5)));
assertTrue(race.enteredFirstStop.await(5, TimeUnit.SECONDS),
"drainAll must be stopping the ready session's pane — proof its initial "
+ "registry snapshot has already been taken");
// Only now does the straggler's spawn complete and register — strictly after the
// snapshot drainAll's main pass is working from.
race.releaseSecondSpawn.countDown();
MemberSession registered = straggler.get(5, TimeUnit.SECONDS);
// Let drainAll finish releasing `ready`; it then re-checks the registry and must find
// (and drain) the straggler that just landed in it.
race.releaseFirstStop.countDown();
drain.get(5, TimeUnit.SECONDS);
assertTrue(sessions.roster().isEmpty(),
"the post-loop sweep must drain the straggler too, not just the initial snapshot");
assertNotNull(registered.paneId());
long paneCloseCalls = herdr.calls.stream().filter(c -> "pane.close".equals(c.method())).count();
assertEquals(2, paneCloseCalls,
"both ready's pane AND the straggler's pane must actually be stopped — a pane "
+ "left running is exactly the orphan this ticket is about");
} finally {
exec.shutdownNow();
}
}
/**
* Delegates every call while blocking the SECOND {@code spawn()} call and the FIRST
* {@code stop()} call until the test releases them — used to force the fleetd #308 race
* deterministically instead of betting on real thread-scheduling timing.
*/
private static final class RaceLauncher implements PeerLauncher {
private final PeerLauncher delegate;
private final AtomicInteger spawnCalls = new AtomicInteger();
private final AtomicInteger stopCalls = new AtomicInteger();
final CountDownLatch enteredSecondSpawn = new CountDownLatch(1);
final CountDownLatch releaseSecondSpawn = new CountDownLatch(1);
final CountDownLatch enteredFirstStop = new CountDownLatch(1);
final CountDownLatch releaseFirstStop = new CountDownLatch(1);
RaceLauncher(PeerLauncher delegate) {
this.delegate = delegate;
}
private static void awaitOrFail(CountDownLatch latch) {
try {
if (!latch.await(5, TimeUnit.SECONDS)) {
throw new AssertionError("RaceLauncher latch timed out");
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new AssertionError("RaceLauncher latch interrupted", e);
}
}
@Override
public Set<Capability> capabilities() {
return delegate.capabilities();
}
@Override
public Set<Capability> capabilitiesFor(String profileName) {
return delegate.capabilitiesFor(profileName);
}
@Override
public PeerHandle spawn(SpawnRequest req) {
if (spawnCalls.incrementAndGet() == 2) {
enteredSecondSpawn.countDown();
awaitOrFail(releaseSecondSpawn);
}
return delegate.spawn(req);
}
@Override
public Set<String> profiles() {
return delegate.profiles();
}
@Override
public String defaultProfile() {
return delegate.defaultProfile();
}
@Override
public String effectiveCwd(SpawnRequest req) {
return delegate.effectiveCwd(req);
}
@Override
public List<String> parityOverlay(String profileName) {
return delegate.parityOverlay(profileName);
}
@Override
public List<?> list() {
return delegate.list();
}
@Override
public int reapOrphanWorkers() {
return delegate.reapOrphanWorkers();
}
@Override
public void stop(String id) {
if (stopCalls.incrementAndGet() == 1) {
enteredFirstStop.countDown();
awaitOrFail(releaseFirstStop);
}
delegate.stop(id);
}
@Override
public boolean clearContext(String id) {
return delegate.clearContext(id);
}
}
private static List<String> promptTexts(FakeHerdr herdr) {
return herdr.calls.stream()
.filter(c -> "agent.prompt".equals(c.method()))