Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 8d2763e354 Add authorization entry-point audit 2026-09-04 10:32:23 +07:00
+65 -85
View File
@@ -1,96 +1,76 @@
# Audit: async ticket / rendezvous lifecycle (`fleetd/src/main/java/dev/ltms/fleet/msg/`)
# Authorization entry-point audit
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.
Scope reviewed: every route registered in `FleetApp.build`, every tool handler
registered in `FleetMcp`, and their shared `CallerResolver` and `Authz` gate.
## Main finding
## Result
```
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
```
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.
### Call sequence that reaches it
`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.
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.
## REST routes
### What goes wrong
| 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 |
- `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).
## MCP tools
### Confidence
| 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 |
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.
## Surface parity
## Secondary (much shorter)
The matching route/tool pairs use the same action:
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.
| 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`.