Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 2dff4b84a4 Audit: async ticket / rendezvous lifecycle in fleetd msg package 2026-09-04 10:35:01 +07:00
+85 -58
View File
@@ -1,69 +1,96 @@
# Teardown/cleanup audit — dev.ltms.fleet.session
# Audit: async ticket / rendezvous lifecycle (`fleetd/src/main/java/dev/ltms/fleet/msg/`)
Scope: `SessionManager.java`, `GitWorktrees.java`, `SessionReaper.java`, `MemberSession.java`,
`Worktrees.java` (interface). Read-only; no code changed.
Scope: `Rendezvous.java` and `MessageService.java` — the lifecycle of an async ticket
(`Task`) and a rendezvous waiter: create, send, ask, answer, resolve, timeout, abandon, prune.
## 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
1. fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java:938
2. issue: answer() completes an async ticket's future with a QUESTION outcome when the worker
asks a second fleet_ask in the same resumed turn, permanently mislabeling a live delegation
as failed and losing its real reply to fleet_poll.
3. fix: guard the finishAsyncTask(turnId, result) call at line 938 the same way sendAsync's
lambda already guards its own call (lines 999-1005): skip it when result.outcome() ==
Outcome.QUESTION, and instead re-associate the task with the new turnId (as markAsyncQuestion
does on the first ask).
4. severity: high
```
## Secondary findings
### Call sequence that reaches it
```
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
```
1. Lead: `fleet_send{sessionId: W, content: "task", wait:false}` → `sendAsync` creates `task1`
/ `ticket1`. Its worker thread calls `send(W, content, ASYNC_TIMEOUT_MS, onAccepted, task1)`,
which does `asyncTasksByWaiter.put(reply, task1)` (line 802) before blocking on
`reply.get()`.
2. Worker `W` calls `fleet_ask{"Q1"}` → `ask(W, "Q1", t)`. `markAsyncQuestion` finds `task1`
via `asyncTasksByWaiter`, stamps `task1.turnId = turnId1`,
`asyncTasksByTurn[turnId1] = task1`. `resolveQuestion` wakes step 1's `send()`, which returns
`Outcome.QUESTION`; `sendAsync`'s lambda sees `QUESTION` and deliberately does **not** call
`finishAsyncTask` (lines 1000-1005) — `ticket1` correctly polls `Phase.ASKING`.
3. Lead polls, sees `ASKING`, answers: `fleet_send{turnId: turnId1, content: "A1"}` →
`answer(turnId1, "A1", t)`. This opens a **new** waiter via `rendezvous.open(workerSession)`
(line 929) but — unlike `send()` — never puts it into `asyncTasksByWaiter`.
`answerAsk(turnId1, "A1")` unblocks the worker's `ask()` call.
`clearAsyncQuestion(turnId1, false)` clears `task1.question` but keeps
`asyncTasksByTurn[turnId1] = task1` (deliberate, per its own javadoc). `answer()` then blocks
on its own `reply.get()` (line 936).
4. Worker `W`, still in the same resumed turn, calls `fleet_ask{"Q2"}` again before replying →
a second `ask(W, "Q2", t)`. `openAsk` mints `turnId2`. `markAsyncQuestion` looks up
`asyncTasksByWaiter.get(waiter)` for the waiter `answer()` opened in step 3 — **not found**
(never registered), so `task == null`; `task1.turnId` stays `turnId1`, no
`asyncTasksByTurn[turnId2]` entry is ever created. `resolveQuestion` still succeeds (it only
needs a live waiter, not a `Task`) and wakes `answer(turnId1,...)`'s blocked `reply.get()`
with `Resolution(QUESTION, "Q2", turnId2)`.
5. `answer(turnId1,...)` (line 936-939): `result = Reply(Outcome.QUESTION, "Q2", turnId2)`;
`finishAsyncTask(turnId1, result)` looks up `task1` by the **original** `turnId1` (still
stamped from step 3) and unconditionally does `task1.future.complete(result)` — completing
`ticket1`'s future with a **QUESTION** outcome, then removes `asyncTasksByTurn[turnId1]`.
`answer()` returns `Outcome.QUESTION` to the lead's own `fleet_send{turnId1,...}` call
(correct, and separately answerable via `turnId2`), but `ticket1` is now terminally done.
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:
### What goes wrong
- `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.
- `fleet_poll{ticket1}` now hits the `f.isDone()` branch in `poll()` permanently. `Outcome.QUESTION`
is not `REPLIED`/`COMPLETED_UNREPLIED` (`r.completed()` is false) and carries no
`WORKER_FAILED`/`BACKEND_EXHAUSTED` reason, so it falls through to
`Phase.FAILED`, `detail = "no reply — question"` — even though the worker is alive and only
waiting on `turnId2`.
- `asyncTasksByTurn` no longer has any entry for `task1`/`target`, so
`hasAsyncQuestion(target)` goes back to `false` immediately, and `hasOrphanedDelegation`
no longer excludes this target's real state correctly either.
- If the worker's eventual real `fleet_reply` (after `turnId2` is answered, or times out and it
finishes on its own) is not captured by a chained direct `answer(turnId2,...)` call,
`reply()`'s fast path (`rendezvous.resolve`) finds no live waiter, `askAnsweredAsyncTasks`
finds no candidate (`task1.future.isDone()` is already true, so it is excluded), and the reply
is silently dropped into the inbox as a **stranded reply** — unreachable from `ticket1` and
from `abandon()`'s stranded-reply recovery (no open `matching` task exists any more).
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.
### Confidence
High. I traced this with no races or interleavings assumed beyond the documented, deterministic
CB-205 chained-ask protocol that `outcomeOf`/`answer()` already generically support (mapping
`Kind.QUESTION` through `answer()`'s own return value is clearly intentional — see the
`Reply.turnId()` javadoc). I did not run the suite, but grepped
`src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` for a test exercising a *second*
`fleet_ask` inside one resumed (answered) turn and found none — `asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter`
and `askSurfacesAsAQuestionAndTheAnswerResumesTheSameTurn` both cover only a single ask per
turn. A `git log -p` on this file also turned up the CB-588 comment (now at
`sendAsync`, describing the `whenComplete` hook) which explicitly frames
`answer()`'s `finishAsyncTask(turnId, result)` as firing "once a QUESTION is resolved" —
i.e. the author modeled that call as inherently terminal, which is exactly the assumption this
bug violates when the resumed turn asks again.
## Secondary (much shorter)
1. **Root-cause detail, same defect as above** — `answer()` (line 929) never puts its freshly
opened waiter into `asyncTasksByWaiter`, unlike `send()` (line 802). Even if the outcome
guard above is added, a chained second ask still can't be re-attached to `task1` via the
normal `markAsyncQuestion` path without also fixing this registration gap.
2. **Low / shape only** — `pruneTerminalTickets()` (line ~1092) is only invoked from inside
`sendAsync()`. A fleet whose sessions stop receiving new async sends (e.g. everything now
goes through blocking `send()`, or the target churns and workers are torn down) never prunes
its already-terminal `tasks` entries past `TICKET_TTL_NANOS`. Not reachable as a "stuck"
ticket (tickets still resolve correctly), only as unbounded `tasks`/`asyncTasksByWaiter`-adjacent
memory growth over a long-lived daemon with no further `sendAsync` traffic; did not verify
this is realistic in production traffic patterns, flagging as a shape only.