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 b8b25cf74c #307: an ask() timeout no longer strands the worker's real reply
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 1m30s
MessageService.reply()'s async-recovery path (askAnsweredAsyncTasks)
required a live Task.turnId, but ask()'s own TimeoutException handler
calls clearAsyncQuestion(turnId, true) — deliberately forgetting turnId
so hasAsyncQuestion() stops reporting the target BUSY. That made a
worker's eventual real fleet_reply, after an unanswered fleet_ask, fall
through to the inbox: fleet_poll{ticket} stayed PENDING forever and was
later force-failed with the false reason "session released before it
replied".

Fix: a new Task.askTimedOut marker is set (markAskTimedOut) right
before the turnId is forgotten, and askAnsweredAsyncTasks accepts it in
place of a live turnId. The marker never touches asyncTasksByTurn, so
the BUSY-release behaviour (invariant 1) is untouched. The existing
ambiguity guard (candidates.size() > 1 -> inbox, never guess) still
applies unchanged, but is now genuinely reachable rather than pure
defence in depth, since an ask timeout frees its target for a fresh,
independent delegation — the affected javadocs are updated to say so.

Tests: MessageServiceTest.aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket
(positive, mutation-proven) and
.twoAskTimedOutTicketsOnOneTargetFallBackToTheInboxRatherThanGuess
(negative/ambiguity). FleetMcpTest's
unansweredAsyncAskReturnsTheTicketToPending was renamed and its final
assertion updated — it had pinned the old (buggy) inbox-stranding
behaviour as expected.
2026-09-04 13:23:01 +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 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 647 additions and 66 deletions
@@ -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)) {
@@ -187,6 +187,24 @@ public final class MessageService {
private volatile Long completedNanos;
private volatile Reply question;
private volatile String turnId;
/**
* Set when this task's {@code fleet_ask} lapsed with no answer (fleetd #307):
* {@link #clearAsyncQuestion} then forgets {@link #turnId} (nulls it and drops the task from
* {@code asyncTasksByTurn}) so {@link #hasAsyncQuestion} stops reporting the target BUSY — a
* later {@code fleet_send} to it must be accepted, not refused. But the worker's turn is
* still genuinely live: it resumed on its own and will eventually call its real
* {@code fleet_reply}. Losing {@link #turnId} loses {@link #askAnsweredAsyncTasks}' only
* signal that such a reply belongs to this task, so that reply used to fall straight to the
* inbox and strand — {@code fleet_poll} stayed {@code PENDING} forever, later force-failed by
* {@link #abandon} with the misleading "session released before it replied". This flag is a
* second, independent signal that survives the forgetting: {@link #askAnsweredAsyncTasks}
* accepts it in place of a live {@link #turnId}, without ever re-adding the task to
* {@code asyncTasksByTurn} (so the BUSY release is untouched). Cleared implicitly once
* {@link #future} resolves — every match in {@link #askAnsweredAsyncTasks} already requires
* {@code !future.isDone()}, so a task that recovered (or was later failed by
* {@link #abandon}) can never match again regardless of this flag's value.
*/
private volatile boolean askTimedOut;
private Task(String ticket, String target, LongSupplier nowNanos) {
this.ticket = ticket;
@@ -395,18 +413,16 @@ public final class MessageService {
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
* are interactive and must never be queued.
*
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks}
* cannot actually return more than one entry today (see its own javadoc for why — in short,
* {@link #hasAsyncQuestion} keeps a target BUSY, so no second task can reach this state, for as
* long as an earlier one's {@code turnId} is still stamped). That is an emergent guarantee from
* two other facts, not one this method enforces, so this branch stays in as defence in depth
* rather than being removed as dead code: if it ever weakens, returning whichever candidate a
* {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the
* <em>wrong</em> ticket — silently handing the lead something that reads like a correct answer to
* a delegation the worker never touched, which is worse than a failure because the lead acts on
* it. When more than one candidate exists, guessing is not safe: fall back to the inbox exactly
* as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery
* deterministically instead.
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks} can
* return more than one entry — a reachable state, not a hypothetical one (see its own javadoc:
* an {@code fleet_ask} that lapsed with no answer, fleetd #307, frees the target for a completely fresh
* delegation, which can itself go on to ask-and-lapse before the first worker's real reply
* arrives). Returning whichever candidate a {@code ConcurrentHashMap} iteration reaches first
* would let a genuine reply complete the <em>wrong</em> ticket — silently handing the lead
* something that reads like a correct answer to a delegation the worker never touched, which is
* worse than a failure because the lead acts on it. When more than one candidate exists, guessing
* is not safe: fall back to the inbox exactly as the zero-candidate case does, and let
* {@link #abandon} apply the eventual recovery deterministically instead.
*
* <p><strong>{@code content} is required (fleetd #302).</strong> Both doors that reach this
* method must reject a missing/blank reply the same way, so the check lives here rather than in
@@ -434,15 +450,18 @@ public final class MessageService {
count(FleetMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path
}
// #137: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming a
// turn that {@link #answer} already gave up waiting on. answer()'s own bounded wait (the
// primary's fleet_send{turnId} call, capped well under a minute) can time out and close its
// waiter long before the worker — now actually resuming real work — finishes and replies. That
// reply used to have nowhere to land but the session inbox, leaving the async ticket's future
// unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's abandon() forced it
// FAILED with a misleading "session released before it replied" reason, even though the reply
// had, in fact, arrived. Completing the matching ticket directly here means fleet_poll{ticket}
// sees the real reply instead.
// #137/fleetd #307: no live rendezvous waiter, but this may be the worker's real fleet_reply resuming
// a turn that either answer() (#137) or ask() (fleetd #307) already gave up waiting on:
// - answer()'s own bounded wait (the primary's fleet_send{turnId} call, capped well under a
// minute) can time out and close its waiter long before the worker — now actually resuming
// real work — finishes and replies.
// - ask()'s own wait for the primary can time out first, with the worker resuming on its own
// and finishing unanswered.
// Either way that reply used to have nowhere to land but the session inbox, leaving the async
// ticket's future unresolved forever: fleet_poll{ticket} stayed PENDING until fleet_stop's
// abandon() forced it FAILED with a misleading "session released before it replied" reason,
// even though the reply had, in fact, arrived. Completing the matching ticket directly here
// means fleet_poll{ticket} sees the real reply instead.
List<Task> candidates = askAnsweredAsyncTasks(session);
if (candidates.size() == 1) {
Task orphan = candidates.get(0);
@@ -473,34 +492,42 @@ public final class MessageService {
}
/**
* Every still-open async task on {@code target} whose {@code fleet_ask} was already answered —
* its {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer}
* — yet whose future is not resolved yet (#137). Empty if no such task exists, including the
* common case where {@code target}'s worker never used {@code fleet_ask} at all (a task that was
* never asked has {@code turnId == null}, so it can never match here and only ever completes
* through the ordinary rendezvous fast path in {@link #reply}).
* Every still-open async task on {@code target} whose worker is genuinely expected to send a
* real {@code fleet_reply} next with nothing left registered to catch it: either its
* {@code fleet_ask} was already answered — {@link Task#turnId} is stamped but {@link
* Task#question} was cleared by {@link #answer} — or its {@code fleet_ask} lapsed unanswered and
* {@link Task#askTimedOut} marks that (fleetd #307; {@link Task#turnId} is {@code null} by then, forgotten
* so the target is not left BUSY — see {@link Task#askTimedOut}'s own javadoc). Either way the
* task's future is not resolved yet. Empty if no such task exists, including the common case
* where {@code target}'s worker never used {@code fleet_ask} at all (a task that was never asked
* has both {@code turnId == null} and {@code askTimedOut == false}, so it can never match here and
* only ever completes through the ordinary rendezvous fast path in {@link #reply}).
*
* <p><strong>Returns at most one entry today — verified, not assumed.</strong> {@link #send}
* refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is true, and that
* check matches ANY task whose {@code turnId} is still stamped in {@code asyncTasksByTurn} —
* not only while its question is still open. {@link #answer} deliberately leaves that stamp in
* place ({@code clearAsyncQuestion(turnId, false)}) until the resumed turn's own future actually
* resolves, at which point {@link #finishAsyncTask} both removes the stamp AND completes that
* task's future in the same call. So a second task can never reach "{@code turnId} stamped, future
* still open" — the exact pair this method matches on — while a first one already holds it: by
* the time the stamp is gone, so is the eligibility. This is an emergent property of those two
* facts holding together, not something this method (or its callers) enforces on its own — flip
* {@code forgetTurn} to {@code true} in that one {@link #answer} call and it silently stops being
* true, with nothing left to fail loudly. The callers below still handle "more than one" as
* defence in depth against exactly that, not because they exercise it today: {@link #reply}
* treats it as unresolvable and falls back to the inbox; {@link #abandon} would pick the oldest
* deterministically (its own {@code matching} list has no such guarantee — see its javadoc).
* <p><strong>Can return more than one entry — reachable, not just defence in depth.</strong>
* {@link #send} refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is
* true, and that check matches ANY task whose {@code turnId} is still stamped in
* {@code asyncTasksByTurn}. While a task's {@code turnId} stays stamped — {@link #answer} leaves
* it in place ({@code clearAsyncQuestion(turnId, false)}) until {@link #finishAsyncTask} removes
* the stamp and completes the future in the same call — no second task on the same target can
* reach an eligible state, because {@link #send} would refuse it as BUSY first. That single-task
* guarantee holds only for the {@code turnId}-stamped half of this method's match: an
* {@link Task#askTimedOut} task is, by construction, no longer stamped in {@code asyncTasksByTurn}
* (that is the whole point of forgetting {@code turnId} in {@link #clearAsyncQuestion}), so the
* target is free the moment one ask lapses. A fresh, independent {@code sendAsync} to the same
* target can then be dispatched, itself pause on {@code fleet_ask}, and itself time out — landing
* a second {@code askTimedOut} task on the very target the first one is still waiting to answer
* for. Two (or more) genuinely open tasks on one target is therefore a real, reachable state
* today, not a hypothetical: {@link #reply} treats it as unresolvable and falls back to the
* inbox rather than guess which task a reply belongs to (guessing wrong would hand the lead a
* plausible-looking answer to a delegation the worker never touched — worse than a failure,
* because the lead acts on it); {@link #abandon} instead picks the oldest deterministically (its
* own {@code matching} list has a different, wider match — see its javadoc).
*/
private List<Task> askAnsweredAsyncTasks(String target) {
List<Task> candidates = new ArrayList<>();
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null && task.turnId != null
&& !task.future.isDone()) {
if (target.equals(task.target) && task.question == null && !task.future.isDone()
&& (task.turnId != null || task.askTimedOut)) {
candidates.add(task);
}
}
@@ -898,6 +925,12 @@ public final class MessageService {
return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) {
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it out of
// asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and stays (it is
// what keeps the target from staying BUSY forever), but it would otherwise also erase
// askAnsweredAsyncTasks' only signal that the worker's eventual real fleet_reply still
// belongs to this task, stranding it in the inbox with a false "never replied" verdict.
markAskTimedOut(ticket.turnId());
clearAsyncQuestion(ticket.turnId(), true);
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
@@ -1162,6 +1195,23 @@ public final class MessageService {
return task;
}
/**
* Mark {@code turnId}'s task as having a {@code fleet_ask} that lapsed with no answer (fleetd #307), so
* {@link #askAnsweredAsyncTasks} still recognizes the worker's eventual real {@code fleet_reply}
* as belonging to it after {@link #clearAsyncQuestion}'s {@code forgetTurn=true} erases
* {@link Task#turnId} — see {@link Task#askTimedOut}. Must be called before that forgetting, while
* {@code turnId} can still resolve the task in {@code asyncTasksByTurn}; a lookup afterward would
* find nothing. Only when it matches the task's current turn — same guard as
* {@link #clearAsyncQuestion} — so a chained second {@code fleet_ask} (#282) that already moved
* the task to a fresh {@code turnId} cannot mark it for a turn that is no longer its own.
*/
private void markAskTimedOut(String turnId) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null && turnId.equals(task.turnId)) {
task.askTimedOut = true;
}
}
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
// CB-582: tell the push loop first — like ticketCollected, a removal for a turnId it never
@@ -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 {
@@ -62,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
@@ -113,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;
}
/**
@@ -291,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) {
@@ -865,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);
@@ -189,7 +189,7 @@ class FleetMcpTest {
}
@Test
void unansweredAsyncAskReturnsTheTicketToPending() throws Exception {
void unansweredAsyncAskReturnsTheTicketToPendingThenAWorkersLateReplyStillCompletesIt() throws Exception {
McpSchema.CallToolResult accepted = FleetMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
@@ -203,8 +203,15 @@ class FleetMcpTest {
assertTrue(textOf(ask).contains("no answer"), textOf(ask));
assertTrue(textOf(FleetMcp.poll(messages, ticket, null)).startsWith("[pending"));
// fleetd #307: the worker resumed on its own after the primary never answered, and its real
// fleet_reply must complete its OWN async ticket — not strand in the inbox with
// fleet_poll{ticket} stuck PENDING forever and later force-failed with a false "session
// released before it replied" reason. This used to land in the inbox instead (see the old
// assertion this replaced: messages.drainReplies("term_a").getFirst()...) — that was the bug.
FleetMcp.reply(messages, "term_a", "finished after timeout");
assertEquals("finished after timeout", messages.drainReplies("term_a").getFirst().content());
assertEquals("finished after timeout", textOf(FleetMcp.poll(messages, ticket, null)));
assertTrue(messages.drainReplies("term_a").isEmpty(),
"the reply completed its own ticket directly and never touched the inbox");
}
@Test
@@ -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;
}
}
@@ -957,6 +957,76 @@ class MessageServiceTest {
assertEquals(MessageService.Phase.DONE, awaitTicketPhase(next, MessageService.Phase.DONE).phase());
}
/**
* fleetd #307: a worker's {@code fleet_ask} can time out because the primary never answers —
* distinct from {@link #aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket}, where the
* primary DID answer and only its own bounded wait for the resumed turn expired.
* {@code ask()}'s timeout path deliberately forgets the task's {@code turnId} (so
* {@code hasAsyncQuestion} stops reporting the target BUSY — see
* {@code unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget} above), which used
* to also erase the one signal {@code askAnsweredAsyncTasks} needed to recognize the worker's
* eventual real {@code fleet_reply}. That reply then had nowhere to land but the inbox, and
* {@code fleet_poll{ticket}} stayed PENDING forever — later force-failed with the false reason
* "session released before it replied", even though the worker had, in fact, replied.
*/
@Test
void aReplyAfterAnAskTimeoutStillCompletesTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks then finishes alone");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT,
messages.ask(T, "which config?", 200).outcome());
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
"only the question wait ended; the delegated turn may still finish");
// The worker keeps working past the timeout and only now calls fleet_reply — with no live
// rendezvous waiter open (ask()'s timeout already closed it) and no new send() having
// reopened one for this target.
assertTrue(messages.reply(T, "PR opened: https://example/pulls/42"));
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("PR opened: https://example/pulls/42", done.reply(),
"fleet_poll{ticket} must return the worker's real reply, not stay pending forever");
assertEquals("reply", done.replySource());
assertFalse(messages.hasStrandedReply(T),
"the reply completed its own ticket directly and never touched the inbox");
}
/**
* fleetd #307's ambiguity guard: an ask timeout frees its target ({@code hasAsyncQuestion}
* becomes false the instant it lapses — proven above), so a second, independent delegation can
* be dispatched to the same target and itself go on to ask-and-lapse before the first worker's
* real reply ever arrives. Two open tasks are then both eligible candidates on one target with
* no live waiter to disambiguate them. A reply arriving now must not guess which one it answers
* — guessing wrong would hand the lead a plausible-looking answer to a delegation the worker
* never touched, worse than a failure because the lead acts on it — so it must fall back to the
* inbox exactly as the zero-candidate case does.
*/
@Test
void twoAskTimedOutTicketsOnOneTargetFallBackToTheInboxRatherThanGuess() throws Exception {
String ticket1 = messages.sendAsync(T, "first task that asks");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q1?", 200).outcome());
String ticket2 = messages.sendAsync(T, "second task that asks");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT, messages.ask(T, "Q2?", 200).outcome());
assertTrue(messages.reply(T, "which task does this answer?"));
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket1).phase(),
"an ambiguous reply must not guess ticket1");
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket2).phase(),
"an ambiguous reply must not guess ticket2");
assertTrue(messages.hasStrandedReply(T));
var drained = messages.drainReplies(T);
assertEquals(1, drained.size());
assertEquals("which task does this answer?", drained.get(0).content());
}
@Test
void asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter() throws Exception {
String first = messages.sendAsync(T, "first task");
@@ -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. ----
@@ -183,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
@@ -675,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};