diff --git a/packages/ai/ai.test.ts b/packages/ai/ai.test.ts index 8af416d8e..451aff9b4 100644 --- a/packages/ai/ai.test.ts +++ b/packages/ai/ai.test.ts @@ -1663,7 +1663,13 @@ describe("mapPiEvent", () => { // OpenCode event mapping // --------------------------------------------------------------------------- -import { mapOpenCodeEvent } from "./providers/opencode-sdk.ts"; +import { + assistantTextFromPartUpdate, + consumeOpenCodeEvent, + createOpenCodeQueryState, + fallbackAssistantText, + mapOpenCodeEvent, +} from "./providers/opencode-sdk.ts"; describe("mapOpenCodeEvent", () => { const SESSION_ID = "oc_session_1"; @@ -1827,4 +1833,124 @@ describe("mapOpenCodeEvent", () => { expect(mapOpenCodeEvent("server.connected", {}, SESSION_ID)).toEqual([]); expect(mapOpenCodeEvent("file.edited", { file: "foo.ts" }, SESSION_ID)).toEqual([]); }); + + // OpenCode 1.18 lands the answer as a text snapshot. Mapping that as a + // delta would duplicate streamed tokens; Ask AI would show "pongpong". + test("message.part.updated text snapshots are not streamed as deltas", () => { + expect( + mapOpenCodeEvent("message.part.updated", { + sessionID: SESSION_ID, + part: { type: "text", text: "pong", sessionID: SESSION_ID, time: { start: 1 } }, + }, SESSION_ID), + ).toEqual([]); + }); + + // User prompt snapshots have no `time`. Using them as the answer would echo + // the question into the Ask AI bubble (#514 empty-bubble / #907). + test("assistantTextFromPartUpdate ignores user prompts and reasoning", () => { + expect( + assistantTextFromPartUpdate({ type: "text", text: "pong", time: { start: 1 } }), + ).toBe("pong"); + expect( + assistantTextFromPartUpdate({ type: "text", text: "Reply with pong" }), + ).toBeUndefined(); + expect( + assistantTextFromPartUpdate({ type: "reasoning", text: "thinking", time: { start: 1 } }), + ).toBeUndefined(); + }); + + test("fallbackAssistantText emits the snapshot only when deltas were missed", () => { + expect( + fallbackAssistantText({ + sawTextDelta: false, + lastAssistantText: "pong", + }), + ).toEqual([{ type: "text_delta", delta: "pong" }]); + expect( + fallbackAssistantText({ + sawTextDelta: true, + lastAssistantText: "pong", + }), + ).toEqual([]); + }); + + // Failure this guards: Ask AI hangs with Failed to fetch (#514 / #907) + // when the turn only lands a snapshot, never message.part.delta. + test("consumeOpenCodeEvent falls back to the assistant snapshot on idle", () => { + const state = createOpenCodeQueryState(); + consumeOpenCodeEvent( + { + type: "message.part.updated", + properties: { + sessionID: SESSION_ID, + part: { type: "text", text: "Reply with pong", sessionID: SESSION_ID }, + }, + }, + SESSION_ID, + state, + ); + consumeOpenCodeEvent( + { + type: "message.part.updated", + properties: { + sessionID: SESSION_ID, + part: { + type: "text", + text: "pong", + sessionID: SESSION_ID, + time: { start: 1, end: 2 }, + }, + }, + }, + SESSION_ID, + state, + ); + const idle = consumeOpenCodeEvent( + { + type: "session.status", + properties: { sessionID: SESSION_ID, status: { type: "idle" } }, + }, + SESSION_ID, + state, + ); + expect(idle.done).toBe(true); + expect(idle.messages).toEqual([ + { type: "text_delta", delta: "pong" }, + { type: "result", sessionId: SESSION_ID, success: true }, + ]); + }); + + test("consumeOpenCodeEvent does not duplicate streamed deltas", () => { + const state = createOpenCodeQueryState(); + consumeOpenCodeEvent( + { + type: "message.part.delta", + properties: { sessionID: SESSION_ID, field: "text", delta: "pong" }, + }, + SESSION_ID, + state, + ); + consumeOpenCodeEvent( + { + type: "message.part.updated", + properties: { + sessionID: SESSION_ID, + part: { type: "text", text: "pong", sessionID: SESSION_ID, time: { start: 1 } }, + }, + }, + SESSION_ID, + state, + ); + const idle = consumeOpenCodeEvent( + { + type: "session.status", + properties: { sessionID: SESSION_ID, status: { type: "idle" } }, + }, + SESSION_ID, + state, + ); + expect(idle.messages).toEqual([ + { type: "result", sessionId: SESSION_ID, success: true }, + ]); + }); }); diff --git a/packages/ai/providers/opencode-sdk.ts b/packages/ai/providers/opencode-sdk.ts index 0cbef8d84..f653f83e8 100644 --- a/packages/ai/providers/opencode-sdk.ts +++ b/packages/ai/providers/opencode-sdk.ts @@ -280,7 +280,7 @@ class OpenCodeSession extends BaseSession { yield BaseSession.BUSY_ERROR; return; } - const { gen } = started; + const { gen, signal } = started; try { // Build model param if specified @@ -293,11 +293,29 @@ class OpenCodeSession extends BaseSession { } } - // Subscribe to SSE events - const { stream } = await this.config.client.event.subscribe(); + // Open the SSE stream and wait for the first frame *before* prompting. + // OpenCode 1.18+ can finish a short turn before a late subscriber + // connects; Ask AI then hangs until the browser reports Failed to fetch. + const { stream } = await this.config.client.event.subscribe( + undefined, + { signal }, + ); + const iterator = stream[Symbol.asyncIterator](); try { - // Send prompt asynchronously + // First frame only proves the subscriber is live. Do not treat it as + // this turn — it can be server.connected or leftover idle from a + // previous query on the same session. + const first = await iterator.next(); + if (first.done) { + yield { + type: "error", + error: "OpenCode event stream closed before the prompt was sent.", + code: "provider_error", + }; + return; + } + try { await this.config.client.session.promptAsync({ path: { id: this.config.sessionId }, @@ -320,29 +338,28 @@ class OpenCodeSession extends BaseSession { } this._firstQuerySent = true; - // Drain SSE events filtered by session ID - for await (const event of stream) { - const eventType = event.type as string; - const props = event.properties as Record | undefined; - if (!props) continue; - - // Filter: only events for our session - const eventSessionId = - (props.sessionID as string) ?? - ((props.info as Record)?.sessionID as string) ?? - ((props.part as Record)?.sessionID as string); - if (eventSessionId && eventSessionId !== this.config.sessionId) continue; - - const mapped = mapOpenCodeEvent(eventType, props, this.id); - for (const msg of mapped) { - yield msg; - if (msg.type === "result" || (msg.type === "error" && isTerminalEvent(eventType))) { - return; - } - } + const state = createOpenCodeQueryState(); + while (!signal.aborted) { + const next = await iterator.next(); + if (next.done) break; + const { messages, done } = consumeOpenCodeEvent( + next.value, + this.config.sessionId, + state, + ); + for (const msg of messages) yield msg; + if (done) return; } + + yield { + type: "error", + error: signal.aborted + ? "OpenCode query aborted." + : "OpenCode event stream ended without a result.", + code: "provider_error", + }; } finally { - stream.return?.(); + await iterator.return?.(); } } catch (err) { yield { @@ -385,6 +402,94 @@ function isTerminalEvent(eventType: string): boolean { return eventType === "session.error" || eventType === "session.status"; } +export function openCodeEventSessionId( + props: Record, +): string | undefined { + return ( + (props.sessionID as string | undefined) ?? + ((props.info as Record | undefined)?.sessionID as + | string + | undefined) ?? + ((props.part as Record | undefined)?.sessionID as + | string + | undefined) + ); +} + +/** + * OpenCode 1.18 still emits `message.part.delta` for streaming text, but some + * turns only land a final `message.part.updated` snapshot (or a late + * subscriber misses the deltas). Capture assistant text parts so Ask AI can + * still render an answer. User-prompt snapshots have no `time` stamp. + */ +export function assistantTextFromPartUpdate( + part: Record | undefined, +): string | undefined { + if (!part || part.type !== "text") return undefined; + if (part.time == null) return undefined; + const text = part.text; + return typeof text === "string" && text.length > 0 ? text : undefined; +} + +export interface OpenCodeQueryState { + sawTextDelta: boolean; + lastAssistantText: string; +} + +export function createOpenCodeQueryState(): OpenCodeQueryState { + return { sawTextDelta: false, lastAssistantText: "" }; +} + +/** + * If this turn never streamed `message.part.delta` text, emit the last + * assistant snapshot so Ask AI still has something to render (#514 / #907). + */ +export function fallbackAssistantText(state: OpenCodeQueryState): AIMessage[] { + if (state.sawTextDelta || !state.lastAssistantText) return []; + return [{ type: "text_delta", delta: state.lastAssistantText }]; +} + +/** + * Consume one OpenCode SSE event after promptAsync. Returns mapped Ask AI + * messages and whether the query should stop. Pre-prompt leftover idle is + * ignored by only calling this after the prompt is sent. + */ +export function consumeOpenCodeEvent( + event: { type?: string; properties?: Record }, + sessionId: string, + state: OpenCodeQueryState, +): { messages: AIMessage[]; done: boolean } { + const eventType = event.type as string; + const props = event.properties; + if (props == null) return { messages: [], done: false }; + + const eventSessionId = openCodeEventSessionId(props); + if (eventSessionId && eventSessionId !== sessionId) { + return { messages: [], done: false }; + } + + if (eventType === "message.part.updated") { + const snapshot = assistantTextFromPartUpdate( + props.part as Record | undefined, + ); + if (snapshot) state.lastAssistantText = snapshot; + } + + const mapped = mapOpenCodeEvent(eventType, props, sessionId); + const messages: AIMessage[] = []; + let done = false; + for (const msg of mapped) { + if (msg.type === "text_delta") state.sawTextDelta = true; + if (msg.type === "result") { + messages.push(...fallbackAssistantText(state)); + done = true; + } + messages.push(msg); + if (msg.type === "error" && isTerminalEvent(eventType)) done = true; + } + return { messages, done }; +} + /** * Map an OpenCode SSE event to AIMessage[]. * @@ -404,6 +509,8 @@ export function mapOpenCodeEvent( case "message.part.delta": { const field = props.field as string; const delta = props.delta as string; + // Reasoning/thinking deltas share this event; Ask AI only renders + // assistant answer text. if (field === "text" && delta) { return [{ type: "text_delta", delta }]; }