diff --git a/AGENTS.md b/AGENTS.md index c622b81a5..74b2896af 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -616,7 +616,32 @@ for a prompt its code submitted (the debug log says "skipped: re-entry"), so the mod's `prompt.submit` hook sees every OTHER prompt (the user's Enter, a notification, another plugin); a turn whose text is exactly one of those is never claimed, so a prompt the user typed can neither be streamed to -Plannotator nor aborted by a cancel, whatever it says. Plan review does not block the session +Plannotator nor aborted by a cancel, whatever it says. Another plugin's prompt +reaches `turn.start` inside the engine's "The plugin sent a message:" +frame while `prompt.submit` saw it bare, so the match also strips that frame. +**Take-over** (`turns.ts`): a prompt typed or delivered while a turn runs fires +`prompt.submit` at Enter with that turn's `turnId`, and the engine folds it +into the turn's next model request. When that turn is our question's and the +origin is on an ALLOWLIST of messages the model must now answer (`composer`, +`bridge`, `slack-ping`, `channel`; never `peer` (owner's call: a peer +session's reply streaming into the panel is better than losing the reviewer's +answer), `task-notification`, `peer-send-message`, `scheduled-trigger`, +`observer*`, `coordinator`, +`projects-relay`, `auto-continuation`, `unclassified` or a kind added later), +streaming stops the moment the prompt reaches the mod's `prompt.submit` hook, +BEFORE `next(e)` runs the hooks beneath (a slow one cannot let the next step +through); the take-over is confirmed when `next(e)` resolves with the prompt, +and a prompt a hook beneath drops releases the hold. The step in flight was +requested before their prompt, so its output is held and released when the +question settles: at the turn's next `turn.step` (that request carries their +prompt) as `done` when the response before it ended the answer (`stop` chunk +`end_turn`, no tool call) and some answer text was shown, else as `taken_over` +(a thinking-only final response would otherwise settle an empty answer); or at `turn.complete` when no +step followed, as `done` with the turn's answer (their prompt then runs as a +turn of its own). The note reads "You typed into this session…" for +`composer` / `bridge` and the neutral "Another message entered this session…" +otherwise. From the take-over on, a cancel closes only the question and an +interrupt answers `ok: false`; Plannotator never aborts that turn. Plan review does not block the session under the mod, so the status is never `blocked` and plan review gets real turns too (verified live). Polls ask for 15 s (750 ms while our question streams, so deltas flush), and stop on 401/403/404/405/503 (`404` = an older CLI without @@ -1067,9 +1092,9 @@ Ask AI providers are detected independently from installed/authenticated local C Automatic resolution is session-only and never writes a preference. Explicit per-origin choices are persisted in cookies, so a user can override the automatic match for one agent without changing the default for another. -**"Ask this session" (session bridge).** A host can let Ask AI be answered by the agent session that opened Plannotator instead of a separate SDK agent. The host implements the small `SessionBridge` interface (`packages/ai/session-bridge.ts`, vendored to Pi: `host`, `status()` = `ready | busy | blocked | gone`, `modes { turn, transient }`, `ask(req, sink, signal)`, optional `interrupt()`) and passes it as `sessionBridge` to the server (`startReviewServer` / `startAnnotateServer` / plan `ServerOptions` on Bun; review, annotate and plan review on Pi), which then registers `SessionBridgeProvider` under id `session-bridge` as the ONLY Ask AI provider, behind the unchanged `/api/ai/*` endpoints: with a bridge (in-process or pull) the SDK providers go into a separate catalog-only registry (`catalogRegistry` on `createAIEndpoints`, set by `createAIRuntime` / `createPiAIRuntime`), so `/api/ai/capabilities` lists only the bridge under `providers` and makes it `defaultProvider`, and `/api/ai/session` answers `503` for any other provider id whatever the client saved. The catalog-only providers exist so the agent-job launchers (Review Agents, Code Tour, Guided Review) keep their discovered Claude / Codex model lists: `?activate=` naming one runs its discovery and reports it under `catalogProviders`, which only `useModelCatalogs` reads; Ask AI never reads that field. Without a bridge (no mod, OpenCode 1, remote mode, `--tailscale`) the SDK providers register exactly as before. `/api/ai/capabilities` adds `label` ("Ask this session · Pi") and live `sessionBridge { host, status, modes }` for it, `models: []` (no model picker). The provider, not the host, enforces one question at a time across threads (`ask_in_flight`), sends no system prompt (the question carries the `SESSION_ASK_HEADER` line plus the surface), and never interrupts implicitly: a busy session answers `agent_busy`, and the client re-asks with `/api/ai/query` `busyPolicy: "wait"` ("Ask when it finishes", streams `status: waiting` until idle) or `"interrupt"` ("Interrupt and ask now", host `interrupt()` then ask). `session_gone` / `session_blocked` offer no other AI (there is none): the answer shows a plain note that the session that opened Plannotator is gone, or waiting on this decision, and Ask AI cannot reach it. Abort cancels only our question; the runtime `detach()`es the bridge before teardown so a decision/exit never stops a turn already running. Client (`resolveAIProviderSelection`): a bridge in the list is the selection in every status, ahead of any saved pick; without one the order is unchanged. The saved provider cookies are never rewritten for this, so a pick still applies to sessions without a bridge; with a bridge the Ask AI bars and Settings → AI show only its label (no provider or model picker). Code review sends the diff's identity (`buildSessionReviewIdentity`), never the patch. **Pi** (`apps/pi-extension/pi-session-bridge.ts`): review, annotate, last and plan review (plan review no longer blocks the session, see "Pi plan review does not block" below), off in remote mode; the question is a `pi.sendMessage({ customType: "plannotator-ask", display: true, details: { askId } }, { triggerTurn: true })` turn read back from `message_start` / `message_update` text deltas / `tool_execution_start` / `agent_end`; listeners register once at load because `pi.on` returns no unsubscribe before Pi 1.0. **OpenCode 2** (`apps/opencode-plugin/opencode-session-bridge.ts`): review, annotate and last ask a real turn — `ctx.session.prompt({ id, text, delivery: "steer", metadata: { source: "plannotator-ask" } })` under a message id generated in OpenCode's own ascending `msg_` format, streamed back from `ctx.event.subscribe()` (`session.inbox.delivered` with our id starts the answer, then `session.text.delta` / `text.ended` / `tool.input.started`, ended by `session.execution.succeeded|failed|interrupted`), with `session.wait` + `session.context` as the fallback when the event stream is missing (upstream #44788). Delivery is "steer", not "queue", on purpose: every V2 command leaves its session-URL notice as a pending steer row, and a queued question would wake the session with that notice promoted alone as its own model turn (seen live on 2.0.22). Busy comes from the execution events cross-checked by a single outstanding `session.wait` probe; abort and "Interrupt and ask now" use `session.interrupt` (whole execution), only on a turn that is ours or on the reviewer's explicit choice. Plan review answers from context only ("Quick answer from this session · OpenCode", `modes { turn: false, transient: true }`): `submit_plan` is a pending tool call, so a real turn cannot run until the decision (a prompt sent then waits for the tool result, verified live), and `markPlanReviewPending` makes EVERY bridge on that session report `blocked` so nothing can interrupt the review. The quick answer is `session.generate` (no transcript write; verified live on 2.0.22 with OpenAI during a pending `submit_plan`, and `@opencode/ai`'s `normalizeToolHistory` fills the unfinished call with an error result for every provider); it still offers the model its tools, so the question carries `SESSION_ASK_TRANSIENT_NOTE` and an empty answer (a tool call) is reported as a failure. The embedded plan server takes the bridge in-process; review/annotate/last and the CLI plan fallback run the `plannotator` CLI as a child and use the pull bridge below. **OpenCode 1** has no bridge (a second adapter over V1's different message/event model, not a small lift). **Claude Code** answers through its mod over the pull bridge below (review, annotate, last and plan review, all real turns; see "Claude Code mod"); without the mod Claude Code has no bridge. +**"Ask this session" (session bridge).** A host can let Ask AI be answered by the agent session that opened Plannotator instead of a separate SDK agent. The host implements the small `SessionBridge` interface (`packages/ai/session-bridge.ts`, vendored to Pi: `host`, `status()` = `ready | busy | blocked | gone`, `modes { turn, transient }`, `ask(req, sink, signal)`, optional `interrupt()`) and passes it as `sessionBridge` to the server (`startReviewServer` / `startAnnotateServer` / plan `ServerOptions` on Bun; review, annotate and plan review on Pi), which then registers `SessionBridgeProvider` under id `session-bridge` as the ONLY Ask AI provider, behind the unchanged `/api/ai/*` endpoints: with a bridge (in-process or pull) the SDK providers go into a separate catalog-only registry (`catalogRegistry` on `createAIEndpoints`, set by `createAIRuntime` / `createPiAIRuntime`), so `/api/ai/capabilities` lists only the bridge under `providers` and makes it `defaultProvider`, and `/api/ai/session` answers `503` for any other provider id whatever the client saved. The catalog-only providers exist so the agent-job launchers (Review Agents, Code Tour, Guided Review) keep their discovered Claude / Codex model lists: `?activate=` naming one runs its discovery and reports it under `catalogProviders`, which only `useModelCatalogs` reads; Ask AI never reads that field. Without a bridge (no mod, OpenCode 1, remote mode, `--tailscale`) the SDK providers register exactly as before. `/api/ai/capabilities` adds `label` ("Ask this session · Pi") and live `sessionBridge { host, status, modes }` for it, `models: []` (no model picker). The provider, not the host, enforces one question at a time across threads (`ask_in_flight`), sends no system prompt (the question carries the `SESSION_ASK_HEADER` line plus the surface), and never interrupts implicitly: a busy session answers `agent_busy`, and the client re-asks with `/api/ai/query` `busyPolicy: "wait"` ("Ask when it finishes", streams `status: waiting` until idle) or `"interrupt"` ("Interrupt and ask now", host `interrupt()` then ask). `session_gone` / `session_blocked` offer no other AI (there is none): the answer shows a plain note that the session that opened Plannotator is gone, or waiting on this decision, and Ask AI cannot reach it. Abort cancels only our question; the runtime `detach()`es the bridge before teardown so a decision/exit never stops a turn already running. **Take-over (every host):** once another message enters the turn answering our question (the person typing into the session, another extension's steer; never a background task's notification, an engine notice, a compaction, or on Claude Code a peer session's message), the rest of that turn answers THAT message, so the host stops streaming at once, keeps what it already sent, and settles with the bridge error code `taken_over` (`SessionBridgeErrorCode`; the provider maps it to `session_taken_over` with `SESSION_ASK_TAKEN_OVER_TEXT`, the neutral "Another message entered this session while it was answering, so the rest of the reply went to that message."; the Claude Code mod sends `SESSION_ASK_TAKEN_OVER_BY_PERSON_TEXT`, "You typed into this session…", when it knows the person typed it). When the answer had already finished before the other message entered (the turn's last model response stopped with no tool call: Claude `end_turn`, Pi `stop`, OpenCode `finish: "stop"` with no `session.tool.input.started` in that step, since some OpenAI-compatible providers report "stop" on a tool-calling step; and some answer text was shown), it settles `done` instead: a follow-up such as Plannotator's own decision never turns a complete answer into a cut one. From then on neither a Stop nor "Interrupt and ask now" stops that turn: Stop only closes the question, and `interrupt()` refuses (`SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT`) while the taken-over run lasts. The client (`useAIChat`) treats `session_taken_over` as a note, not an error: the partial answer stays and `SessionAskNote` renders the text under it (`response.notice`). Per host: Claude Code mod, `prompt.submit` with the question turn's `turnId` from an allowlisted origin (see "Claude Code mod"); Pi, a `user` message, or a custom message other than ours delivered right after a `turn_start` (a steer or follow-up), in the run answering our question; a non-triggering display-only custom message (e.g. `plannotator-handoff`) is appended at `turn_end`, outside that window, and takes nothing over; the run stays protected until `agent_settled` (idle on Pi without it), so an error retry of the person's run is covered; OpenCode 2, a `session.inbox.delivered` for another row whose kind, read from `session.inbox.enqueued` (`item.type`), is `user` (synthetic, compaction, move and unseen rows never take over), after ours was delivered AND the model began answering (`session.step.started` or any answer output), so a command's session-URL notice promoted in the same batch as the question is not a take-over. **Older servers:** the pull server advertises the codes it knows as `features` on every poll answer (`SESSION_BRIDGE_POLL_FEATURES`, today `["taken_over"]`). A pull host (the mod's `bridge.ts`, `runPullSessionBridgeClient`) talking to a server without `taken_over` settles a take-over as `done` with the partial answer plus the note as its last paragraph (`takenOverFallback`), because an older server reads the unknown code as `failed` and its UI would replace the partial answer with the error. Client (`resolveAIProviderSelection`): a bridge in the list is the selection in every status, ahead of any saved pick; without one the order is unchanged. The saved provider cookies are never rewritten for this, so a pick still applies to sessions without a bridge; with a bridge the Ask AI bars and Settings → AI show only its label (no provider or model picker). Code review sends the diff's identity (`buildSessionReviewIdentity`), never the patch. **Pi** (`apps/pi-extension/pi-session-bridge.ts`): review, annotate, last and plan review (plan review no longer blocks the session, see "Pi plan review does not block" below), off in remote mode; the question is a `pi.sendMessage({ customType: "plannotator-ask", display: true, details: { askId } }, { triggerTurn: true })` turn read back from `message_start` / `message_update` text deltas / `tool_execution_start` / `agent_end`; listeners register once at load because `pi.on` returns no unsubscribe before Pi 1.0. **OpenCode 2** (`apps/opencode-plugin/opencode-session-bridge.ts`): review, annotate and last ask a real turn — `ctx.session.prompt({ id, text, delivery: "steer", metadata: { source: "plannotator-ask" } })` under a message id generated in OpenCode's own ascending `msg_` format, streamed back from `ctx.event.subscribe()` (`session.inbox.delivered` with our id starts the answer, then `session.text.delta` / `text.ended` / `tool.input.started`, ended by `session.execution.succeeded|failed|interrupted`), with `session.wait` + `session.context` as the fallback when the event stream is missing (upstream #44788). Delivery is "steer", not "queue", on purpose: every V2 command leaves its session-URL notice as a pending steer row, and a queued question would wake the session with that notice promoted alone as its own model turn (seen live on 2.0.22). Busy comes from the execution events cross-checked by a single outstanding `session.wait` probe; abort and "Interrupt and ask now" use `session.interrupt` (whole execution), only on a turn that is ours or on the reviewer's explicit choice. Plan review answers from context only ("Quick answer from this session · OpenCode", `modes { turn: false, transient: true }`): `submit_plan` is a pending tool call, so a real turn cannot run until the decision (a prompt sent then waits for the tool result, verified live), and `markPlanReviewPending` makes EVERY bridge on that session report `blocked` so nothing can interrupt the review. The quick answer is `session.generate` (no transcript write; verified live on 2.0.22 with OpenAI during a pending `submit_plan`, and `@opencode/ai`'s `normalizeToolHistory` fills the unfinished call with an error result for every provider); it still offers the model its tools, so the question carries `SESSION_ASK_TRANSIENT_NOTE` and an empty answer (a tool call) is reported as a failure. The embedded plan server takes the bridge in-process; review/annotate/last and the CLI plan fallback run the `plannotator` CLI as a child and use the pull bridge below. **OpenCode 1** has no bridge (a second adapter over V1's different message/event model, not a small lift). **Claude Code** answers through its mod over the pull bridge below (review, annotate, last and plan review, all real turns; see "Claude Code mod"); without the mod Claude Code has no bridge. -**Pull bridge (host-neutral, `packages/ai/session-bridge-pull.ts` server half, vendored to Pi; `session-bridge-pull-client.ts` host half).** For a host that runs the Plannotator server as a SEPARATE process and must not open a listener of its own (the OpenCode plugin's CLI child and the Claude Code mod). The host generates a per-launch secret and starts the server with `PLANNOTATOR_SESSION_BRIDGE_TOKEN` (>= 32 chars), `PLANNOTATOR_SESSION_BRIDGE_HOST` (`opencode` / `claude-code` / `pi`) and optional `PLANNOTATOR_SESSION_BRIDGE_MODES` (`turn,transient`; default `turn`). The Bun CLI takes that config at its very first line (`takeEnvPullSessionBridgeConfig`, cached once per process; `createAIRuntime` reads the cache) and deletes the three variables, so nothing it spawns (git/gh before the server starts, agent jobs, terminals, the auto-update wrapper) inherits the token; `--tailscale` discards it outright (`discardEnvPullSessionBridgeConfig`). The runtime registers the same `SessionBridgeProvider` over it; it is off in remote mode and under `--tailscale` (and an in-process bridge wins). Pi's `createPiAIRuntime` accepts the same config as its `pullSessionBridge` option (no env takeover: Pi itself always bridges in-process). The host learns the port (OpenCode: `PLANNOTATOR_READY_FILE`) and then talks to two endpoints on `http://127.0.0.1:`, both `POST` + JSON with `Authorization: Bearer `: `/api/ai/bridge/poll` `{ status?, modes?, waitMs? }` long-polls (clamped to 25s; answers early when there is work) and returns `{ commands, closing?, superseded? }`, commands being `{ type: "ask", askId, text, mode }`, `{ type: "cancel", askId }` and `{ type: "interrupt", interruptId }`; `/api/ai/bridge/event` takes one event or `{ events: [...] }` — `started`, `delta`, `tool`, `done`, `error` (per question, `code` one of `busy|blocked|gone|aborted|failed`), `status`, and `interrupted { interruptId, ok, message? }` — and answers `409 { code: "ask_not_active" }` for a question that is no longer running, which tells the host to stop it. Commands are re-sent every ~5s until acknowledged (any event for that question; `interrupted` for an interrupt), so hosts dedupe by id. Guards, in order: a loopback Host with the server's own port (the runtime's `authorizeSessionBridgeRequest`, `403 session_bridge_forbidden_host`), no `Origin` header (`403`, a browser is never the host), the bearer token (`401`); without a pull bridge both paths answer `404`, which an older-binary-aware host treats as "no bridge". Liveness: until the host's first request the bridge reports `ready` (a question waits for the first poll); no first request within 30s, or 30s with no open poll after that, is `gone`, and a running question then fails `session_gone`. The newest poll supersedes an open one. A reviewer's Stop drops a question the host never confirmed, else queues `cancel` and frees the slot when the host confirms (or after 15s); `detach()` (decision / shutdown) drops only unconfirmed questions, and `dispose()` answers the open poll `closing`. Both runtimes route the two paths through `createAIEndpoints` (`pullBridge` dep), and Bun lifts the idle timeout for the poll (`isLongLivedAIEndpointPath`). Tests: `packages/ai/session-bridge-pull.test.ts` (fake host over the real client: streaming, early question, busy wait / interrupt, abort before and after pickup, host disappears, never connects, transient, bad token / Origin / rebinding Host, supersede, dispose, detach), `packages/server/ai-runtime.sessionBridge.test.ts` (env takeover, scrub, remote-off), `apps/opencode-plugin/session-bridge-cli.test.ts` (plugin ↔ CLI child end to end). +**Pull bridge (host-neutral, `packages/ai/session-bridge-pull.ts` server half, vendored to Pi; `session-bridge-pull-client.ts` host half).** For a host that runs the Plannotator server as a SEPARATE process and must not open a listener of its own (the OpenCode plugin's CLI child and the Claude Code mod). The host generates a per-launch secret and starts the server with `PLANNOTATOR_SESSION_BRIDGE_TOKEN` (>= 32 chars), `PLANNOTATOR_SESSION_BRIDGE_HOST` (`opencode` / `claude-code` / `pi`) and optional `PLANNOTATOR_SESSION_BRIDGE_MODES` (`turn,transient`; default `turn`). The Bun CLI takes that config at its very first line (`takeEnvPullSessionBridgeConfig`, cached once per process; `createAIRuntime` reads the cache) and deletes the three variables, so nothing it spawns (git/gh before the server starts, agent jobs, terminals, the auto-update wrapper) inherits the token; `--tailscale` discards it outright (`discardEnvPullSessionBridgeConfig`). The runtime registers the same `SessionBridgeProvider` over it; it is off in remote mode and under `--tailscale` (and an in-process bridge wins). Pi's `createPiAIRuntime` accepts the same config as its `pullSessionBridge` option (no env takeover: Pi itself always bridges in-process). The host learns the port (OpenCode: `PLANNOTATOR_READY_FILE`) and then talks to two endpoints on `http://127.0.0.1:`, both `POST` + JSON with `Authorization: Bearer `: `/api/ai/bridge/poll` `{ status?, modes?, waitMs? }` long-polls (clamped to 25s; answers early when there is work) and returns `{ commands, closing?, superseded? }`, commands being `{ type: "ask", askId, text, mode }`, `{ type: "cancel", askId }` and `{ type: "interrupt", interruptId }`; `/api/ai/bridge/event` takes one event or `{ events: [...] }` — `started`, `delta`, `tool`, `done`, `error` (per question, `code` one of `busy|blocked|gone|aborted|failed|taken_over`; anything else reads as `failed`); every poll answer also carries `features` (optional protocol features the server knows, `SESSION_BRIDGE_POLL_FEATURES`), `status`, and `interrupted { interruptId, ok, message? }` — and answers `409 { code: "ask_not_active" }` for a question that is no longer running, which tells the host to stop it. Commands are re-sent every ~5s until acknowledged (any event for that question; `interrupted` for an interrupt), so hosts dedupe by id. Guards, in order: a loopback Host with the server's own port (the runtime's `authorizeSessionBridgeRequest`, `403 session_bridge_forbidden_host`), no `Origin` header (`403`, a browser is never the host), the bearer token (`401`); without a pull bridge both paths answer `404`, which an older-binary-aware host treats as "no bridge". Liveness: until the host's first request the bridge reports `ready` (a question waits for the first poll); no first request within 30s, or 30s with no open poll after that, is `gone`, and a running question then fails `session_gone`. The newest poll supersedes an open one. A reviewer's Stop drops a question the host never confirmed, else queues `cancel` and frees the slot when the host confirms (or after 15s); `detach()` (decision / shutdown) drops only unconfirmed questions, and `dispose()` answers the open poll `closing`. Both runtimes route the two paths through `createAIEndpoints` (`pullBridge` dep), and Bun lifts the idle timeout for the poll (`isLongLivedAIEndpointPath`). Tests: `packages/ai/session-bridge-pull.test.ts` (fake host over the real client: streaming, early question, busy wait / interrupt, abort before and after pickup, host disappears, never connects, transient, bad token / Origin / rebinding Host, supersede, dispose, detach), `packages/server/ai-runtime.sessionBridge.test.ts` (env takeover, scrub, remote-off), `apps/opencode-plugin/session-bridge-cli.test.ts` (plugin ↔ CLI child end to end). **Model lists come from the installed tools, not hand lists.** Claude's list is the SDK's `supportedModels()` against the installed `claude` (a throwaway process with no prompt, `settingSources: []` so no user hooks run, no MCP servers, 10s cap, ~0.5s measured); Codex's is the app-server `model/list`. Both run lazily behind the provider initializer (`?activate=` or the first session), never at startup; a success is kept for the process, a failure may be retried after 60s (`createBestEffortOnce`), and the capabilities answer marks each such provider's list `modelsSource: 'fallback' | 'discovered'` so the client forgets a fallback answer and retries on its next load. The same discovery captures the tool version once (`claude --version` run alongside it; codex from the app-server initialize `userAgent`, via `cliVersionFrom`), kept even when discovery fails, and the answer carries it as an optional `toolVersion`; every Claude/Codex model picker (Agents tab launchers, Guided Review, Ask AI bars, Settings) shows a muted `ModelSourceHint` line ("From your installed Codex 0.155.1", or the built-in-list variant on `fallback`) only when `toolVersion` is present, so a host that omits it gets no hint. Both runtimes share `createDeferredModelDiscovery`: an `?activate=` probe waits for discovery, and so does a session for every provider except Claude, whose sessions resolve the model against the current list while discovery finishes in the background, so the first Ask AI answer does not wait on it — except when that list is still the fallback and lacks the requested pick (e.g. `opus[1m]`), where the session waits for discovery rather than silently running a different model once. One shape, `CatalogModel` in `packages/core/model-catalog.ts` (id, label, efforts + default effort, fast mode, `resolvedId`), serves Ask AI AND the review / Code Tour / Guided Review launchers: `useModelCatalogs` (`packages/ui/hooks/`) fetches the same `/api/ai/capabilities?activate=` answer for the ONE engine a launcher is set to, and `useAgentSettings` resolves saved picks against it at read time (the cookie is never rewritten). ONE resolver, `resolveModelChoice`, is used by the launchers, Ask AI's client (`aiProvider.ts`) and the AI session endpoint: exact id → the alias whose `resolvedId` covers it → the same family's alias (`claude-opus-5` → `opus`) → the surface default (Claude `opus` for review, `sonnet` for tour/guide; Codex `''` = the model the Codex list marks default, except Guided Review, which prefers `gpt-6-luna` when offered: `PREFERRED_GUIDE_CODEX_MODEL`) → the catalog default; it never moves a pick onto a `[1m]` id unless the pick was one. Efforts clamp to the model's own levels, launches wait until the catalog settles, and a loading row shows meanwhile. The Claude catalog drops the SDK's `default` pointer row and adds a bare latest alias per family offered, and always offers `opus` / `sonnet` / `haiku` (some CLIs, e.g. Claude Code 2.1.141, name Opus only through the `default` row, whose description then supplies the version); alias labels carry the version parsed from the row's `resolvedModel` (`claudeModelVersion`: "Opus 5.5 (latest)"), since the SDK's `displayName` has none. Codex fast mode is dropped when resolution replaces a saved model with a different one; changing Fast or reasoning then re-keys the section to the model shown. Codex fast support is read from `serviceTiers` (a `priority`/`fast` tier) as well as the deprecated `additionalSpeedTiers`. The only static lists are the small `CLAUDE_FALLBACK_MODELS` / `CODEX_FALLBACK_MODELS`, used when discovery fails. diff --git a/apps/hook/hooks/mod/bridge.test.ts b/apps/hook/hooks/mod/bridge.test.ts index 9543cb5f0..e2cb815c7 100644 --- a/apps/hook/hooks/mod/bridge.test.ts +++ b/apps/hook/hooks/mod/bridge.test.ts @@ -5,9 +5,9 @@ */ import { afterEach, describe, expect, test } from 'bun:test' import { createPullSessionBridge, type PullSessionBridge } from '../../../../packages/ai/session-bridge-pull.ts' -import { createBridge } from './bridge' +import { createBridge, TAKEN_OVER_INTERRUPT_TEXT, takenOverFallback } from './bridge' import { fakeHost, type FakeHost } from './testing/fake-host' -import { TurnTracker } from './turns' +import { TAKEN_OVER_BY_PERSON_TEXT, TurnTracker } from './turns' const TOKEN = 'k'.repeat(64) let server: PullSessionBridge | null = null @@ -19,7 +19,7 @@ afterEach(() => { server = null }) -function wire(host: FakeHost, bridge: PullSessionBridge) { +function wire(host: FakeHost, bridge: PullSessionBridge, options: { olderServer?: boolean } = {}) { host.onFetch = async (url, body) => { const response = await bridge.handle( new Request(url, { @@ -29,7 +29,13 @@ function wire(host: FakeHost, bridge: PullSessionBridge) { }), ) if (!response) return { status: 404, ok: false, text: '' } - return { status: response.status, ok: response.ok, text: await response.text() } + let text = await response.text() + if (options.olderServer && url.endsWith('/poll')) { + // A server from before the `features` advert. + const { features: _features, ...rest } = JSON.parse(text) + text = JSON.stringify(rest) + } + return { status: response.status, ok: response.ok, text } } } @@ -42,7 +48,7 @@ async function until(check: () => boolean, ms = 3_000) { } function collector() { - const seen = { deltas: '', done: null as string | null, error: null as string | null } + const seen = { deltas: '', done: null as string | null, error: null as string | null, message: null as string | null } return { seen, sink: { @@ -52,8 +58,9 @@ function collector() { done: (answer: string) => { seen.done = answer }, - error: (code: string) => { + error: (code: string, message?: string) => { seen.error = code + seen.message = message ?? null }, }, } @@ -135,6 +142,79 @@ describe('Ask this session over the pull bridge', () => { await running }) + // The failure this guards: the person typed into the question's turn, and + // the reply to their prompt streamed into Plannotator as the answer, and + // "Interrupt and ask now" then aborted the person's own work. + test('a prompt typed into the question\'s turn settles it as taken over, and the turn is never interrupted', async () => { + live = true + server = createPullSessionBridge({ token: TOKEN, host: 'claude-code', modes: { turn: true, transient: false } }) + const host = fakeHost() + wire(host, server) + const turns = new TurnTracker() + const client = createBridge({ host, baseUrl: 'http://127.0.0.1:4321', token: TOKEN, turns, isLive: () => live }) + const running = client.run() + + await until(() => server!.bridge.status() === 'ready') + const { seen, sink } = collector() + const controller = new AbortController() + server.bridge.ask({ askId: 'ask-3', text: '[Plannotator Ask AI] Why step 2?', mode: 'turn' }, sink, controller.signal) + await until(() => host.submits.length === 1) + + turns.onTurnStart('turn-3', 'The plannotator plugin sent a message:\n[Plannotator Ask AI] Why step 2?') + turns.onStep('turn-3') + turns.onText('turn-3', 'Because ') + turns.onPromptEntered({ text: 'also fix the tests', fromUs: false, turnId: 'turn-3', originKind: 'composer' }) + turns.onStep('turn-3') + turns.onText('turn-3', 'Fixed the tests.') + + await until(() => seen.error !== null) + expect(seen.error).toBe('taken_over') + expect(seen.message).toBe(TAKEN_OVER_BY_PERSON_TEXT) + expect(seen.deltas).toBe('Because ') + + // A late Stop and "Interrupt and ask now" both leave the person's turn alone. + controller.abort() + await expect(server.bridge.interrupt!()).rejects.toThrow(TAKEN_OVER_INTERRUPT_TEXT) + expect(host.aborted).toEqual([]) + + live = false + server.dispose() + await running + }) + + // The failure this guards: a newer mod against an older CLI, whose server + // reads `taken_over` as `failed` and whose UI then replaces the partial + // answer with the error. + test('against a server that does not advertise taken_over, a take-over settles as the partial answer plus the note', async () => { + live = true + server = createPullSessionBridge({ token: TOKEN, host: 'claude-code', modes: { turn: true, transient: false } }) + const host = fakeHost() + wire(host, server, { olderServer: true }) + const turns = new TurnTracker() + const client = createBridge({ host, baseUrl: 'http://127.0.0.1:4321', token: TOKEN, turns, isLive: () => live }) + const running = client.run() + + await until(() => server!.bridge.status() === 'ready') + const { seen, sink } = collector() + server.bridge.ask({ askId: 'ask-4', text: '[Plannotator Ask AI] Why step 2?', mode: 'turn' }, sink, new AbortController().signal) + await until(() => host.submits.length === 1) + turns.onTurnStart('turn-4', 'The plannotator plugin sent a message:\n[Plannotator Ask AI] Why step 2?') + turns.onStep('turn-4') + turns.onText('turn-4', 'Because ') + turns.onPromptEntered({ text: 'also fix the tests', fromUs: false, turnId: 'turn-4', originKind: 'composer' }) + turns.onStep('turn-4') + + await until(() => seen.done !== null) + expect(seen.error).toBeNull() + expect(seen.done).toBe(takenOverFallback('Because ', TAKEN_OVER_BY_PERSON_TEXT).answer) + expect(seen.deltas).toBe(seen.done) + expect(seen.done).toContain(TAKEN_OVER_BY_PERSON_TEXT) + + live = false + server.dispose() + await running + }) + test('a wrong token stops the client instead of retrying forever', async () => { live = true server = createPullSessionBridge({ token: 'x'.repeat(64), host: 'claude-code', modes: { turn: true, transient: false } }) diff --git a/apps/hook/hooks/mod/bridge.ts b/apps/hook/hooks/mod/bridge.ts index 8485c28ca..b1374ba59 100644 --- a/apps/hook/hooks/mod/bridge.ts +++ b/apps/hook/hooks/mod/bridge.ts @@ -14,7 +14,8 @@ * (it waits for idle), the turn's streamed text goes back as deltas, and * `turn.complete`'s answer as `done`. Busy = Claude is mid-turn: reported as * `busy`, so the reviewer chooses wait or interrupt; an interrupt aborts the - * running turn with `$.turn.abort`. Plan review does not block the session + * running turn with `$.turn.abort`, except a question's turn the person typed + * into (turns.ts, take-over), which is theirs. Plan review does not block the session * under the mod, so the status is never `blocked`. */ @@ -27,6 +28,28 @@ export const BRIDGE_HOST = 'claude-code' export const BRIDGE_MODES = 'turn' /** Long-poll wait we ask for; below the server's 25 s cap and any fetch timeout. */ export const BRIDGE_POLL_WAIT_MS = 15_000 +/** + * Why "Interrupt and ask now" refuses a turn another message took over. Same + * text as `SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT` (packages/ai/session-bridge.ts). + */ +export const TAKEN_OVER_INTERRUPT_TEXT = + 'The session is now answering another message, so Plannotator will not stop it. Ask when it finishes instead.' + +/** Used when a `taken_over` comes without a message. */ +const TAKEN_OVER_FALLBACK_NOTE = + 'Another message entered this session while it was answering, so the rest of the reply went to that message.' + +/** + * A server that does not advertise `taken_over` (poll `features`) reads it as + * `failed`, and its UI then replaces the partial answer with the error. Settle + * as an answer instead: what streamed, plus the note as its last paragraph. + * Mirrors `takenOverFallback` in packages/ai/session-bridge-pull-client.ts. + */ +export function takenOverFallback(streamed: string, message: string | undefined): { delta: string; answer: string } { + const note = `_${(message || TAKEN_OVER_FALLBACK_NOTE).trim()}_` + const delta = streamed ? `\n\n${note}` : note + return { delta, answer: `${streamed}${delta}` } +} type BridgeCommand = | { type: 'ask'; askId: string; text: string; mode: string } @@ -57,16 +80,17 @@ export function bridgeBaseUrl(port: number): string { return `http://127.0.0.1:${port}` } -export function parseBridgeCommands(text: string): { commands: BridgeCommand[]; closing: boolean } { +export function parseBridgeCommands(text: string): { commands: BridgeCommand[]; closing: boolean; features: string[] } { try { - const body = JSON.parse(text) as { commands?: unknown; closing?: unknown } + const body = JSON.parse(text) as { commands?: unknown; closing?: unknown; features?: unknown } const commands = Array.isArray(body.commands) ? body.commands.filter((command): command is BridgeCommand => !!command && typeof command === 'object' && typeof (command as { type?: unknown }).type === 'string') : [] - return { commands, closing: body.closing === true } + const features = Array.isArray(body.features) ? body.features.filter((feature): feature is string => typeof feature === 'string') : [] + return { commands, closing: body.closing === true, features } } catch { - return { commands: [], closing: false } + return { commands: [], closing: false, features: [] } } } @@ -87,6 +111,8 @@ export function createBridge(options: BridgeOptions): BridgeHandle { let outbox: BridgeEvent[] = [] let sending: Promise = Promise.resolve() let lastStatus: 'ready' | 'busy' = turns.busy ? 'busy' : 'ready' + /** The server knows the `taken_over` code (poll `features`). */ + let serverTakesTakenOver = false const post = async (events: BridgeEvent[]): Promise => { if (events.length === 0) return @@ -133,11 +159,23 @@ export function createBridge(options: BridgeOptions): BridgeHandle { if (seenAsks.has(command.askId)) return seenAsks.add(command.askId) const askId = command.askId + let streamed = '' const sink: AskSink = { - delta: (text) => emit({ type: 'delta', askId, text }, false), + delta: (text) => { + streamed += text + emit({ type: 'delta', askId, text }, false) + }, tool: (name) => emit({ type: 'tool', askId, name }), done: (answer) => emit({ type: 'done', askId, answer }), - error: (code, message) => emit({ type: 'error', askId, code, ...(message ? { message } : {}) }), + error: (code, message) => { + if (code === 'taken_over' && !serverTakesTakenOver) { + const fallback = takenOverFallback(streamed, message) + emit({ type: 'delta', askId, text: fallback.delta }, false) + emit({ type: 'done', askId, answer: fallback.answer }) + return + } + emit({ type: 'error', askId, code, ...(message ? { message } : {}) }) + }, } if (!turns.beginAsk(askId, command.text, sink)) { emit({ type: 'error', askId, code: 'busy', message: 'Another question is already running in this session.' }) @@ -157,6 +195,11 @@ export function createBridge(options: BridgeOptions): BridgeHandle { emit({ type: 'interrupted', interruptId, ok: true }) return } + // A question's turn that the person typed into is theirs now: never stopped from Plannotator. + if (turns.isTakenOver(running)) { + emit({ type: 'interrupted', interruptId, ok: false, message: TAKEN_OVER_INTERRUPT_TEXT }) + return + } try { await host.abortTurn(running) emit({ type: 'interrupted', interruptId, ok: true }) @@ -215,7 +258,8 @@ export function createBridge(options: BridgeOptions): BridgeHandle { continue } failures = 0 - const { commands, closing } = parseBridgeCommands(response.text) + const { commands, closing, features } = parseBridgeCommands(response.text) + serverTakesTakenOver = features.includes('taken_over') for (const command of commands) { host.debug(`bridge ${base}: ${command.type}`) handle(command) diff --git a/apps/hook/hooks/mod/controller.ts b/apps/hook/hooks/mod/controller.ts index 9373934c9..b7edaad27 100644 --- a/apps/hook/hooks/mod/controller.ts +++ b/apps/hook/hooks/mod/controller.ts @@ -52,7 +52,7 @@ import { scriptOnlyAnnotateFlag, scriptOnlyAnnotateFlagText, } from './tool' -import { TurnTracker } from './turns' +import { TurnTracker, type EnteredPrompt } from './turns' /** Persisted in `$.store` so open reviews reattach after a restart or `--resume`. */ export interface LaunchRecord { @@ -674,9 +674,34 @@ export class PlannotatorMod { } /** A prompt entered the session (prompt.submit), from register.ts. */ - onPromptEntered(text: string, fromUs: boolean): void { - this.host.debug(`prompt.submit${fromUs ? ' (ours)' : ''}: ${JSON.stringify(text.slice(0, 160))}`) - this.turns.onPromptEntered(text, fromUs) + onPromptEntered(prompt: EnteredPrompt): void { + const { text, fromUs, turnId, originKind } = prompt + this.host.debug( + `prompt.submit${fromUs ? ' (ours)' : ''}${originKind ? ` [${originKind}]` : ''}${turnId ? ` into ${turnId}` : ''}: ${JSON.stringify(text.slice(0, 160))}`, + ) + const wasOurs = !!turnId && this.turns.ownsTurn(turnId) + this.turns.onPromptEntered(prompt) + if (wasOurs && turnId && this.turns.isTakenOver(turnId)) this.host.debug(`ask turn ${turnId} taken over`) + } + + /** A prompt reached prompt.submit, before the hooks beneath it ran, from register.ts. */ + onPromptSubmitting(prompt: Omit): void { + this.turns.onPromptSubmitting(prompt) + } + + /** A prompt announced by onPromptSubmitting did not enter, from register.ts. */ + onPromptDropped(prompt: Omit): void { + this.turns.onPromptDropped(prompt) + } + + /** A model request of a turn is about to go out (turn.step), from register.ts. */ + onTurnStep(turnId: string): void { + this.turns.onStep(turnId) + } + + /** A model response of a turn finished (turn.step's `stop` chunk), from register.ts. */ + onTurnStepStop(turnId: string, stopReason: string | null): void { + this.turns.onStepStop(turnId, stopReason) } /** Turn events, from register.ts. */ diff --git a/apps/hook/hooks/mod/register.ts b/apps/hook/hooks/mod/register.ts index 35017e006..23739ff73 100644 --- a/apps/hook/hooks/mod/register.ts +++ b/apps/hook/hooks/mod/register.ts @@ -298,14 +298,34 @@ export function register(on: On) { }) on('prompt.submit', async ($: Engine, e: any, next: Next) => { - const result = await next(e) const instance = allowed ? mod : null - if (instance && !instance.isDisposed && result && typeof result.text === 'string') { + const live = !!instance && !instance.isDisposed + // A prompt typed (or delivered) while a turn ran carries that turn's id: + // when the turn is a question's and a person sent it, the rest of the turn + // answers this prompt instead (turns.ts, take-over). Streaming stops HERE, + // before the hooks beneath run, so a slow one cannot let the next step + // through; it resumes if one of them drops the prompt. + const origin = e.origin + const ref = { + fromUs: !!origin && origin.kind === 'plugin' && origin.name === PLUGIN_NAME, + ...(typeof e.turnId === 'string' ? { turnId: e.turnId } : {}), + ...(origin && typeof origin.kind === 'string' ? { originKind: origin.kind } : {}), + } + if (live) instance.onPromptSubmitting(ref) + let result + try { + result = await next(e) + } catch (error) { + if (live) instance.onPromptDropped(ref) + throw error + } + if (!live || instance.isDisposed) return result + if (result && typeof result.text === 'string') { // Every prompt seen here is someone else's (the engine skips our hooks // for prompts our own code submitted): its turn is never a question's. - const origin = result.origin ?? e.origin - const fromUs = !!origin && origin.kind === 'plugin' && origin.name === PLUGIN_NAME - instance.onPromptEntered(result.text, fromUs) + instance.onPromptEntered({ ...ref, text: result.text }) + } else { + instance.onPromptDropped(ref) } return result }) @@ -321,12 +341,17 @@ export function register(on: On) { on('turn.step', async function* ($: Engine, e: any, next: Next) { const instance = allowed ? mod : null if (!instance || !instance.turns.ownsTurn(e.turnId)) return yield* next(e) + // A question's turn that someone else's prompt entered settles here: this + // request carries their prompt, so nothing it says is the question's answer. + instance.onTurnStep(e.turnId) + if (!instance.turns.ownsTurn(e.turnId)) return yield* next(e) const stream = next(e) let step = await stream.next() while (!step.done) { const chunk = step.value if (chunk && chunk.kind === 'text' && typeof chunk.text === 'string') instance.turns.onText(e.turnId, chunk.text) else if (chunk && chunk.kind === 'tool' && typeof chunk.name === 'string') instance.turns.onTool(e.turnId, chunk.name) + else if (chunk && chunk.kind === 'stop') instance.onTurnStepStop(e.turnId, typeof chunk.stopReason === 'string' ? chunk.stopReason : null) yield chunk step = await stream.next() } diff --git a/apps/hook/hooks/mod/turns.test.ts b/apps/hook/hooks/mod/turns.test.ts index 651fb5ac6..14d323ea7 100644 --- a/apps/hook/hooks/mod/turns.test.ts +++ b/apps/hook/hooks/mod/turns.test.ts @@ -1,10 +1,30 @@ import { describe, expect, test } from 'bun:test' -import { TurnTracker, type AskSink } from './turns' +import { + SESSION_ASK_TAKEN_OVER_BY_PERSON_TEXT, + SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT, + SESSION_ASK_TAKEN_OVER_TEXT, +} from '../../../../packages/ai/session-bridge.ts' +import { TAKEN_OVER_INTERRUPT_TEXT } from './bridge' +import { TAKEN_OVER_BY_PERSON_TEXT, TAKEN_OVER_TEXT, TurnTracker, type AskSink } from './turns' -function sink(): AskSink & { deltas: string[]; errors: string[] } { +function sink(): AskSink & { deltas: string[]; errors: string[]; messages: (string | undefined)[]; answers: string[] } { const deltas: string[] = [] const errors: string[] = [] - return { deltas, errors, delta: (t) => deltas.push(t), tool: () => undefined, done: () => undefined, error: (c) => errors.push(c) } + const messages: (string | undefined)[] = [] + const answers: string[] = [] + return { + deltas, + errors, + messages, + answers, + delta: (t) => deltas.push(t), + tool: () => undefined, + done: (answer) => answers.push(answer), + error: (c, m) => { + errors.push(c) + messages.push(m) + }, + } } // The failure these guard: a prompt the USER typed that happens to contain the @@ -15,7 +35,7 @@ describe('TurnTracker claims only its own prompt', () => { const turns = new TurnTracker() const s = sink() turns.beginAsk('a1', 'why?', s) - turns.onPromptEntered('why? explain the parser', false) + turns.onPromptEntered({ text: 'why? explain the parser', fromUs: false }) turns.onTurnStart('user-turn', 'why? explain the parser') turns.onText('user-turn', 'secret user answer') expect(turns.ownsTurn('user-turn')).toBe(false) @@ -26,7 +46,7 @@ describe('TurnTracker claims only its own prompt', () => { const turns = new TurnTracker() turns.beginAsk('a1', 'why?', sink()) expect(turns.cancelAsk('a1')).toBeNull() - turns.onPromptEntered('why?', false) + turns.onPromptEntered({ text: 'why?', fromUs: false }) expect(turns.onTurnStart('user-turn', 'why?')).toBeNull() // Our own submission still gets dropped when it arrives. // Our own submission never reaches our prompt.submit hook (the engine @@ -39,9 +59,203 @@ describe('TurnTracker claims only its own prompt', () => { const s = sink() turns.beginAsk('a1', 'why?', s) // If an engine ever does show it, with our origin, it is still ours. - turns.onPromptEntered('why?', true) + turns.onPromptEntered({ text: 'why?', fromUs: true }) turns.onTurnStart('t1', 'The plannotator plugin sent a message:\nwhy?') turns.onText('t1', 'because') expect(s.deltas).toEqual(['because']) }) }) + +/** A question that is running as turn `t1`, with some answer already streamed. */ +function running() { + const turns = new TurnTracker() + const s = sink() + turns.beginAsk('a1', '[Plannotator Ask AI] Why step 2?', s) + turns.onTurnStart('t1', 'The plannotator plugin sent a message:\n[Plannotator Ask AI] Why step 2?') + turns.onStep('t1') + turns.onText('t1', 'Because ') + return { turns, s } +} + +const typed = (originKind: string, text = 'also fix the tests') => ({ text, fromUs: false, turnId: 't1', originKind }) + +// The failure these guard: the person typed into the turn Plannotator's +// question started, and the reply to THEIR prompt was streamed into +// Plannotator as the answer, and a Plannotator Stop aborted their work. +describe("TurnTracker: a prompt typed into the question's turn takes it over", () => { + test('streaming stops at once and the question settles at the next step with the note', () => { + const { turns, s } = running() + turns.onPromptEntered(typed('composer')) + // The step in flight was requested before their prompt: held, not streamed. + turns.onText('t1', 'it is needed.') + expect(s.deltas).toEqual(['Because ']) + expect(turns.isTakenOver('t1')).toBe(true) + + // The next request carries their prompt: the question settles here. + turns.onStepStop('t1', 'tool_use') + turns.onStep('t1') + expect(s.deltas).toEqual(['Because ', 'it is needed.']) + expect(s.errors).toEqual(['taken_over']) + expect(s.messages).toEqual([TAKEN_OVER_BY_PERSON_TEXT]) + expect(turns.ownsTurn('t1')).toBe(false) + + // Their reply is never streamed, and a late cancel aborts nothing. + turns.onText('t1', 'Fixed the tests.') + expect(s.deltas).toEqual(['Because ', 'it is needed.']) + expect(turns.cancelAsk('a1')).toBeNull() + expect(turns.isTakenOver('t1')).toBe(true) + turns.onTurnComplete('t1', 'Fixed the tests.', false) + expect(turns.isTakenOver('t1')).toBe(false) + expect(s.answers).toEqual([]) + }) + + test('streaming stops when the prompt reaches our hook, before the hooks beneath it settle', () => { + const { turns, s } = running() + turns.onPromptSubmitting(typed('composer')) + turns.onText('t1', 'it is needed.') + expect(s.deltas).toEqual(['Because ']) + // A Stop while it is on its way closes only the question. + expect(turns.isTakenOver('t1')).toBe(true) + expect(turns.cancelAsk('a1')).toBeNull() + expect(s.errors).toEqual(['aborted']) + }) + + test('a prompt a hook beneath dropped releases the hold and the answer streams on', () => { + const { turns, s } = running() + turns.onPromptSubmitting(typed('composer')) + turns.onText('t1', 'it is ') + turns.onPromptDropped(typed('composer')) + turns.onText('t1', 'needed.') + expect(s.deltas).toEqual(['Because ', 'it is ', 'needed.']) + expect(turns.isTakenOver('t1')).toBe(false) + turns.onTurnComplete('t1', 'Because it is needed.', false) + expect(s.answers).toEqual(['Because it is needed.']) + }) + + test('when the response before their prompt ended the answer, the question settles as done', () => { + const { turns, s } = running() + turns.onText('t1', 'it is needed.') + turns.onPromptEntered(typed('composer', 'thanks')) + turns.onStepStop('t1', 'end_turn') + turns.onStep('t1') + expect(s.answers).toEqual(['Because it is needed.']) + expect(s.errors).toEqual([]) + }) + + test('a thinking-only final response (nothing streamed) settles as taken over, not an empty done', () => { + const turns = new TurnTracker() + const s = sink() + turns.beginAsk('a1', 'why?', s) + turns.onTurnStart('t1', 'The plannotator plugin sent a message:\nwhy?') + turns.onStep('t1') + turns.onPromptEntered(typed('composer')) + turns.onStepStop('t1', 'end_turn') + turns.onStep('t1') + expect(s.answers).toEqual([]) + expect(s.errors).toEqual(['taken_over']) + }) + + test('a Stop between the take-over and the next step closes the question, never the turn', () => { + const { turns, s } = running() + turns.onPromptEntered(typed('composer')) + expect(turns.cancelAsk('a1')).toBeNull() + expect(s.errors).toEqual(['aborted']) + turns.onStep('t1') + expect(s.errors).toEqual(['aborted']) + }) + + test('when the turn ends before another step, everything it said answered the question', () => { + const { turns, s } = running() + turns.onPromptEntered(typed('composer', 'thanks')) + turns.onText('t1', 'it is needed.') + turns.onTurnComplete('t1', 'Because it is needed.', false) + expect(s.deltas).toEqual(['Because ', 'it is needed.']) + expect(s.answers).toEqual(['Because it is needed.']) + expect(s.errors).toEqual([]) + // Their prompt then runs as a turn of its own, which is not claimed. + expect(turns.onTurnStart('t2', 'thanks')).toBeNull() + expect(turns.ownsTurn('t2')).toBe(false) + }) + + test('the person (prompt box, Remote Control) gets "You typed"; a channel or Slack ping the neutral note', () => { + for (const [originKind, note] of [ + ['composer', TAKEN_OVER_BY_PERSON_TEXT], + ['bridge', TAKEN_OVER_BY_PERSON_TEXT], + ['channel', TAKEN_OVER_TEXT], + ['slack-ping', TAKEN_OVER_TEXT], + ] as const) { + const { turns, s } = running() + turns.onPromptEntered(typed(originKind)) + turns.onStep('t1') + expect([originKind, s.errors, s.messages]).toEqual([originKind, ['taken_over'], [note]]) + } + }) + + test("the agent's own work and engine notices never take the turn over", () => { + for (const originKind of [ + // The owner's call: a peer session's reply streaming into the panel is + // better than losing the reviewer's answer. + 'peer', + 'task-notification', + 'peer-send-message', + 'scheduled-trigger', + 'observer', + 'observer-activity', + 'coordinator', + 'projects-relay', + 'auto-continuation', + 'unclassified', + 'some-future-kind', + ]) { + const { turns, s } = running() + turns.onPromptSubmitting(typed(originKind)) + turns.onPromptEntered(typed(originKind)) + turns.onStep('t1') + turns.onText('t1', 'it is needed.') + expect([originKind, turns.isTakenOver('t1')]).toEqual([originKind, false]) + expect(s.deltas).toEqual(['Because ', 'it is needed.']) + turns.onTurnComplete('t1', 'Because it is needed.', false) + expect(s.answers).toEqual(['Because it is needed.']) + expect(s.errors).toEqual([]) + } + }) + + test("a prompt typed over someone else's turn changes nothing for a queued question", () => { + const turns = new TurnTracker() + const s = sink() + turns.onTurnStart('user-turn', 'refactor the parser') + turns.beginAsk('a1', 'why?', s) + turns.onPromptEntered({ text: 'and the lexer', fromUs: false, turnId: 'user-turn', originKind: 'composer' }) + expect(turns.isTakenOver('user-turn')).toBe(false) + turns.onTurnComplete('user-turn', 'done', false) + turns.onTurnStart('t1', 'The plannotator plugin sent a message:\nwhy?') + turns.onText('t1', 'because') + expect(s.deltas).toEqual(['because']) + }) +}) + +describe("TurnTracker never claims another plugin's turn", () => { + test("another plugin's prompt, framed at turn.start, is not the question's turn", () => { + const turns = new TurnTracker() + const s = sink() + turns.beginAsk('a1', 'why?', s) + // prompt.submit sees the other plugin's text bare; turn.start sees it framed. + turns.onPromptEntered({ text: 'why? summarize the diff', fromUs: false, originKind: 'plugin' }) + expect(turns.onTurnStart('ws-turn', 'The workspaces plugin sent a message:\nwhy? summarize the diff')).toBeNull() + turns.onText('ws-turn', 'the diff adds a parser') + expect(turns.ownsTurn('ws-turn')).toBe(false) + expect(s.deltas).toEqual([]) + turns.onTurnComplete('ws-turn', 'the diff adds a parser', false) + // Our own question still claims its turn afterwards. + turns.onTurnStart('t1', 'The plannotator plugin sent a message:\nwhy?') + turns.onText('t1', 'because') + expect(s.deltas).toEqual(['because']) + }) +}) + +test('the mod sends the same take-over notes the provider would', () => { + // A hooks module imports only its own files, so the texts are copies. + expect(TAKEN_OVER_TEXT).toBe(SESSION_ASK_TAKEN_OVER_TEXT) + expect(TAKEN_OVER_BY_PERSON_TEXT).toBe(SESSION_ASK_TAKEN_OVER_BY_PERSON_TEXT) + expect(TAKEN_OVER_INTERRUPT_TEXT).toBe(SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT) +}) diff --git a/apps/hook/hooks/mod/turns.ts b/apps/hook/hooks/mod/turns.ts index c807c1d92..a8a0bcab4 100644 --- a/apps/hook/hooks/mod/turns.ts +++ b/apps/hook/hooks/mod/turns.ts @@ -8,13 +8,83 @@ * streamed text goes back to Plannotator as deltas and its final answer as * `done`. A question cancelled while still queued is confirmed at once and its * turn is aborted the moment it starts. + * + * Take-over: a prompt a PERSON put into the question's turn while it ran + * (`prompt.submit` carrying that turn's id, from an origin in + * TAKEOVER_ORIGINS) makes the rest of the turn theirs. Streaming stops when + * the prompt reaches the mod's hook (`onPromptSubmitting`, before the hooks + * beneath run), and the take-over is confirmed once it entered + * (`onPromptEntered`); a prompt a hook beneath dropped releases the hold + * (`onPromptDropped`). The step in flight was requested before their prompt + * existed, so what it says is still the question's answer: it is held and + * released when the ask settles. The ask settles at the turn's next step (the + * engine folds a prompt typed mid-turn into the next model request, so from + * there on the output answers THEM): as `done` when the response before it + * ended the answer (`end_turn`, no tool call) with text shown, else as `taken_over` with the + * note. When the turn ends first (their prompt then runs as a turn of its + * own), it settles with the turn's answer as `done`. Either way Plannotator + * never aborts that turn again: a Stop only closes the question, and + * "Interrupt and ask now" refuses. */ export interface AskSink { delta(text: string): void tool(name: string): void done(answer: string): void - error(code: 'busy' | 'blocked' | 'gone' | 'aborted' | 'failed', message?: string): void + error(code: AskErrorCode, message?: string): void +} + +export type AskErrorCode = 'busy' | 'blocked' | 'gone' | 'aborted' | 'failed' | 'taken_over' + +/** + * Sent with `taken_over`, so a server older than that code still shows why the + * answer stopped. Same texts as `SESSION_ASK_TAKEN_OVER_BY_PERSON_TEXT` and + * `SESSION_ASK_TAKEN_OVER_TEXT` in packages/ai/session-bridge.ts (a hooks + * module imports only its own files; turns.test.ts holds them equal). + */ +export const TAKEN_OVER_BY_PERSON_TEXT = + 'You typed into this session while it was answering, so the rest of the reply went to your prompt.' +export const TAKEN_OVER_TEXT = + 'Another message entered this session while it was answering, so the rest of the reply went to that message.' + +/** What `prompt.submit` reported for a prompt that entered the session. */ +export interface EnteredPrompt { + text: string + /** Its origin is this plugin. */ + fromUs: boolean + /** The turn it was typed over or delivered into (`e.turnId`); absent when the session was idle. */ + turnId?: string + /** `origin.kind` (`composer`, `task-notification`, `plugin`, ...). */ + originKind?: string +} + +/** + * Origins whose prompt, delivered into a running turn, takes that turn over: + * something a person (or another person-driven session) says that the model + * must now answer. An allowlist, so an origin the engine adds later, or one + * that is the agent's own work, takes nothing over: + * - `composer`: the person's Enter in the terminal. `bridge`: the person + * through Remote Control. `slack-ping`: the session's owner from Slack. + * `channel`: a message an MCP channel relays (Slack, Telegram), a person on + * the other end. + * Not listed, so never a take-over: `peer` (another Claude session's message: + * the owner's call is that losing the reviewer's answer is worse than a peer's + * reply streaming into the panel), `task-notification` and + * `peer-send-message` (notifications framed for the agent: a background task + * or another session's SendMessage finishing), `scheduled-trigger`, + * `observer`, `observer-activity`, `coordinator`, `projects-relay`, + * `auto-continuation`, `unclassified` (the engine's idle notices and delivery + * receipts), `sdk`, `plugin` (a plugin's prompt runs once idle and carries no + * turn id). + */ +const TAKEOVER_ORIGINS: ReadonlySet = new Set(['composer', 'bridge', 'slack-ping', 'channel']) +/** Of those, the person typing into this session: the note says "You typed". */ +const PERSON_ORIGINS: ReadonlySet = new Set(['composer', 'bridge']) + +/** The turn a prompt would take over, or null. */ +function takeoverTurn(prompt: Omit): string | null { + if (prompt.fromUs || !prompt.turnId || !prompt.originKind) return null + return TAKEOVER_ORIGINS.has(prompt.originKind) ? prompt.turnId : null } interface ActiveAsk { @@ -24,6 +94,18 @@ interface ActiveAsk { turnId: string | null cancelled: boolean finished: boolean + /** Someone else's prompt entered this turn: nothing more is streamed. */ + takenOver: boolean + /** The take-over came from the person typing (note wording). */ + byPerson: boolean + /** Prompts on their way into this turn (between our prompt.submit hook and next(e) resolving). */ + pending: number + /** Output held while a prompt is pending or after a take-over, released when it settles. */ + held: { kind: 'text' | 'tool'; value: string }[] + /** Text sent to Plannotator so far. */ + streamed: string + /** How the turn's last finished model response stopped (`end_turn`: the answer was complete). */ + lastStop: string | null } export class TurnTracker { @@ -43,25 +125,75 @@ export class TurnTracker { * by a cancel. */ private foreign: string[] = [] + /** Turns that started as a question's and were taken over: Plannotator never aborts them. */ + private takenOverTurns = new Set() + + /** + * A prompt reached the mod's prompt.submit hook, before the hooks beneath it + * run: when it would take the question's turn over, stop streaming now, so a + * slow hook beneath cannot let more of the turn through. + */ + onPromptSubmitting(prompt: Omit): void { + const ask = this.askOn(takeoverTurn(prompt)) + if (ask) ask.pending += 1 + } + + /** A prompt entered (prompt.submit resolved with it). */ + onPromptEntered(prompt: EnteredPrompt): void { + if (prompt.fromUs) return + const value = prompt.text.trim() + if (value) { + this.foreign.push(value) + if (this.foreign.length > 16) this.foreign.shift() + } + const ask = this.askOn(takeoverTurn(prompt)) + if (!ask) return + if (ask.pending > 0) ask.pending -= 1 + if (ask.takenOver) return + ask.takenOver = true + ask.byPerson = !!prompt.originKind && PERSON_ORIGINS.has(prompt.originKind) + this.takenOverTurns.add(ask.turnId!) + } + + /** A prompt announced by onPromptSubmitting did not enter (a hook beneath dropped it). */ + onPromptDropped(prompt: Omit): void { + const ask = this.askOn(takeoverTurn(prompt)) + if (!ask || ask.pending === 0) return + ask.pending -= 1 + if (!ask.takenOver && ask.pending === 0) this.releaseHeld(ask) + } - /** A prompt entered. `fromUs`: its origin is this plugin. */ - onPromptEntered(text: string, fromUs: boolean): void { - if (fromUs) return - const value = text.trim() - if (!value) return - this.foreign.push(value) - if (this.foreign.length > 16) this.foreign.shift() + /** The unfinished question running as `turnId`, if any. */ + private askOn(turnId: string | null): ActiveAsk | null { + const ask = this.ask + return turnId && ask && !ask.finished && ask.turnId === turnId ? ask : null + } + + private holding(ask: ActiveAsk): boolean { + return ask.takenOver || ask.pending > 0 } /** Whether the turn is a prompt someone else submitted (consumed). */ private takeForeign(turnText: string): boolean { const value = turnText.trim() - const index = this.foreign.lastIndexOf(value) + // Another plugin's prompt reaches turn.start inside the engine's frame + // ("The plugin sent a message: ..."), while prompt.submit saw it bare. + const unframed = value.replace(PLUGIN_FRAME, '').trim() + let index = this.foreign.lastIndexOf(value) + if (index < 0 && unframed !== value) index = this.foreign.lastIndexOf(unframed) if (index < 0) return false this.foreign.splice(index, 1) return true } + /** + * Whether a turn started as a question's and is now someone else's (or a + * person's prompt is on its way into it): Plannotator never aborts it. + */ + isTakenOver(turnId: string): boolean { + return this.takenOverTurns.has(turnId) || (this.askOn(turnId)?.pending ?? 0) > 0 + } + get busy(): boolean { return this.runningTurnId !== null || (this.ask !== null && !this.ask.finished) } @@ -73,7 +205,20 @@ export class TurnTracker { /** Register a question about to be submitted. False when one is already in flight. */ beginAsk(askId: string, text: string, sink: AskSink): boolean { if (this.askInFlight) return false - this.ask = { askId, text, sink, turnId: null, cancelled: false, finished: false } + this.ask = { + askId, + text, + sink, + turnId: null, + cancelled: false, + finished: false, + takenOver: false, + byPerson: false, + pending: 0, + held: [], + streamed: '', + lastStop: null, + } return true } @@ -101,23 +246,80 @@ export class TurnTracker { } onText(turnId: string, text: string): void { - if (this.ownsTurn(turnId) && text) this.ask?.sink.delta(text) + const ask = this.ask + if (!ask || !this.ownsTurn(turnId) || !text) return + if (this.holding(ask)) ask.held.push({ kind: 'text', value: text }) + else this.forward(ask, text) } onTool(turnId: string, name: string): void { - if (this.ownsTurn(turnId) && name) this.ask?.sink.tool(name) + const ask = this.ask + if (!ask || !this.ownsTurn(turnId) || !name) return + if (this.holding(ask)) ask.held.push({ kind: 'tool', value: name }) + else ask.sink.tool(name) + } + + /** + * A model request of `turnId` is about to go out (turn.step). After a + * take-over it carries the other prompt, so the question settles here. + */ + onStep(turnId: string): void { + const ask = this.askOn(turnId) + if (!ask) return + const lastStop = ask.lastStop + ask.lastStop = null + if (!ask.takenOver) return + ask.finished = true + this.releaseHeld(ask) + // The response before this request ended the answer (no tool call): the + // engine runs on only for their prompt, and the question was answered whole. + // A thinking-only final response streamed nothing: an empty `done` would + // show nothing at all, so that case reads as taken over (as on Pi / OpenCode). + if (lastStop === 'end_turn' && ask.streamed.trim()) ask.sink.done(ask.streamed) + else ask.sink.error('taken_over', ask.byPerson ? TAKEN_OVER_BY_PERSON_TEXT : TAKEN_OVER_TEXT) + } + + /** A model response of `turnId` finished (the turn.step `stop` chunk). */ + onStepStop(turnId: string, stopReason: string | null): void { + const ask = this.askOn(turnId) + if (ask) ask.lastStop = stopReason } onTurnComplete(turnId: string, answer: string, aborted: boolean): void { if (this.runningTurnId === turnId) this.runningTurnId = null + this.takenOverTurns.delete(turnId) const ask = this.ask if (!ask || ask.finished || ask.turnId !== turnId) return ask.finished = true + if (ask.takenOver) { + // No step ran after the other prompt entered: everything the turn said + // answered the question (their prompt runs as a turn of its own). + this.releaseHeld(ask) + if (aborted || !answer.trim()) ask.sink.error('taken_over', ask.byPerson ? TAKEN_OVER_BY_PERSON_TEXT : TAKEN_OVER_TEXT) + else ask.sink.done(answer) + return + } + // A prompt still on its way never entered this turn: what was held is the answer. + this.releaseHeld(ask) if (aborted || ask.cancelled) ask.sink.error('aborted') else if (answer.trim()) ask.sink.done(answer) else ask.sink.error('failed', 'Claude finished the turn without a text answer.') } + private releaseHeld(ask: ActiveAsk): void { + const held = ask.held + ask.held = [] + for (const item of held) { + if (item.kind === 'text') this.forward(ask, item.value) + else ask.sink.tool(item.value) + } + } + + private forward(ask: ActiveAsk, text: string): void { + ask.streamed += text + ask.sink.delta(text) + } + /** * Cancel OUR question. Returns the turn id to abort when its turn is * running; a question still queued is dropped when its turn starts. @@ -126,6 +328,13 @@ export class TurnTracker { const ask = this.ask if (!ask || ask.askId !== askId || ask.finished) return null ask.cancelled = true + if (ask.takenOver || ask.pending > 0) { + // The turn carries someone else's prompt now (or is about to): close the question, never the turn. + ask.finished = true + ask.held = [] + ask.sink.error('aborted') + return null + } if (ask.turnId) return ask.turnId // Still queued behind another turn: confirm now, abort its turn when it starts. ask.finished = true @@ -143,6 +352,9 @@ export class TurnTracker { } } +/** The engine's frame around a plugin's prompt at turn.start. */ +const PLUGIN_FRAME = /^The [^\n]{1,200}? plugin sent a message:\s*/ + /** * Whether a turn's text is our question. The engine may wrap a plugin's * prompt in a "sent a message" frame, expand pastes or trim, so the question's diff --git a/apps/opencode-plugin/opencode-session-bridge.test.ts b/apps/opencode-plugin/opencode-session-bridge.test.ts index e4e104728..ddbd22cfb 100644 --- a/apps/opencode-plugin/opencode-session-bridge.test.ts +++ b/apps/opencode-plugin/opencode-session-bridge.test.ts @@ -4,9 +4,15 @@ import { createOpenCodeSessionBridge, markPlanReviewPending, readAnswerAfter, + TAKEN_OVER_INTERRUPT_TEXT, + TAKEN_OVER_TEXT, type OpenCodeSessionBridge, } from "./opencode-session-bridge"; -import type { SessionBridgeErrorCode } from "@plannotator/ai/session-bridge"; +import { + SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT, + SESSION_ASK_TAKEN_OVER_TEXT, + type SessionBridgeErrorCode, +} from "@plannotator/ai/session-bridge"; const SESSION = "ses_test"; @@ -249,6 +255,137 @@ describe("OpenCode session bridge", () => { await waitFor(() => host.interrupts === 1); }); + // The failure these guard: the person typed into the OpenCode run answering + // the reviewer's question, and the reply to THEIR prompt streamed into + // Plannotator as the answer, and a Plannotator Stop interrupted their work. + test("a prompt delivered into the run answering our question takes it over: streaming stops, Stop and interrupt leave the run alone", async () => { + const host = fakeHost(); + const bridge = bridgeFor(host); + const sink = recordingSink(); + const controller = new AbortController(); + bridge.ask({ askId: "a1", text: "q", mode: "turn" }, sink.sink, controller.signal); + await waitFor(() => host.prompts.length === 1); + host.emit("session.execution.started"); + host.emit("session.inbox.delivered", { inboxID: host.prompts[0].id }); + host.emit("session.step.started", { assistantMessageID: "m1" }); + host.emit("session.text.started", { assistantMessageID: "m1", ordinal: 0 }); + host.emit("session.text.delta", { assistantMessageID: "m1", ordinal: 0, delta: "Because " }); + // The person steers a prompt of their own into the running execution. + host.emit("session.inbox.enqueued", { inboxID: "msg_person", item: { type: "user", delivery: "steer", payload: {} } }); + host.emit("session.inbox.delivered", { inboxID: "msg_person" }); + host.emit("session.step.started", { assistantMessageID: "m2" }); + host.emit("session.text.started", { assistantMessageID: "m2", ordinal: 0 }); + host.emit("session.text.delta", { assistantMessageID: "m2", ordinal: 0, delta: "Fixed the tests." }); + + await waitFor(() => sink.error !== undefined); + expect(sink.error).toEqual({ code: "taken_over", message: TAKEN_OVER_TEXT }); + expect(sink.deltas).toEqual(["Because "]); + + controller.abort(); + await expect(Promise.resolve(bridge.interrupt?.())).rejects.toThrow(TAKEN_OVER_INTERRUPT_TEXT); + await new Promise((resolve) => setTimeout(resolve, 20)); + expect(host.interrupts).toBe(0); + + // Once that execution ends, interrupting the session works again. + host.emit("session.execution.succeeded"); + host.emit("session.execution.started"); + await bridge.interrupt?.(); + expect(host.interrupts).toBe(1); + }); + + test("a row promoted in the same batch as our question (a command's notice) is not a take-over", async () => { + const host = fakeHost(); + const bridge = bridgeFor(host); + const sink = recordingSink(); + bridge.ask({ askId: "a1", text: "q", mode: "turn" }, sink.sink, new AbortController().signal); + await waitFor(() => host.prompts.length === 1); + host.emit("session.execution.started"); + // Even when the notices are user rows, the same batch as ours takes nothing over. + for (const id of ["msg_notice", "msg_notice_after"]) { + host.emit("session.inbox.enqueued", { inboxID: id, item: { type: "user", delivery: "steer", payload: {} } }); + } + host.emit("session.inbox.delivered", { inboxID: "msg_notice" }); + host.emit("session.inbox.delivered", { inboxID: host.prompts[0].id }); + host.emit("session.inbox.delivered", { inboxID: "msg_notice_after" }); + host.emit("session.step.started", { assistantMessageID: "m1" }); + host.emit("session.text.started", { assistantMessageID: "m1", ordinal: 0 }); + host.emit("session.text.delta", { assistantMessageID: "m1", ordinal: 0, delta: "Because of X." }); + host.emit("session.execution.succeeded"); + await waitFor(() => sink.done !== undefined); + expect(sink.done).toBe("Because of X."); + expect(sink.error).toBeUndefined(); + }); + + // The failure this guards: a synthetic notice or a compaction delivered + // mid-answer cut the answer off with a take-over note. + test("only a USER row takes the run over: synthetic, compaction, move and unseen rows do not", async () => { + for (const kind of ["synthetic", "compaction", "move", null]) { + const host = fakeHost(); + const bridge = bridgeFor(host); + const sink = recordingSink(); + bridge.ask({ askId: "a1", text: "q", mode: "turn" }, sink.sink, new AbortController().signal); + await waitFor(() => host.prompts.length === 1); + host.emit("session.execution.started"); + host.emit("session.inbox.delivered", { inboxID: host.prompts[0].id }); + host.emit("session.step.started", { assistantMessageID: "m1" }); + host.emit("session.text.delta", { assistantMessageID: "m1", ordinal: 0, delta: "Because " }); + if (kind) host.emit("session.inbox.enqueued", { inboxID: "msg_row", item: { type: kind, delivery: "steer", payload: {} } }); + host.emit("session.inbox.delivered", { inboxID: "msg_row" }); + host.emit("session.text.delta", { assistantMessageID: "m1", ordinal: 0, delta: "of X." }); + host.emit("session.execution.succeeded"); + await waitFor(() => sink.done !== undefined || sink.error !== undefined); + expect([kind, sink.done, sink.error]).toEqual([kind, "Because of X.", undefined]); + } + }); + + test("a user row arriving after the answer finished (last step stopped, no tool calls) settles the answer as done", async () => { + const host = fakeHost(); + const bridge = bridgeFor(host); + const sink = recordingSink(); + bridge.ask({ askId: "a1", text: "q", mode: "turn" }, sink.sink, new AbortController().signal); + await waitFor(() => host.prompts.length === 1); + host.emit("session.execution.started"); + host.emit("session.inbox.delivered", { inboxID: host.prompts[0].id }); + host.emit("session.step.started", { assistantMessageID: "m1" }); + host.emit("session.text.delta", { assistantMessageID: "m1", ordinal: 0, delta: "Done." }); + host.emit("session.step.ended", { assistantMessageID: "m1", finish: "stop" }); + host.emit("session.inbox.enqueued", { inboxID: "msg_next", item: { type: "user", delivery: "queue", payload: {} } }); + host.emit("session.inbox.delivered", { inboxID: "msg_next" }); + host.emit("session.step.started", { assistantMessageID: "m2" }); + host.emit("session.text.delta", { assistantMessageID: "m2", ordinal: 0, delta: "Next thing." }); + await waitFor(() => sink.done !== undefined); + expect(sink.done).toBe("Done."); + expect(sink.error).toBeUndefined(); + expect(sink.deltas).toEqual(["Done."]); + }); + + // The failure this guards: an OpenAI-compatible provider reporting "stop" + // on a step that called tools made a cut answer read as finished. + test('a step that called tools is never a finished answer, even when its finish says "stop"', async () => { + const host = fakeHost(); + const bridge = bridgeFor(host); + const sink = recordingSink(); + bridge.ask({ askId: "a1", text: "q", mode: "turn" }, sink.sink, new AbortController().signal); + await waitFor(() => host.prompts.length === 1); + host.emit("session.execution.started"); + host.emit("session.inbox.delivered", { inboxID: host.prompts[0].id }); + host.emit("session.step.started", { assistantMessageID: "m1" }); + host.emit("session.text.delta", { assistantMessageID: "m1", ordinal: 0, delta: "Let me look." }); + host.emit("session.tool.input.started", { assistantMessageID: "m1", id: "t1", name: "read" }); + host.emit("session.step.ended", { assistantMessageID: "m1", finish: "stop" }); + host.emit("session.inbox.enqueued", { inboxID: "msg_person", item: { type: "user", delivery: "steer", payload: {} } }); + host.emit("session.inbox.delivered", { inboxID: "msg_person" }); + await waitFor(() => sink.error !== undefined || sink.done !== undefined); + expect(sink.done).toBeUndefined(); + expect(sink.error).toEqual({ code: "taken_over", message: TAKEN_OVER_TEXT }); + }); + + test("the plugin sends the same take-over texts the provider would", () => { + // Spelled out in the plugin so an older CLI server still shows them. + expect(TAKEN_OVER_TEXT).toBe(SESSION_ASK_TAKEN_OVER_TEXT); + expect(TAKEN_OVER_INTERRUPT_TEXT).toBe(SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT); + }); + test("refuses an interrupt while the session waits on a plan review", async () => { const host = fakeHost(); const bridge = bridgeFor(host); diff --git a/apps/opencode-plugin/opencode-session-bridge.ts b/apps/opencode-plugin/opencode-session-bridge.ts index c3108ccb8..aafc26483 100644 --- a/apps/opencode-plugin/opencode-session-bridge.ts +++ b/apps/opencode-plugin/opencode-session-bridge.ts @@ -35,6 +35,21 @@ * - `session.interrupt` stops the WHOLE execution. It is only called for a * turn that is ours, or when the reviewer chose "Interrupt and ask now", and * never while the session waits on a plan review (status `blocked`). + * - Take-over: once our question is delivered and the model has started + * answering it (`session.step.started`, or any answer output), a + * `session.inbox.delivered` for another USER row is a prompt we did not send + * (the person steering, a queued prompt promoted mid-run) entering the run, + * so the rest of the run answers it. The delivered event carries only the + * row id, so each row's kind is read from `session.inbox.enqueued` + * (`item.type`: user | synthetic | compaction | move); a synthetic notice, + * a compaction, a move, or a row whose kind this bridge never saw takes + * nothing over. Rows promoted in the same batch as ours (a command's + * session-URL notice) are delivered before the model starts and take + * nothing over either. On a take-over the question settles at once: `done` + * when the last model step had finished the answer (`finish: "stop"` AND no + * `session.tool.input.started` in that step, since some OpenAI-compatible + * providers report "stop" on a tool-calling step), else `taken_over` (deltas already sent stand). Neither a Stop + * nor "Interrupt and ask now" interrupts that execution afterwards. */ import type { @@ -128,8 +143,30 @@ interface ActiveTurn { streamed: Set; needsSeparator: boolean; watchdog: ReturnType | null; + /** The model began answering after our row was delivered: a later delivery is someone else's prompt. */ + answering: boolean; + /** How the last model step ended (`stop`: the answer was complete, no tool calls). */ + lastFinish: string | undefined; + /** The current (or last) model step started a tool call: never a finished answer, whatever `finish` says. */ + stepCalledTools: boolean; } +/** + * Sent with `taken_over` (a user row can be the person or another plugin's + * prompt, so the wording is neutral). Equal to `SESSION_ASK_TAKEN_OVER_TEXT` + * (packages/ai/session-bridge.ts; a test holds them together), spelled out to + * keep this module's imports type-only. + */ +export const TAKEN_OVER_TEXT = + "Another message entered this session while it was answering, so the rest of the reply went to that message."; + +/** Why "Interrupt and ask now" refuses a taken-over execution. Equal to `SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT`. */ +export const TAKEN_OVER_INTERRUPT_TEXT = + "The session is now answering another message, so Plannotator will not stop it. Ask when it finishes instead."; + +/** Inbox row kinds remembered per session (`session.inbox.enqueued`), bounded. */ +const MAX_REMEMBERED_ROWS = 256; + function isRecord(value: unknown): value is Record { return !!value && typeof value === "object"; } @@ -180,6 +217,10 @@ export function createOpenCodeSessionBridge(options: OpenCodeSessionBridgeOption const controller = new AbortController(); let probeTimer: ReturnType | null = null; let probing = false; + /** The execution that answered our question was taken over: never interrupted from Plannotator. */ + let takenOverRun = false; + /** `item.type` of each inbox row seen enqueued (user | synthetic | compaction | move). */ + const rowKinds = new Map(); const status = (): SessionBridgeStatus => { if (gone || disposed) return "gone"; @@ -320,18 +361,55 @@ export function createOpenCodeSessionBridge(options: OpenCodeSessionBridgeOption lastStartedAt = Date.now(); running = true; return; + case "session.inbox.enqueued": { + const item = isRecord(data.item) ? data.item : undefined; + if (typeof data.inboxID === "string" && typeof item?.type === "string") { + rowKinds.set(data.inboxID, item.type); + if (rowKinds.size > MAX_REMEMBERED_ROWS) rowKinds.delete(rowKinds.keys().next().value!); + } + return; + } case "session.inbox.delivered": if (turn && data.inboxID === turn.messageID) { turn.delivered = true; clearWatchdog(turn); // Stopped before it reached the model: stop it now that it is ours. if (turn.cancelled) void session?.interrupt?.({ sessionID }).catch(() => {}); + } else if ( + turn?.delivered && + turn.answering && + !turn.cancelled && + typeof data.inboxID === "string" && + rowKinds.get(data.inboxID) === "user" + ) { + // Someone else's prompt entered the run answering our question. + takenOverRun = true; + // Some OpenAI-compatible providers report "stop" on a step that + // called tools, so a tool call in the step rules it out too. + const complete = turn.lastFinish === "stop" && !turn.stepCalledTools && !!turn.answer; + finishTurn(turn, () => (complete ? turn.sink.done(turn.answer) : turn.sink.error("taken_over", TAKEN_OVER_TEXT))); + } + if (typeof data.inboxID === "string") rowKinds.delete(data.inboxID); + return; + case "session.step.started": + if (turn?.delivered) { + turn.answering = true; + turn.lastFinish = undefined; + turn.stepCalledTools = false; } return; + case "session.step.ended": + if (turn?.delivered && typeof data.finish === "string") turn.lastFinish = data.finish; + return; + case "session.reasoning.started": + if (turn?.delivered) turn.answering = true; + return; case "session.text.started": + if (turn?.delivered) turn.answering = true; if (turn?.delivered && turn.answer) turn.needsSeparator = true; return; case "session.text.delta": + if (turn?.delivered) turn.answering = true; if (turn?.delivered && typeof data.delta === "string") { turn.streamed.add(`${String(data.assistantMessageID)}:${String(data.ordinal)}`); appendText(turn, data.delta); @@ -345,12 +423,17 @@ export function createOpenCodeSessionBridge(options: OpenCodeSessionBridgeOption } return; case "session.tool.input.started": + if (turn?.delivered) { + turn.answering = true; + turn.stepCalledTools = true; + } if (turn?.delivered && !turn.cancelled && typeof data.name === "string") turn.sink.tool?.(data.name); return; case "session.execution.succeeded": case "session.execution.failed": case "session.execution.interrupted": { running = false; + takenOverRun = false; if (!turn?.delivered) return; if (turn.cancelled) { finishTurn(turn, () => turn.sink.error("aborted")); @@ -399,6 +482,9 @@ export function createOpenCodeSessionBridge(options: OpenCodeSessionBridgeOption streamed: new Set(), needsSeparator: false, watchdog: null, + answering: false, + lastFinish: undefined, + stepCalledTools: false, }; active = turn; @@ -528,6 +614,8 @@ export function createOpenCodeSessionBridge(options: OpenCodeSessionBridgeOption // The provider never asks while blocked, but a plan review may have // started since: interrupting now would kill it. if (isPlanReviewPending(sessionID)) throw new Error("The session is waiting on a plan review."); + // The person typed into the run that answered a question: it is theirs. + if (takenOverRun && running) throw new Error(TAKEN_OVER_INTERRUPT_TEXT); const interrupt = session?.interrupt; if (typeof interrupt !== "function") throw new Error("This OpenCode host cannot interrupt a session from a plugin."); await interrupt({ sessionID }); diff --git a/apps/pi-extension/pi-session-bridge.test.ts b/apps/pi-extension/pi-session-bridge.test.ts index fe4052668..dcb72fb25 100644 --- a/apps/pi-extension/pi-session-bridge.test.ts +++ b/apps/pi-extension/pi-session-bridge.test.ts @@ -9,7 +9,12 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import { request as httpRequest } from "node:http"; import { createPiSessionBridgeHub, PLANNOTATOR_ASK_CUSTOM_TYPE } from "./pi-session-bridge.ts"; -import { SESSION_ASK_HEADER, SessionBridgeProvider } from "./generated/ai/session-bridge.ts"; +import { + SESSION_ASK_HEADER, + SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT, + SESSION_ASK_TAKEN_OVER_TEXT, + SessionBridgeProvider, +} from "./generated/ai/session-bridge.ts"; import type { AIMessage } from "./generated/ai/types.ts"; import { startAnnotateServer } from "./server/serverAnnotate.ts"; @@ -63,6 +68,7 @@ function fakePi() { startTurn(askId: string) { idle = false; emit("agent_start", { type: "agent_start" }); + emit("turn_start", { type: "turn_start" }); emit("message_start", { type: "message_start", message: { role: "custom", customType: PLANNOTATOR_ASK_CUSTOM_TYPE, details: { askId } }, @@ -78,6 +84,11 @@ function fakePi() { }); } }, + /** The loop finishes one assistant turn (tool results in) and starts the next. */ + nextTurn() { + emit("turn_end", { type: "turn_end" }); + emit("turn_start", { type: "turn_start" }); + }, /** One agent run ends; Pi may still retry it before the prompt settles. */ endRun(stopReason = "stop", errorMessage?: string) { emit("agent_end", { type: "agent_end", messages: [{ role: "assistant", stopReason, errorMessage, content: [] }] }); @@ -253,6 +264,148 @@ describe("Pi session bridge", () => { expect(first.calls.at(-1)?.slice(0, 2)).toEqual(["error", "gone"]); }); + // The failure these guard: the person typed into Pi while it answered the + // reviewer's question, and the reply to THEIR prompt streamed into + // Plannotator as the answer, and a Plannotator Stop aborted their work. + test("a message the person steers into our run takes it over: streaming stops, Stop and interrupt leave the run alone", () => { + const host = fakePi(); + const bridge = createPiSessionBridgeHub(host.pi).createBridge(host.ctx as never, {}); + const out = sink(); + const controller = new AbortController(); + bridge.ask({ askId: "ask-1", text: "Question text", mode: "turn" }, out.sink, controller.signal); + host.startTurn("ask-1"); + host.assistantText("Because "); + // The person steers; the loop delivers it right after the next turn_start. + host.nextTurn(); + host.emit("message_start", { type: "message_start", message: { role: "user", content: "also fix the tests" } }); + host.assistantText("Fixed the tests."); + + expect(out.calls).toEqual([ + ["delta", "Because "], + ["error", "taken_over", undefined], + ]); + controller.abort(); + expect(host.aborts).toHaveLength(0); + expect(() => bridge.interrupt?.()).toThrow(SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT); + expect(host.aborts).toHaveLength(0); + + // Once that run ends, interrupting the session works again. + host.endTurn(); + host.setIdle(false); + void bridge.interrupt?.(); + expect(host.aborts).toHaveLength(1); + }); + + test("a triggering custom message from another extension, steered into our run, takes it over", () => { + const host = fakePi(); + const bridge = createPiSessionBridgeHub(host.pi).createBridge(host.ctx as never, {}); + const out = sink(); + bridge.ask({ askId: "ask-1", text: "Question text", mode: "turn" }, out.sink, new AbortController().signal); + host.startTurn("ask-1"); + host.assistantText("Because "); + host.nextTurn(); + host.emit("message_start", { type: "message_start", message: { role: "custom", customType: "other-ext", details: {} } }); + expect(out.calls.at(-1)).toEqual(["error", "taken_over", undefined]); + }); + + // The failure this guards: Plannotator's own decision follow-up (or any + // follow-up) arriving after the answer had finished turned a complete + // answer into a "taken over" one. + test("a follow-up arriving after the answer finished settles the answer as done", () => { + const host = fakePi(); + const bridge = createPiSessionBridgeHub(host.pi).createBridge(host.ctx as never, {}); + const out = sink(); + bridge.ask({ askId: "ask-1", text: "Question text", mode: "turn" }, out.sink, new AbortController().signal); + host.startTurn("ask-1"); + host.assistantText("Done."); + // The answer's last turn: a stop with no tool call. The loop then pulls a follow-up. + host.emit("turn_end", { type: "turn_end", message: { role: "assistant", stopReason: "stop", content: [{ type: "text" }] }, toolResults: [] }); + host.emit("turn_start", { type: "turn_start" }); + host.emit("message_start", { type: "message_start", message: { role: "user", content: "plan approved" } }); + host.assistantText("Implementing."); + expect(out.calls).toEqual([ + ["delta", "Done."], + ["done", "Done."], + ]); + }); + + test("a turn that called tools before the steer is not a finished answer", () => { + const host = fakePi(); + const bridge = createPiSessionBridgeHub(host.pi).createBridge(host.ctx as never, {}); + const out = sink(); + bridge.ask({ askId: "ask-1", text: "Question text", mode: "turn" }, out.sink, new AbortController().signal); + host.startTurn("ask-1"); + host.assistantText("Let me look."); + host.emit("turn_end", { type: "turn_end", message: { role: "assistant", stopReason: "toolUse", content: [{ type: "toolCall" }] }, toolResults: [{}] }); + host.emit("turn_start", { type: "turn_start" }); + host.emit("message_start", { type: "message_start", message: { role: "user", content: "stop" } }); + expect(out.calls.at(-1)).toEqual(["error", "taken_over", undefined]); + }); + + test("the taken-over run stays protected across an error retry until it settles", () => { + const host = fakePi(); + const bridge = createPiSessionBridgeHub(host.pi).createBridge(host.ctx as never, {}); + bridge.ask({ askId: "ask-1", text: "Question text", mode: "turn" }, sink().sink, new AbortController().signal); + host.startTurn("ask-1"); + host.assistantText("Because "); + host.nextTurn(); + host.emit("message_start", { type: "message_start", message: { role: "user", content: "also fix the tests" } }); + // The person's run errors; Pi retries it (agent.continue) before settling. + host.endRun("error", "overloaded"); + expect(() => bridge.interrupt?.()).toThrow(SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT); + host.emit("agent_start", { type: "agent_start" }); + expect(() => bridge.interrupt?.()).toThrow(SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT); + host.endRun(); + host.settle(); + host.setIdle(false); + void bridge.interrupt?.(); + expect(host.aborts).toHaveLength(1); + }); + + test("a display-only custom message appended at turn_end is not a take-over", () => { + const host = fakePi(); + const bridge = createPiSessionBridgeHub(host.pi).createBridge(host.ctx as never, {}); + const out = sink(); + bridge.ask({ askId: "ask-1", text: "Question text", mode: "turn" }, out.sink, new AbortController().signal); + host.startTurn("ask-1"); + host.assistantText("Because "); + // Pi appends a non-triggering custom message (e.g. plannotator-handoff) + // while handling turn_end, before the next turn_start. + host.emit("turn_end", { type: "turn_end" }); + host.emit("message_start", { type: "message_start", message: { role: "custom", customType: "plannotator-handoff", details: {} } }); + host.emit("turn_start", { type: "turn_start" }); + host.assistantText("of X."); + host.endTurn(); + expect(out.calls).toEqual([ + ["delta", "Because "], + ["delta", "\n\nof X."], + ["done", "Because \n\nof X."], + ]); + }); + + test("through the provider: a take-over keeps the partial answer and ends with the note", async () => { + const host = fakePi(); + const bridge = createPiSessionBridgeHub(host.pi).createBridge(host.ctx as never, {}); + const provider = new SessionBridgeProvider(bridge, { pollIntervalMs: 5 }); + const session = await provider.createSession({ + context: { mode: "annotate", annotate: { content: "", filePath: "last-message" } }, + }); + const messages: AIMessage[] = []; + const done = (async () => { + for await (const message of session.query("What did you mean?")) messages.push(message); + })(); + while (host.sent.length === 0) await new Promise((resolve) => setTimeout(resolve, 2)); + host.startTurn(host.sent[0].message.details!.askId!); + host.assistantText("I meant "); + host.nextTurn(); + host.emit("message_start", { type: "message_start", message: { role: "user", content: "never mind" } }); + await done; + expect(messages).toEqual([ + { type: "text_delta", delta: "I meant " }, + { type: "error", code: "session_taken_over", error: SESSION_ASK_TAKEN_OVER_TEXT }, + ]); + }); + test("interrupt stops the session's own turn", () => { const host = fakePi(); const bridge = createPiSessionBridgeHub(host.pi).createBridge(host.ctx as never, {}); diff --git a/apps/pi-extension/pi-session-bridge.ts b/apps/pi-extension/pi-session-bridge.ts index 65d6f815c..bf9c13f5c 100644 --- a/apps/pi-extension/pi-session-bridge.ts +++ b/apps/pi-extension/pi-session-bridge.ts @@ -23,6 +23,20 @@ * a continuation after an overflow compaction. An error end therefore waits * for `agent_settled` (Pi >= 0.80.4), or for the session to go idle on older * Pi, before it is reported; a retry that starts answering clears it. + * - Take-over: a message we did not send that enters the run answering our + * question makes the rest of that run someone else's. Pi emits such a + * message's `message_start` from the agent loop right after a `turn_start`, + * before the next assistant message: a `user` message (the person steering + * or following up, or `sendUserMessage`) or a triggering custom message + * (steered). A NON-triggering custom message sent while the run streams + * (display-only, e.g. `plannotator-handoff`) is appended at `turn_end` + * instead, outside that window, and takes nothing over. On a take-over the + * question settles at once: `done` when the turn before it ended the answer + * (stop, no tool call; e.g. Plannotator's own decision follow-up arriving + * after the answer finished), `taken_over` otherwise (the deltas already + * sent stand). From then on until the run settles (`agent_settled`, or idle + * on older Pi, so an error retry stays covered), neither a Stop nor + * "Interrupt and ask now" aborts that run: it answers someone else. * * Only type imports here: this module is loaded eagerly by index.ts. */ @@ -100,6 +114,22 @@ export function createPiSessionBridgeHub( ): PiSessionBridgeHub { const startWatchdogMs = options.startWatchdogMs ?? START_WATCHDOG_MS; let active: ActiveAsk | null = null; + /** Between a `turn_start` and that turn's assistant message: where the loop delivers queued messages. */ + let turnPrelude = false; + /** The run that answered our question was taken over: never aborted from Plannotator until it settles. */ + let takenOverRun = false; + /** Older Pi (no `agent_settled`): polls idleness after the taken-over run ends. */ + let takenOverPoll: ReturnType | null = null; + let takenOverIsIdle: (() => boolean) | null = null; + /** The last turn of the run ended the answer: a stop with no tool call. */ + let lastTurnFinal = false; + + const clearTakenOver = () => { + takenOverRun = false; + takenOverIsIdle = null; + if (takenOverPoll) clearInterval(takenOverPoll); + takenOverPoll = null; + }; const finish = (ask: ActiveAsk) => { clearWatchdog(ask); @@ -114,18 +144,47 @@ export function createPiSessionBridgeHub( ask.sink.error("failed", message); }; + pi.on("turn_start", () => { + turnPrelude = true; + }); + + pi.on("turn_end", (event) => { + turnPrelude = false; + const message = event?.message as { stopReason?: string; content?: unknown } | undefined; + const calledTools = + (Array.isArray(event?.toolResults) && event.toolResults.length > 0) || + (Array.isArray(message?.content) && message.content.some((part: { type?: string } | null) => part?.type === "toolCall")); + lastTurnFinal = message?.stopReason === "stop" && !calledTools; + }); + pi.on("message_start", (event) => { - const ask = active; - if (!ask) return; const message = event?.message as { role?: string; customType?: string; details?: { askId?: unknown } } | undefined; - if (!message) return; - if (message.role === "custom" && message.customType === PLANNOTATOR_ASK_CUSTOM_TYPE) { - if (message.details?.askId === ask.askId) { - ask.started = true; - clearWatchdog(ask); - } + const inPrelude = turnPrelude; + const answerFinished = lastTurnFinal; + if (message?.role === "assistant") { + turnPrelude = false; + lastTurnFinal = false; + } + const ask = active; + if (!ask || !message) return; + const ours = message.role === "custom" && message.customType === PLANNOTATOR_ASK_CUSTOM_TYPE && message.details?.askId === ask.askId; + if (ours) { + ask.started = true; + clearWatchdog(ask); return; } + if (ask.started && (message.role === "user" || (message.role === "custom" && inPrelude))) { + // Someone else's prompt entered the run answering our question: the + // rest of it answers them. Settle now and leave the run alone. + clearTakenOver(); + takenOverRun = true; + takenOverIsIdle = ask.isIdle; + finish(ask); + if (answerFinished && ask.answer) ask.sink.done(ask.answer); + else ask.sink.error("taken_over"); + return; + } + if (message.role === "custom" && message.customType === PLANNOTATOR_ASK_CUSTOM_TYPE) return; if (ask.started && message.role === "assistant") { // Pi retried (or continued after compaction): the run is still answering. ask.pendingError = null; @@ -155,6 +214,23 @@ export function createPiSessionBridgeHub( }); pi.on("agent_end", (event) => { + turnPrelude = false; + lastTurnFinal = false; + if (takenOverRun && !takenOverPoll) { + // Pi may retry this run (agent.continue): the taken-over flag holds + // until it settles. Without `agent_settled` (older Pi), until idle. + const isIdle = takenOverIsIdle; + takenOverPoll = setInterval(() => { + let idle = true; + try { + idle = isIdle ? isIdle() : true; + } catch { + idle = true; + } + if (idle) clearTakenOver(); + }, SETTLE_POLL_MS); + (takenOverPoll as { unref?: () => void }).unref?.(); + } const ask = active; if (!ask?.started) return; if (ask.cancelled) { @@ -187,6 +263,7 @@ export function createPiSessionBridgeHub( }); pi.on("agent_settled", () => { + clearTakenOver(); const ask = active; if (ask) reportPendingError(ask); }); @@ -294,6 +371,10 @@ export function createPiSessionBridgeHub( } }, interrupt() { + if (takenOverRun) { + // Same text as SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT (type-only imports here). + throw new Error("The session is now answering another message, so Plannotator will not stop it. Ask when it finishes instead."); + } ctx.abort(); }, }; diff --git a/packages/ai/index.ts b/packages/ai/index.ts index 43de72703..052fb10ce 100644 --- a/packages/ai/index.ts +++ b/packages/ai/index.ts @@ -100,6 +100,9 @@ export { SESSION_BRIDGE_ERROR, SESSION_ASK_HEADER, SESSION_ASK_TRANSIENT_NOTE, + SESSION_ASK_TAKEN_OVER_TEXT, + SESSION_ASK_TAKEN_OVER_BY_PERSON_TEXT, + SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT, formatSessionAskText, sessionBridgeLabel, } from "./session-bridge.ts"; @@ -125,6 +128,7 @@ export { SESSION_BRIDGE_HOST_ENV, SESSION_BRIDGE_MODES_ENV, SESSION_BRIDGE_MAX_POLL_MS, + SESSION_BRIDGE_POLL_FEATURES, } from "./session-bridge-pull.ts"; export type { BridgeCommand, diff --git a/packages/ai/session-bridge-pull-client.ts b/packages/ai/session-bridge-pull-client.ts index 548db4a4a..3366403dc 100644 --- a/packages/ai/session-bridge-pull-client.ts +++ b/packages/ai/session-bridge-pull-client.ts @@ -41,6 +41,22 @@ export interface PullSessionBridgeClientOptions { interface HostAsk { controller: AbortController; finished: boolean; + /** What streamed so far: the answer a `taken_over` fallback settles with. */ + text: string; +} + +/** Used when a host sends `taken_over` without a message (in-process hosts rely on the provider's). */ +const TAKEN_OVER_FALLBACK_NOTE = "Another message entered this session while it was answering, so the rest of the reply went to that message."; + +/** + * A server that does not advertise `taken_over` reads it as `failed`, and its + * UI then replaces the partial answer with the error. Settle as an answer + * instead: what streamed, plus the note as its last paragraph. + */ +export function takenOverFallback(streamed: string, message: string | undefined): { delta: string; answer: string } { + const note = `_${(message || TAKEN_OVER_FALLBACK_NOTE).trim()}_`; + const delta = streamed ? `\n\n${note}` : note; + return { delta, answer: `${streamed}${delta}` }; } function sleep(ms: number, signal: AbortSignal): Promise { @@ -78,6 +94,8 @@ export async function runPullSessionBridgeClient(options: PullSessionBridgeClien const asks = new Map(); const interrupts = new Set(); let stopped = false; + /** From the latest poll answer's `features` (SESSION_BRIDGE_POLL_FEATURES). */ + let serverTakesTakenOver = false; // One ordered outbox: events reach the server in the order they happened. let outbox: BridgeHostEvent[] = []; @@ -134,7 +152,7 @@ export async function runPullSessionBridgeClient(options: PullSessionBridgeClien const runAsk = (command: Extract) => { if (asks.has(command.askId)) return; const controller = new AbortController(); - const ask: HostAsk = { controller, finished: false }; + const ask: HostAsk = { controller, finished: false, text: "" }; asks.set(command.askId, ask); emit({ type: "started", askId: command.askId }); const finish = (event: BridgeHostEvent) => { @@ -147,14 +165,24 @@ export async function runPullSessionBridgeClient(options: PullSessionBridgeClien { askId: command.askId, text: command.text, mode: command.mode }, { delta: (text) => { - if (!ask.finished && text) emit({ type: "delta", askId: command.askId, text }, false); + if (ask.finished || !text) return; + ask.text += text; + emit({ type: "delta", askId: command.askId, text }, false); }, tool: (name) => { if (!ask.finished && name) emit({ type: "tool", askId: command.askId, name }); }, done: (answer) => finish({ type: "done", askId: command.askId, answer }), - error: (code: SessionBridgeErrorCode, message?: string) => - finish({ type: "error", askId: command.askId, code, ...(message ? { message } : {}) }), + error: (code: SessionBridgeErrorCode, message?: string) => { + if (code === "taken_over" && !serverTakesTakenOver) { + if (ask.finished) return; + const fallback = takenOverFallback(ask.text, message); + emit({ type: "delta", askId: command.askId, text: fallback.delta }, false); + finish({ type: "done", askId: command.askId, answer: fallback.answer }); + return; + } + finish({ type: "error", askId: command.askId, code, ...(message ? { message } : {}) }); + }, }, controller.signal, ); @@ -249,12 +277,13 @@ export async function runPullSessionBridgeClient(options: PullSessionBridgeClien continue; } failures = 0; - let body: { commands?: BridgeCommand[]; closing?: boolean } = {}; + let body: { commands?: BridgeCommand[]; closing?: boolean; features?: unknown } = {}; try { body = (await res.json()) as typeof body; } catch { // Treat as an empty answer. } + serverTakesTakenOver = Array.isArray(body.features) && body.features.includes("taken_over"); for (const command of Array.isArray(body.commands) ? body.commands : []) { if (command && typeof command === "object" && typeof command.type === "string") handleCommand(command); } diff --git a/packages/ai/session-bridge-pull.test.ts b/packages/ai/session-bridge-pull.test.ts index ac1892308..beae06a3e 100644 --- a/packages/ai/session-bridge-pull.test.ts +++ b/packages/ai/session-bridge-pull.test.ts @@ -20,7 +20,7 @@ import { SESSION_BRIDGE_POLL_PATH, type PullSessionBridgeOptions, } from "./session-bridge-pull.ts"; -import { runPullSessionBridgeClient } from "./session-bridge-pull-client.ts"; +import { runPullSessionBridgeClient, takenOverFallback } from "./session-bridge-pull-client.ts"; import type { AIContext, AIMessage } from "./types.ts"; const TOKEN = "t".repeat(43); @@ -87,7 +87,7 @@ afterEach(() => { * A Plannotator server's AI runtime with a pull bridge, plus a host process * simulated by the real pull client talking to it through `fetch`. */ -function setup(options: Partial = {}, host = fakeHost()) { +function setup(options: Partial = {}, host = fakeHost(), wire: { olderServer?: boolean } = {}) { const pull = createPullSessionBridge({ token: TOKEN, host: "opencode", @@ -110,8 +110,14 @@ function setup(options: Partial = {}, host = fakeHost( (endpoints as Record Promise>)[path]( new Request(`http://${HOST}${path}`, { ...init, headers: { host: HOST, ...(init.headers as Record), ...headers } }), ); - const fetchShim = (async (url: string | URL | Request, init?: RequestInit) => - call(new URL(String(url)).pathname, init ?? {})) as typeof fetch; + const fetchShim = (async (url: string | URL | Request, init?: RequestInit) => { + const path = new URL(String(url)).pathname; + const res = await call(path, init ?? {}); + if (!wire.olderServer || path !== SESSION_BRIDGE_POLL_PATH) return res; + // A server from before the `features` advert. + const { features: _features, ...rest } = (await res.json()) as Record; + return Response.json(rest, { status: res.status }); + }) as typeof fetch; const controller = new AbortController(); const startHost = (overrides: { token?: string } = {}) => runPullSessionBridgeClient({ @@ -157,6 +163,39 @@ describe("pull session bridge", () => { expect(run.messages.at(-1)).toMatchObject({ type: "result", success: true, result: "Because of #12." }); }); + // The failure this guards: a host's `taken_over` read as an unknown code, + // so the reviewer saw "the session could not answer" instead of the note. + test("a host's taken_over crosses the pull bridge as session_taken_over with the partial answer kept", async () => { + const { provider, host, startHost } = setup(); + void startHost(); + const session = await provider.createSession({ context: CONTEXT }); + const run = collect(session.query("Why this change?")); + await waitFor(() => host.asks.length === 1); + host.asks[0].sink.delta("Because "); + host.asks[0].sink.error("taken_over", "the person typed"); + await run.done; + const text = run.messages.filter((m) => m.type === "text_delta").map((m) => (m as { delta: string }).delta).join(""); + expect(text).toBe("Because "); + expect(run.messages.at(-1)).toEqual({ type: "error", code: SESSION_BRIDGE_ERROR.takenOver, error: "the person typed" }); + }); + + // The failure this guards: a newer OpenCode plugin against an older CLI, + // whose UI would replace the partial answer with a "failed" error. + test("against a server that does not advertise taken_over, the host settles the partial answer plus the note", async () => { + const { provider, host, startHost } = setup({}, fakeHost(), { olderServer: true }); + void startHost(); + const session = await provider.createSession({ context: CONTEXT }); + const run = collect(session.query("Why this change?")); + await waitFor(() => host.asks.length === 1); + host.asks[0].sink.delta("Because "); + host.asks[0].sink.error("taken_over", "NOTE"); + await run.done; + const fallback = takenOverFallback("Because ", "NOTE"); + const text = run.messages.filter((m) => m.type === "text_delta").map((m) => (m as { delta: string }).delta).join(""); + expect(text).toBe(fallback.answer); + expect(run.messages.at(-1)).toMatchObject({ type: "result", success: true, result: fallback.answer }); + }); + test("a question asked before the host connects waits for its first poll", async () => { const { provider, host, startHost } = setup(); const session = await provider.createSession({ context: CONTEXT }); @@ -305,10 +344,10 @@ describe("pull session bridge", () => { const first = call(SESSION_BRIDGE_POLL_PATH, { method: "POST", body: JSON.stringify({ waitMs: 5_000 }) }, auth); await new Promise((resolve) => setTimeout(resolve, 10)); const second = call(SESSION_BRIDGE_POLL_PATH, { method: "POST", body: JSON.stringify({ waitMs: 5_000 }) }, auth); - expect(await (await first).json()).toEqual({ commands: [], superseded: true }); + expect(await (await first).json()).toEqual({ commands: [], superseded: true, features: ["taken_over"] }); await new Promise((resolve) => setTimeout(resolve, 10)); pull.dispose(); - expect(await (await second).json()).toEqual({ commands: [], closing: true }); + expect(await (await second).json()).toEqual({ commands: [], closing: true, features: ["taken_over"] }); expect(pull.bridge.status()).toBe("gone"); }); diff --git a/packages/ai/session-bridge-pull.ts b/packages/ai/session-bridge-pull.ts index 9389bd867..b68946c24 100644 --- a/packages/ai/session-bridge-pull.ts +++ b/packages/ai/session-bridge-pull.ts @@ -17,7 +17,8 @@ * POST /api/ai/bridge/poll body { status?, modes?, waitMs? } * Long-poll for commands. Answers as soon as there is something to do, or * with an empty list after `waitMs` (clamped to 0..25s, default 25s). - * -> 200 { commands: BridgeCommand[], closing?: true, superseded?: true } + * -> 200 { commands: BridgeCommand[], features: string[], closing?: true, superseded?: true } + * `features` lists optional protocol features (SESSION_BRIDGE_POLL_FEATURES). * Each poll is also the host's heartbeat and status report. * * POST /api/ai/bridge/event body BridgeHostEvent | { events: BridgeHostEvent[] } @@ -67,6 +68,18 @@ export const SESSION_BRIDGE_HOST_ENV = "PLANNOTATOR_SESSION_BRIDGE_HOST"; export const SESSION_BRIDGE_MODES_ENV = "PLANNOTATOR_SESSION_BRIDGE_MODES"; export const SESSION_BRIDGE_MAX_POLL_MS = 25_000; + +/** + * Optional protocol features this server understands, sent as `features` on + * every poll answer. A host checks them before relying on one; a server that + * predates the list sends none. + * - `taken_over`: the `taken_over` error code (the session's turn was taken + * over by another message). Without it a host settles such a question as + * `done` with the partial answer and the note appended, because an older + * server reads the unknown code as `failed` and its UI then replaces the + * partial answer with the error. + */ +export const SESSION_BRIDGE_POLL_FEATURES: readonly string[] = ["taken_over"]; const MAX_EVENT_BODY_BYTES = 1_000_000; export type BridgeCommand = @@ -112,13 +125,18 @@ export interface PullSessionBridge { } const STATUSES: ReadonlySet = new Set(["ready", "busy", "blocked", "gone"]); -const ERROR_CODES: ReadonlySet = new Set(["busy", "blocked", "gone", "aborted", "failed"]); +const ERROR_CODES: ReadonlySet = new Set(["busy", "blocked", "gone", "aborted", "failed", "taken_over"]); const BRIDGE_HOSTS: ReadonlySet = new Set(["pi", "opencode", "claude-code"]); function isRecord(value: unknown): value is Record { return !!value && typeof value === "object" && !Array.isArray(value); } +/** A poll answer: always advertises the protocol features this server knows. */ +function pollReply(body: Record): Response { + return json({ ...body, features: SESSION_BRIDGE_POLL_FEATURES }); +} + function json(body: unknown, status = 200): Response { return new Response(JSON.stringify(body), { status, headers: { "content-type": "application/json" } }); } @@ -283,7 +301,7 @@ export function createPullSessionBridge(options: PullSessionBridgeOptions): Pull openPoll = null; clearTimeout(poll.timer); lastSeen = now(); - poll.resolve(json({ commands })); + poll.resolve(pollReply({ commands })); }; const bridge: SessionBridge = { @@ -407,7 +425,7 @@ export function createPullSessionBridge(options: PullSessionBridgeOptions): Pull if (typeof body.modes.transient === "boolean") modes.transient = body.modes.transient; } } - if (disposed) return json({ commands: [], closing: true }); + if (disposed) return pollReply({ commands: [], closing: true }); if (hostStatus === "gone") failForGone(); const requested = isRecord(body) && typeof body.waitMs === "number" && Number.isFinite(body.waitMs) ? body.waitMs : SESSION_BRIDGE_MAX_POLL_MS; @@ -418,11 +436,11 @@ export function createPullSessionBridge(options: PullSessionBridgeOptions): Pull if (previous) { openPoll = null; clearTimeout(previous.timer); - previous.resolve(json({ commands: [], superseded: true })); + previous.resolve(pollReply({ commands: [], superseded: true })); } const commands = takeDueCommands(); - if (commands.length > 0 || waitMs === 0) return json({ commands }); + if (commands.length > 0 || waitMs === 0) return pollReply({ commands }); return await new Promise((resolve) => { const poll: OpenPoll = { @@ -431,7 +449,7 @@ export function createPullSessionBridge(options: PullSessionBridgeOptions): Pull if (openPoll !== poll) return; openPoll = null; lastSeen = now(); - resolve(json({ commands: takeDueCommands() })); + resolve(pollReply({ commands: takeDueCommands() })); }, waitMs), }; openPoll = poll; @@ -443,7 +461,7 @@ export function createPullSessionBridge(options: PullSessionBridgeOptions): Pull openPoll = null; clearTimeout(poll.timer); lastSeen = now(); - resolve(json({ commands: [] })); + resolve(pollReply({ commands: [] })); }, { once: true }, ); @@ -532,7 +550,7 @@ export function createPullSessionBridge(options: PullSessionBridgeOptions): Pull openPoll = null; if (poll) { clearTimeout(poll.timer); - poll.resolve(json({ commands: [], closing: true })); + poll.resolve(pollReply({ commands: [], closing: true })); } failForGone(); }, diff --git a/packages/ai/session-bridge.test.ts b/packages/ai/session-bridge.test.ts index 675b91c05..583915e54 100644 --- a/packages/ai/session-bridge.test.ts +++ b/packages/ai/session-bridge.test.ts @@ -5,6 +5,7 @@ import { BaseSession } from "./base-session.ts"; import { SessionManager } from "./session-manager.ts"; import { SESSION_ASK_HEADER, + SESSION_ASK_TAKEN_OVER_TEXT, SESSION_BRIDGE_ERROR, SESSION_BRIDGE_PROVIDER_NAME, SessionBridgeProvider, @@ -228,6 +229,20 @@ describe("SessionBridgeProvider", () => { expect(run.messages).toEqual([expect.objectContaining({ code: SESSION_BRIDGE_ERROR.gone })]); }); + test("a taken-over answer keeps what streamed and ends with session_taken_over and the note", async () => { + const host = fakeBridge("ready"); + const run = collect((await newSession(new SessionBridgeProvider(host.bridge))).query("q")); + await waitFor(() => host.asks.length === 1); + host.asks[0].sink.delta("Because "); + host.asks[0].sink.error("taken_over"); + host.asks[0].sink.delta("their reply"); + await run.done; + expect(run.messages).toEqual([ + { type: "text_delta", delta: "Because " }, + { type: "error", code: SESSION_BRIDGE_ERROR.takenOver, error: SESSION_ASK_TAKEN_OVER_TEXT }, + ]); + }); + test("blocked: refuses without a transient mode, asks transiently with one", async () => { const turnOnly = fakeBridge("blocked"); const refused = collect((await newSession(new SessionBridgeProvider(turnOnly.bridge))).query("q")); diff --git a/packages/ai/session-bridge.ts b/packages/ai/session-bridge.ts index f6c134b00..d26e22e70 100644 --- a/packages/ai/session-bridge.ts +++ b/packages/ai/session-bridge.ts @@ -44,7 +44,14 @@ export type SessionBridgeHost = "pi" | "opencode" | "claude-code"; export type SessionBridgeAskMode = "turn" | "transient"; -export type SessionBridgeErrorCode = "busy" | "blocked" | "gone" | "aborted" | "failed"; +/** + * `taken_over`: someone else's prompt (the person typing into the session, a + * peer, another plugin) entered the turn that was answering our question, so + * the rest of that turn answers them. The host stops streaming at once, keeps + * what it already sent, and settles with this code; the client shows the + * partial answer with a note, not an error. + */ +export type SessionBridgeErrorCode = "busy" | "blocked" | "gone" | "aborted" | "failed" | "taken_over"; export interface SessionBridgeSink { delta(text: string): void; @@ -69,6 +76,13 @@ export interface SessionBridge { * Send one question. The host reports through `sink` exactly once with * `done` or `error`. Aborting `signal` cancels OUR question only: drop it if * it was not delivered yet, stop the turn only if the running turn is ours. + * + * Take-over rule (every host): once a prompt we did not send enters the + * turn that is answering the question, that turn is no longer ours. The + * host stops streaming, settles with `error("taken_over")`, and from then + * on neither an abort of this question nor `interrupt()` stops that turn. + * When the answer had already finished (the turn's last model response + * called no tools) before the other message entered, it settles `done`. */ ask(req: SessionBridgeAskRequest, sink: SessionBridgeSink, signal: AbortSignal): void; /** @@ -98,8 +112,25 @@ export const SESSION_BRIDGE_ERROR = { gone: "session_gone", inFlight: "ask_in_flight", failed: "session_ask_failed", + takenOver: "session_taken_over", } as const; +/** + * The note shown under a partial answer whose turn another message took over. + * Neutral: most hosts cannot tell the person typing from another extension's + * steer or a peer session's message. + */ +export const SESSION_ASK_TAKEN_OVER_TEXT = + "Another message entered this session while it was answering, so the rest of the reply went to that message."; + +/** The same note when the host knows the person typed it (Claude Code: the prompt box, Remote Control). */ +export const SESSION_ASK_TAKEN_OVER_BY_PERSON_TEXT = + "You typed into this session while it was answering, so the rest of the reply went to your prompt."; + +/** Why "Interrupt and ask now" refuses a turn another message took over. */ +export const SESSION_ASK_TAKEN_OVER_INTERRUPT_TEXT = + "The session is now answering another message, so Plannotator will not stop it. Ask when it finishes instead."; + const HOST_LABELS: Record = { pi: "Pi", opencode: "OpenCode", @@ -450,7 +481,9 @@ export class SessionBridgeSession extends BaseSession { ? errorMessage(SESSION_BRIDGE_ERROR.blocked, message || BLOCKED_TEXT) : code === "busy" ? errorMessage(SESSION_BRIDGE_ERROR.agentBusy, message || BUSY_TEXT) - : errorMessage(SESSION_BRIDGE_ERROR.failed, message || "The session could not answer."); + : code === "taken_over" + ? errorMessage(SESSION_BRIDGE_ERROR.takenOver, message || SESSION_ASK_TAKEN_OVER_TEXT) + : errorMessage(SESSION_BRIDGE_ERROR.failed, message || "The session could not answer."); settle(mapped); }, }; diff --git a/packages/review-editor/components/AITab.tsx b/packages/review-editor/components/AITab.tsx index a733539af..e64779431 100644 --- a/packages/review-editor/components/AITab.tsx +++ b/packages/review-editor/components/AITab.tsx @@ -12,7 +12,7 @@ import { AIConfigBar } from './AIConfigBar'; import { submitHint } from '@plannotator/ui/utils/platform'; import { OverlayScrollArea } from '@plannotator/ui/components/OverlayScrollArea'; import type { AIProviderOption } from '@plannotator/ui/utils/aiProvider'; -import { SessionAskActions, SessionAskStatus, sessionAskErrorTone, type SessionAskAction } from '@plannotator/ui/components/ai/SessionAskNotice'; +import { SessionAskActions, SessionAskNote, SessionAskStatus, sessionAskErrorTone, type SessionAskAction } from '@plannotator/ui/components/ai/SessionAskNotice'; interface AITabProps { messages: AIChatEntry[]; @@ -419,6 +419,7 @@ const QAPair = memo<{ Thinking... ) : null} + {!response.error && } ); diff --git a/packages/ui/components/ai/DocumentAIChatPanel.tsx b/packages/ui/components/ai/DocumentAIChatPanel.tsx index e066f136d..c00ddc9a2 100644 --- a/packages/ui/components/ai/DocumentAIChatPanel.tsx +++ b/packages/ui/components/ai/DocumentAIChatPanel.tsx @@ -5,7 +5,7 @@ import { formatRelativeTime, renderChatMarkdown } from '../../utils/aiChatFormat import { OverlayScrollArea } from '../OverlayScrollArea'; import { SparklesIcon } from '../SparklesIcon'; import { AIProviderBar } from './AIProviderBar'; -import { SessionAskActions, SessionAskStatus, sessionAskErrorTone, type SessionAskAction } from './SessionAskNotice'; +import { SessionAskActions, SessionAskNote, SessionAskStatus, sessionAskErrorTone, type SessionAskAction } from './SessionAskNotice'; import { submitHint } from '../../utils/platform'; interface DocumentAIChatPanelProps { @@ -220,6 +220,7 @@ const DocumentQAPair = memo<{ Thinking... ) : null} + {!response.error && } ); diff --git a/packages/ui/components/ai/SessionAskNotice.dom.test.tsx b/packages/ui/components/ai/SessionAskNotice.dom.test.tsx index 64c711bb9..304edb3ff 100644 --- a/packages/ui/components/ai/SessionAskNotice.dom.test.tsx +++ b/packages/ui/components/ai/SessionAskNotice.dom.test.tsx @@ -78,6 +78,14 @@ describe.if(hasDom)('Ask this session notices in the document chat panel', () => expect(el.querySelector('[data-session-ask-actions]')).toBeNull(); }); + // The failure this guards: a taken-over answer rendered as an error, which + // replaces the partial answer the reviewer already read. + test('taken over: the partial answer stays, with the note under it', () => { + const el = render({ messages: [entry({ text: 'Because of X', notice: 'NOTE-SENTINEL' })] }); + expect(el.textContent).toContain('Because of X'); + expect(el.querySelector('[data-session-ask-note="taken-over"]')?.textContent).toBe('NOTE-SENTINEL'); + }); + test('a waiting question shows the waiting status instead of "Thinking"', () => { const el = render({ messages: [entry({ isStreaming: true, status: 'waiting' })] }); expect(el.querySelector('[data-session-ask-status="waiting"]')).not.toBeNull(); diff --git a/packages/ui/components/ai/SessionAskNotice.tsx b/packages/ui/components/ai/SessionAskNotice.tsx index 37c0196b2..f860f97e5 100644 --- a/packages/ui/components/ai/SessionAskNotice.tsx +++ b/packages/ui/components/ai/SessionAskNotice.tsx @@ -84,3 +84,18 @@ export const SessionAskActions: React.FC<{ return null; }; + +/** + * A muted note under an "Ask this session" answer that stopped early because + * the person typed into the session while it answered (the rest of that turn + * went to their prompt). The partial answer above it stands. Renders nothing + * for any other response. + */ +export const SessionAskNote: React.FC<{ response: Pick }> = ({ response }) => { + if (!response.notice) return null; + return ( +

+ {response.notice} +

+ ); +}; diff --git a/packages/ui/hooks/useAIChat.sessionAsk.test.tsx b/packages/ui/hooks/useAIChat.sessionAsk.test.tsx index 991b3f7b8..1abacae42 100644 --- a/packages/ui/hooks/useAIChat.sessionAsk.test.tsx +++ b/packages/ui/hooks/useAIChat.sessionAsk.test.tsx @@ -98,6 +98,29 @@ describe.if(hasDom)('useAIChat — Ask this session', () => { expect(answered.status).toBeUndefined(); }); + // The failure this guards: session_taken_over handled like any error, which + // hides the partial answer and raises the panel's error banner. + test('session_taken_over keeps the streamed answer and records the note, not an error', async () => { + setAITransport({ + session: async () => Response.json({ sessionId: 's1' }), + query: async () => + sse( + { type: 'text_delta', delta: 'Because ' }, + { type: 'error', code: 'session_taken_over', error: 'NOTE-SENTINEL' }, + ), + abort: async () => {}, + permission: () => {}, + } satisfies AITransport); + const chat = await mount(); + await act(async () => { await chat.current!.ask({ prompt: 'why?' }); }); + const { response } = chat.current!.messages[0]; + expect(response.text).toBe('Because '); + expect(response.notice).toBe('NOTE-SENTINEL'); + expect(response.error).toBeUndefined(); + expect(response.isStreaming).toBe(false); + expect(chat.current!.error).toBeNull(); + }); + test('the fallback re-asks on a fresh session of the other provider', async () => { const sessions: Array> = []; const queries: Array> = []; diff --git a/packages/ui/hooks/useAIChat.ts b/packages/ui/hooks/useAIChat.ts index 5da728a28..e1f69774c 100644 --- a/packages/ui/hooks/useAIChat.ts +++ b/packages/ui/hooks/useAIChat.ts @@ -2,6 +2,7 @@ import { useCallback, useEffect, useRef, useState } from 'react'; import type { AIContext } from '@plannotator/core'; import type { AIQuestion, AIResponse } from '../types'; import { generateId } from '../utils/generateId'; +import { SESSION_ASK_ERROR_CODES } from '../utils/aiProvider'; export interface AIChatEntry { question: AIQuestion; @@ -421,6 +422,25 @@ export function useAIChat({ description: msg.description, toolUseId: msg.toolUseId, }]); + } else if (msg.type === 'error' && msg.code === SESSION_ASK_ERROR_CODES.takenOver) { + // "Ask this session": the person typed into the session while it + // answered. The answer so far stands; the reason is a note under + // it, not an error that would replace it. + updateMessages(prev => + prev.map(m => + m.question.id === questionId + ? { + ...m, + response: { + ...m.response, + notice: typeof msg.error === 'string' ? msg.error : undefined, + ...(m.response.status && { status: undefined }), + isStreaming: false, + }, + } + : m + ) + ); } else if (msg.type === 'error') { updateMessages(prev => prev.map(m => diff --git a/packages/ui/types.ts b/packages/ui/types.ts index 95426f299..19ed26e99 100644 --- a/packages/ui/types.ts +++ b/packages/ui/types.ts @@ -473,6 +473,12 @@ export interface AIResponse { errorCode?: string; /** "Ask this session": the question is waiting for a busy session, or interrupting it. */ status?: 'waiting' | 'interrupting'; + /** + * "Ask this session": a note shown under the answer (not an error). Set when + * the person typed into the session while it answered, so the answer stops + * where their prompt took the turn over. + */ + notice?: string; createdAt: number; } diff --git a/packages/ui/utils/aiProvider.ts b/packages/ui/utils/aiProvider.ts index 140812ad7..c965c7825 100644 --- a/packages/ui/utils/aiProvider.ts +++ b/packages/ui/utils/aiProvider.ts @@ -45,6 +45,8 @@ export const SESSION_ASK_ERROR_CODES = { agentBusy: 'agent_busy', blocked: 'session_blocked', gone: 'session_gone', + /** The person typed into the session while it answered: a note under the partial answer, not an error. */ + takenOver: 'session_taken_over', } as const; export function isSessionBridgeProvider(provider: Pick | null | undefined): boolean {