6.5 KiB
Audit: async ticket / rendezvous lifecycle (fleetd/src/main/java/dev/ltms/fleet/msg/)
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. 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
Call sequence that reaches it
- Lead:
fleet_send{sessionId: W, content: "task", wait:false}→sendAsynccreatestask1/ticket1. Its worker thread callssend(W, content, ASYNC_TIMEOUT_MS, onAccepted, task1), which doesasyncTasksByWaiter.put(reply, task1)(line 802) before blocking onreply.get(). - Worker
Wcallsfleet_ask{"Q1"}→ask(W, "Q1", t).markAsyncQuestionfindstask1viaasyncTasksByWaiter, stampstask1.turnId = turnId1,asyncTasksByTurn[turnId1] = task1.resolveQuestionwakes step 1'ssend(), which returnsOutcome.QUESTION;sendAsync's lambda seesQUESTIONand deliberately does not callfinishAsyncTask(lines 1000-1005) —ticket1correctly pollsPhase.ASKING. - Lead polls, sees
ASKING, answers:fleet_send{turnId: turnId1, content: "A1"}→answer(turnId1, "A1", t). This opens a new waiter viarendezvous.open(workerSession)(line 929) but — unlikesend()— never puts it intoasyncTasksByWaiter.answerAsk(turnId1, "A1")unblocks the worker'sask()call.clearAsyncQuestion(turnId1, false)clearstask1.questionbut keepsasyncTasksByTurn[turnId1] = task1(deliberate, per its own javadoc).answer()then blocks on its ownreply.get()(line 936). - Worker
W, still in the same resumed turn, callsfleet_ask{"Q2"}again before replying → a secondask(W, "Q2", t).openAskmintsturnId2.markAsyncQuestionlooks upasyncTasksByWaiter.get(waiter)for the waiteranswer()opened in step 3 — not found (never registered), sotask == null;task1.turnIdstaysturnId1, noasyncTasksByTurn[turnId2]entry is ever created.resolveQuestionstill succeeds (it only needs a live waiter, not aTask) and wakesanswer(turnId1,...)'s blockedreply.get()withResolution(QUESTION, "Q2", turnId2). answer(turnId1,...)(line 936-939):result = Reply(Outcome.QUESTION, "Q2", turnId2);finishAsyncTask(turnId1, result)looks uptask1by the originalturnId1(still stamped from step 3) and unconditionally doestask1.future.complete(result)— completingticket1's future with a QUESTION outcome, then removesasyncTasksByTurn[turnId1].answer()returnsOutcome.QUESTIONto the lead's ownfleet_send{turnId1,...}call (correct, and separately answerable viaturnId2), butticket1is now terminally done.
What goes wrong
fleet_poll{ticket1}now hits thef.isDone()branch inpoll()permanently.Outcome.QUESTIONis notREPLIED/COMPLETED_UNREPLIED(r.completed()is false) and carries noWORKER_FAILED/BACKEND_EXHAUSTEDreason, so it falls through toPhase.FAILED,detail = "no reply — question"— even though the worker is alive and only waiting onturnId2.asyncTasksByTurnno longer has any entry fortask1/target, sohasAsyncQuestion(target)goes back tofalseimmediately, andhasOrphanedDelegationno longer excludes this target's real state correctly either.- If the worker's eventual real
fleet_reply(afterturnId2is answered, or times out and it finishes on its own) is not captured by a chained directanswer(turnId2,...)call,reply()'s fast path (rendezvous.resolve) finds no live waiter,askAnsweredAsyncTasksfinds 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 fromticket1and fromabandon()'s stranded-reply recovery (no openmatchingtask exists any more).
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)
- Root-cause detail, same defect as above —
answer()(line 929) never puts its freshly opened waiter intoasyncTasksByWaiter, unlikesend()(line 802). Even if the outcome guard above is added, a chained second ask still can't be re-attached totask1via the normalmarkAsyncQuestionpath without also fixing this registration gap. - Low / shape only —
pruneTerminalTickets()(line ~1092) is only invoked from insidesendAsync(). A fleet whose sessions stop receiving new async sends (e.g. everything now goes through blockingsend(), or the target churns and workers are torn down) never prunes its already-terminaltasksentries pastTICKET_TTL_NANOS. Not reachable as a "stuck" ticket (tickets still resolve correctly), only as unboundedtasks/asyncTasksByWaiter-adjacent memory growth over a long-lived daemon with no furthersendAsynctraffic; did not verify this is realistic in production traffic patterns, flagging as a shape only.