From f2f8efa4119c60857933cc1d88daa1bdacd63b63 Mon Sep 17 00:00:00 2001 From: Simon Klee Date: Thu, 14 May 2026 21:54:28 +0200 Subject: [PATCH] cli: fix question recovery matching wrong session The recovery logic matched questions by checking if the list was non-empty, which caused it to pick up stale questions from earlier turns. When a re-ask fired, the wrong question could resolve the blocker, leaving the real question undelivered and the process stuck. Match questions by messageID and callID from the originating tool part so only the correct question unblocks the prompt. Fixes #27503 --- .../src/cli/cmd/run/stream.transport.ts | 26 +++++++++++-------- .../test/cli/run/stream.transport.test.ts | 11 ++++++-- 2 files changed, 24 insertions(+), 13 deletions(-) diff --git a/packages/opencode/src/cli/cmd/run/stream.transport.ts b/packages/opencode/src/cli/cmd/run/stream.transport.ts index 22240ebf5..f1f64d7c4 100644 --- a/packages/opencode/src/cli/cmd/run/stream.transport.ts +++ b/packages/opencode/src/cli/cmd/run/stream.transport.ts @@ -15,7 +15,7 @@ // The tick counter prevents stale idle events from resolving the wrong turn. // We also re-check live session status before resolving an idle event so a // delayed idle from an older turn cannot complete a newer busy turn. -import type { Event, GlobalEvent, OpencodeClient } from "@opencode-ai/sdk/v2" +import type { Event, GlobalEvent, OpencodeClient, ToolPart } from "@opencode-ai/sdk/v2" import { Context, Deferred, Effect, Exit, Layer, Scope, Stream } from "effect" import { makeRuntime } from "@/effect/run-service" import { @@ -505,7 +505,10 @@ function createLayer(input: StreamInput) { state.footerView = current } - const recoverQuestion = Effect.fn("RunStreamTransport.recoverQuestion")(function* (partID: string) { + const recoverQuestion = Effect.fn("RunStreamTransport.recoverQuestion")(function* (part: ToolPart) { + const partID = part.id + const matches = (request: SessionData["questions"][number]) => + request.tool?.messageID === part.messageID && request.tool?.callID === part.callID if (recovering.has(partID)) { return } @@ -513,7 +516,7 @@ function createLayer(input: StreamInput) { recovering.add(partID) try { while (!closed && !abort.signal.aborted && !input.footer.isClosed) { - if (state.data.questions.length > 0 || !state.data.tools.has(partID)) { + if (state.data.questions.some(matches) || !state.data.tools.has(partID)) { return } @@ -521,23 +524,25 @@ function createLayer(input: StreamInput) { Effect.map((item) => (item.data ?? []).filter((request) => request.sessionID === input.sessionID)), Effect.orElseSucceed(() => []), ) - if (state.data.questions.length > 0 || !state.data.tools.has(partID)) { + if (state.data.questions.some(matches) || !state.data.tools.has(partID)) { return } - if (questions.length > 0) { + const matching = questions.filter(matches) + if (matching.length > 0) { + state.data.questions = state.data.questions.filter(matches) bootstrapSessionData({ data: state.data, messages: [], permissions: [], - questions, + questions: matching, }) - for (const request of questions) { + for (const request of matching) { seedBlocker(request.id) } input.trace?.write("question.recover", { sessionID: input.sessionID, - requests: questions.map((request) => request.id), + requests: matching.map((request) => request.id), }) syncFooter([]) return @@ -784,10 +789,9 @@ function createLayer(input: StreamInput) { event.properties.part.sessionID === input.sessionID && event.properties.part.type === "tool" && event.properties.part.tool === "question" && - event.properties.part.state.status === "running" && - state.data.questions.length === 0 + event.properties.part.state.status === "running" ) { - yield* recoverQuestion(event.properties.part.id).pipe( + yield* recoverQuestion(event.properties.part).pipe( Effect.forkIn(scope, { startImmediately: true }), Effect.asVoid, ) diff --git a/packages/opencode/test/cli/run/stream.transport.test.ts b/packages/opencode/test/cli/run/stream.transport.test.ts index 3358ae774..54eaee38f 100644 --- a/packages/opencode/test/cli/run/stream.transport.test.ts +++ b/packages/opencode/test/cli/run/stream.transport.test.ts @@ -835,12 +835,17 @@ describe("run stream transport", () => { callID: "call-question-1", }, } + const stale = { + ...request, + id: "question-old", + tool: { messageID: "msg-old", callID: "call-question-old" }, + } const transport = await createSessionTransport({ sdk: sdk({ stream: src.stream, questions: async () => { questionCalls += 1 - return ok(questionCalls > 1 ? [request] : []) + return ok(questionCalls === 1 ? [stale] : [stale, request]) }, promptAsync: async () => { queueMicrotask(() => { @@ -885,7 +890,9 @@ describe("run stream transport", () => { const view = await waitFor(() => { const item = ui.events.findLast((event) => event.type === "stream.view") - return item?.type === "stream.view" && item.view.type === "question" ? item.view : undefined + return item?.type === "stream.view" && item.view.type === "question" && item.view.request.id === request.id + ? item.view + : undefined }) expect(view).toEqual({