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 77 deletions
-76
View File
@@ -1,76 +0,0 @@
# Authorization entry-point audit
Scope reviewed: every route registered in `FleetApp.build`, every tool handler
registered in `FleetMcp`, and their shared `CallerResolver` and `Authz` gate.
## Result
No authorization-action mismatch was found. The only handler whose operation
changes with its arguments is `fleet_poll`. It selects `DRAIN` when `target` is
present and `READ` when it is absent, before it calls either service branch.
`Fleetd` creates one `CallerResolver` and passes that same instance to both
`FleetMcp` and `FleetApp` (`Fleetd.java:615-650`, `720-722`). REST resolves it
in the Javalin pre-handler. MCP resolves it in the transport context extractor.
## REST routes
| Route | Gate and choice location | Branch/target review | Verdict |
|---|---|---|---|
| `GET /healthz` | None | Liveness probe only; intentionally open. | ok |
| `GET /metrics` | `METRICS`, start of `metrics` | No branch or target. | ok |
| `GET /sessions` | `READ`, start of `sessions` | No caller-selected target. | ok |
| `GET /agents` | `READ`, start of `agents` | No caller-selected target. | ok |
| `GET /members` | `READ`, start of `listMembers` | No caller-selected target. | ok |
| `GET /profiles` | `READ`, start of `profiles` | No caller-selected target. | ok |
| `GET /member-credentials` | `READ`, start of `memberCredentials` | No caller-selected target. The view exposes policy names and counts, not values. | ok |
| `POST /members` | `SPAWN`, before query/body parsing in `spawnMember` | Arguments select role, profile, cwd, and worktree. They do not select a different authority type. | ok |
| `DELETE /members/{paneId}` | `STOP`, after reading `paneId` in `stopMember` | Caller can select another pane, but only a primary has `STOP`. | ok |
| `POST /sessions/{id}/message` | `SEND`, after reading path `id` in `sendMessage` | Normal send, async send, and `turnId` answer all deliver a turn/message. `turnId` does not widen the roles allowed to send. | ok |
| `POST /sessions/{id}/reply` | `REPLY`, after reading path `id` in `replyMessage` | Caller can name a target, and `Authz` requires it to equal the connection-resolved terminal. | ok |
| `GET /sessions/{id}/replies` | `DRAIN`, after reading path `id` in `drainReplies` | Removes inbox entries; only a primary has `DRAIN`. | ok |
| `POST /sessions/{id}/ask` | `ASK`, after reading path `id` in `askMessage` | Caller can name a target, and `Authz` requires it to equal the connection-resolved terminal. | ok |
| `GET /sessions/{id}/status` | `READ`, after reading path `id` in `sessionStatus` | Any authenticated role may observe any session. This matches the `READ` policy, which intentionally does not use target ownership. | ok |
| `GET /tasks/{ticket}` | `READ`, start of `taskStatus` | Ticket polling is read-only; no service branch changes the action. | ok |
## MCP tools
| Tool | Gate and choice location | Branch/target review | Verdict |
|---|---|---|---|
| `fleet_send` | `SEND`, at the start of `sendHandler` | `coordId`, `turnId`, synchronous, and async forms all deliver a message or a turn answer. `sessionId` is read before the branch for audit target only. | ok |
| `fleet_reply` | `REPLY`, after deriving the connection terminal in `replyHandler` | No target argument exists. The service always receives the caller's own terminal. | ok |
| `fleet_ask` | `ASK`, after deriving the connection terminal in `askHandler` | No target argument exists. The service always receives the caller's own terminal. | ok |
| `fleet_status` | `READ`, start of `statusHandler` | Any authenticated role may query any session. This is the same intentional `READ` policy as the REST route. | ok |
| `fleet_poll` with `ticket` | `READ`, `pollAction(target)` before dispatch in `pollHandler` | Reads a task only. | ok |
| `fleet_poll` with `target` | `DRAIN`, `pollAction(target)` before dispatch in `pollHandler` | Drains and removes a target inbox. Only a primary has `DRAIN`. | ok |
| `fleet_ack` | `DRAIN`, start of `ackHandler` | Removes one target inbox entry. Only a primary has `DRAIN`. | ok |
| `fleet_spawn` | `SPAWN`, start of `spawnHandler` | Arguments select member configuration only. | ok |
| `fleet_list` | `READ`, start of `listHandler` | No caller-selected target. | ok |
| `fleet_stop` | `STOP`, after reading `paneId` in `stopHandler` | Caller can select another pane, but only a primary has `STOP`. | ok |
| `fleet_profiles` | `READ`, start of `profilesHandler` | No caller-selected target. | ok |
| `fleet_whoami` | `READ`, start of `whoamiHandler` | Reports the connection-resolved caller, not an argument. | ok |
## Surface parity
The matching route/tool pairs use the same action:
| Operation | REST | MCP | Result |
|---|---|---|---|
| spawn | `SPAWN` | `SPAWN` | match |
| stop | `STOP` | `STOP` | match |
| send and ask answer | `SEND` | `SEND` | match |
| reply | `REPLY` | `REPLY` | match |
| ask | `ASK` | `ASK` | match |
| drain replies | `DRAIN` | `DRAIN` for `fleet_poll{target}` and `fleet_ack` | match |
| status | `READ` | `READ` | match |
| task/ticket polling | `READ` | `READ` | match |
| list/profiles | `READ` | `READ` | match |
## Test risk
`FleetMcpAuthzTest` has a direct regression test for the argument-dependent
`fleet_poll` action. It tests the shared role table for the other tools, but it
does not pin each handler's chosen action. This is not a current defect because
the reviewed handlers choose the matching action. A future change that adds an
argument-dependent operation should add a handler-level action-selection test,
like `pollingByTargetIsADrainAndPollingByTicketIsARead`.
@@ -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