Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
128 changes: 127 additions & 1 deletion packages/ai/ai.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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 },
]);
});
});
157 changes: 132 additions & 25 deletions packages/ai/providers/opencode-sdk.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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 },
Expand All @@ -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<string, unknown> | undefined;
if (!props) continue;

// Filter: only events for our session
const eventSessionId =
(props.sessionID as string) ??
((props.info as Record<string, unknown>)?.sessionID as string) ??
((props.part as Record<string, unknown>)?.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 {
Expand Down Expand Up @@ -385,6 +402,94 @@ function isTerminalEvent(eventType: string): boolean {
return eventType === "session.error" || eventType === "session.status";
}

export function openCodeEventSessionId(
props: Record<string, unknown>,
): string | undefined {
return (
(props.sessionID as string | undefined) ??
((props.info as Record<string, unknown> | undefined)?.sessionID as
| string
| undefined) ??
((props.part as Record<string, unknown> | 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<string, unknown> | 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<string, unknown> },
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<string, unknown> | 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[].
*
Expand All @@ -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 }];
}
Expand Down