Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 6ed70700a0 #282: don't let a chained fleet_ask kill its own async ticket
CI / contract (pull_request) Successful in 1m6s
CI / build (pull_request) Successful in 1m38s
answer() opened a fresh forward waiter but, unlike send(), never
registered it in asyncTasksByWaiter. So when a worker chained a
second fleet_ask inside the same resumed turn (before calling
fleet_reply), markAsyncQuestion had no Task to re-associate, and
answer() then completed the async ticket's future with the second
QUESTION as if it were a terminal reply — fleet_poll reported FAILED
while the worker was still alive and mid-conversation.

Fix: register answer()'s waiter in asyncTasksByWaiter (mirroring
send()) so a chained ask can re-arm the ticket under its new turnId,
and guard answer()'s finishAsyncTask call the same way sendAsync's
own lambda already does (skip on Outcome.QUESTION). Also drop the
stale asyncTasksByTurn entry left behind when markAsyncQuestion
re-arms a task under a new turnId, a leak the fix makes reachable
for the first time.

Reachability confirmed by driving the exact sequence through the
public API (sendAsync -> ask -> answer -> ask again) in a new test;
reverting the production change makes it fail with
"expected: <ASKING> but was: <FAILED>", confirming it catches the
regression.
2026-09-04 10:52:54 +07:00
3 changed files with 82 additions and 70 deletions
-69
View File
@@ -1,69 +0,0 @@
# Teardown/cleanup audit — dev.ltms.fleet.session
Scope: `SessionManager.java`, `GitWorktrees.java`, `SessionReaper.java`, `MemberSession.java`,
`Worktrees.java` (interface). Read-only; no code changed.
## Main finding
```
1. SessionManager.java:336-338
2. issue: the final worktree removal in release() is the one step in the whole method
that is not wrapped in try/catch. Every other cleanup step here (hasUncommitted check,
snapshot, listener notification) is defended because a `git` call can throw — exec()'s
own javadoc documents both a non-zero exit and its 30-second timeout as normal failure
modes, and every sibling worktrees.* call in this class is guarded against exactly that.
By the time this line runs, registry.remove(paneId) and handles.remove(paneId) have
already happened and launcher.stop(paneId) has already run, so if worktrees.remove()
throws here (e.g. `git worktree remove --force` times out on a stale lock file or a
slow/network filesystem, or exits non-zero), the exception escapes release() with no
way to retry: the paneId is already gone from the registry, so a second stop call is a
no-op and never re-attempts the removal. The worktree directory is now leaked forever,
invisible to `fleet_list`. The caller sees a stop failure — FleetMcp.stop() only catches
HerdrException, and FleetApp.stopMember() catches nothing — even though the session was
in fact fully torn down (pane stopped, deregistered, listeners notified).
3. fix: wrap the `worktrees.remove(...)` call at the end of release() in a try/catch that
logs a warning, matching the pattern already used for every other cleanup step in this
method (e.g. cleanupAfterAddFailure's own worktree/branch removal, or the dirty-check
catch above it).
4. severity: medium
```
## Secondary findings
```
1. SessionManager.java:503-514 (acquireWithWorktree's catch block)
2. issue: after worktrees.add() succeeds, if overlayParity(), shareWithGroup(), or
launcher.spawn() then throws, the catch block removes only the worktree
(worktrees.remove(repoRoot, path)) and never deletes the branch `git worktree add`
created. GitWorktrees.cleanupAfterAddFailure — the sibling cleanup for failures inside
add() itself — explicitly deletes the branch too, with a `-D` and a documented reason
("a branch that never finished provisioning has no session, no PR, nothing else
pointing at it"). That reasoning applies equally here, but this later catch block (the
one covering the three post-add() steps) omits it. Since spawn failures are a normal,
recurring event (this very branch already logs "spawn failed for profile=..."), this
leaks an orphan `worker/<slug>-<nonce>` branch in the shared repo on every such failure,
with nothing pointing at it once the (failed) session is never registered.
3. fix: after worktrees.remove(...) succeeds in this catch, also delete the branch with
`git branch -D branch` (best-effort, log-only on failure), matching
cleanupAfterAddFailure's own two-step cleanup.
4. severity: low
```
No other issue in this scope survived a read of every exit of `add()`, `remove()`,
`snapshot()`, `overlayParity()`, `shareWithGroup()`/`shareRootWithGroup()`, `release()`,
`reapIdle()`, `drainAll()`, and the `SessionReaper` loop. Two shapes I checked and ruled
out as not reachable / not defects:
- `release()`'s `worktrees.repoRoot(removed.cwd())` looked suspicious because `cwd` for a
worktree session is the worktree path itself, so `repoRoot` would equal `worktreePath` —
but I verified with a live git repo (`git --version` 2.53.0) that
`git -C <worktree> worktree remove --force <same worktree>` works correctly: git
resolves `-C` against the common git dir regardless of which linked worktree it's given,
so this is not a bug.
- `git worktree remove --force` on a worktree containing a nested `.git` directory: I
expected this to need a double `--force` per older git docs, but tested it live and a
single `--force` succeeds on git 2.53.0. Not a live failure mode on this stack.
The `if (x != null)` guard-in-catch shape from fleetd #274 (guard assigned only at the end)
does not recur elsewhere in this scope: every `catch` block that guards on a local now
assigns that local before the risky call it protects, not after.
@@ -927,7 +927,16 @@ public final class MessageService {
}
try {
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(workerSession);
// #282: mirror send()'s registration (:802) so a SECOND fleet_ask inside this same
// resumed turn can re-associate the async ticket with its new turnId via
// markAsyncQuestion — without this, that second ask has no Task to attach to, and
// markAsyncQuestion silently returns null.
Task task = asyncTasksByTurn.get(turnId);
if (task != null) {
asyncTasksByWaiter.put(reply, task);
}
if (!rendezvous.answerAsk(turnId, content)) {
asyncTasksByWaiter.remove(reply);
rendezvous.close(workerSession, reply);
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
}
@@ -935,7 +944,13 @@ public final class MessageService {
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
finishAsyncTask(turnId, result);
// #282: this waiter can resolve with a FRESH question rather than a terminal reply —
// the worker chained a second fleet_ask before replying. Mirror sendAsync's own guard
// (:1000) and leave the ticket open (markAsyncQuestion above already re-armed it under
// the new turnId) instead of completing it here with a QUESTION "reply".
if (result.outcome() != Outcome.QUESTION) {
finishAsyncTask(turnId, result);
}
return result;
} catch (TimeoutException e) {
// The worker resumed but hasn't replied yet — no completion fallback arms an answered
@@ -948,6 +963,7 @@ public final class MessageService {
Thread.currentThread().interrupt();
throw new IllegalStateException("interrupted awaiting reply from " + workerSession, e);
} finally {
asyncTasksByWaiter.remove(reply);
rendezvous.close(workerSession, reply);
}
} finally {
@@ -1106,9 +1122,16 @@ public final class MessageService {
private Task markAsyncQuestion(CompletableFuture<Rendezvous.Resolution> waiter, String text, String turnId) {
Task task = waiter == null ? null : asyncTasksByWaiter.get(waiter);
if (task != null) {
String previousTurnId = task.turnId;
task.question = new Reply(Outcome.QUESTION, text, turnId);
task.turnId = turnId;
asyncTasksByTurn.put(turnId, task);
// #282: a second fleet_ask in the same resumed turn re-arms an already-answered task
// (answer() re-registers it in asyncTasksByWaiter) under a FRESH turnId — drop the old
// key so asyncTasksByTurn does not keep growing by one stale entry per chained ask.
if (previousTurnId != null && !previousTurnId.equals(turnId)) {
asyncTasksByTurn.remove(previousTurnId, task);
}
}
return task;
}
@@ -975,6 +975,64 @@ class MessageServiceTest {
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
/**
* fleetd #282: a worker that chains a SECOND {@code fleet_ask} inside the same resumed turn —
* before it ever calls {@code fleet_reply} — used to kill its own async ticket. {@code answer()}
* opens a fresh forward waiter but (unlike {@code send()}) never registered it in
* {@code asyncTasksByWaiter}, so the second ask's {@code markAsyncQuestion} found no {@code Task}
* to re-associate. That waiter still resolved with the second {@code QUESTION} once the worker
* asked again, and {@code answer()} completed the ticket's future with that QUESTION "reply"
* unconditionally — so {@code fleet_poll} reported FAILED while the worker was still alive and
* the primary was mid-conversation with it.
*
* <p>Driven entirely through {@code MessageService}'s public API (sendAsync/ask/answer/poll) —
* never by reaching into {@link Rendezvous} or the task maps directly, so this test cannot pass
* for a reason unrelated to the real bug.
*/
@Test
void secondFleetAskInTheSameResumedTurnDoesNotKillTheAsyncTicket() throws Exception {
String ticket = messages.sendAsync(T, "task that asks twice");
awaitWaiting();
// The worker's first fleet_ask.
CompletableFuture<MessageService.AskResult> ask1 =
CompletableFuture.supplyAsync(() -> messages.ask(T, "Q1", 5000));
MessageService.TaskView asking1 = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertEquals("Q1", asking1.reply());
// The primary answers it — answer() resumes the turn and blocks for what comes next.
CompletableFuture<MessageService.Reply> answer1 = CompletableFuture.supplyAsync(
() -> messages.answer(asking1.turnId(), "a1", 5000));
assertEquals("a1", ask1.get(5, TimeUnit.SECONDS).answer());
// Still in the SAME resumed turn — before replying — the worker asks again.
CompletableFuture<MessageService.AskResult> ask2 =
CompletableFuture.supplyAsync(() -> messages.ask(T, "Q2", 5000));
// answer1's own call unblocks with the second QUESTION (documented QUESTION-chaining
// behaviour — see FleetMcp.answer's javadoc: "Answer it by calling fleet_send again with
// turnId=..."). The bug: this used to also kill the async ticket in the process.
MessageService.Reply firstAnswerResult = answer1.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.QUESTION, firstAnswerResult.outcome());
String turnId2 = firstAnswerResult.turnId();
MessageService.TaskView asking2 = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertEquals("Q2", asking2.reply(),
"the ticket must surface the SECOND question, not be dead/FAILED");
assertEquals(turnId2, asking2.turnId());
// The primary answers the second question; the worker finally sends its real fleet_reply.
CompletableFuture<MessageService.Reply> answer2 = CompletableFuture.supplyAsync(
() -> messages.answer(turnId2, "a2", 5000));
assertEquals("a2", ask2.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
assertEquals(MessageService.Outcome.REPLIED, answer2.get(5, TimeUnit.SECONDS).outcome());
MessageService.TaskView done = awaitTicketPhase(ticket, MessageService.Phase.DONE);
assertEquals("done", done.reply());
}
// --- CB-582: fleet_status pendingAsk() ------------------------------------------------------
@Test