diff --git a/packages/opencode/src/session/llm/ai-sdk.ts b/packages/opencode/src/session/llm/ai-sdk.ts index 8db8985d7b09..1cff5e1458af 100644 --- a/packages/opencode/src/session/llm/ai-sdk.ts +++ b/packages/opencode/src/session/llm/ai-sdk.ts @@ -13,6 +13,10 @@ export function adapterState() { reasoning: 0, currentTextID: undefined as string | undefined, currentReasoningID: undefined as string | undefined, + // Block ids for which a *-start event has already been emitted downstream. + // Used to synthesize a missing start when a provider streams orphan deltas. + startedText: {} as Record, + startedReasoning: {} as Record, toolNames: {} as Record, copilotTotalNanoAiu: undefined as number | undefined, } @@ -73,6 +77,40 @@ function currentReasoningID(state: ReturnType, id: string | return state.currentReasoningID } +// The AI SDK enqueues a non-fatal in-band error part +// { type: "error", error: `${"reasoning" | "text"} part ${id} not found` } +// when a reasoning/text delta (or its end) arrives with no preceding *-start +// block. This is common with OpenAI-compatible and Anthropic proxies that omit +// the start events. The SDK returns and keeps streaming, so this must not be +// promoted to a fatal turn failure. Verified against vercel/ai@6 stream-text.ts +// (text part: L1171/L1190, reasoning part: L1220/L1239). +function isOrphanStreamStateError(error: unknown) { + const message = errorMessage(error).trim() + return message.endsWith(" not found") && (message.startsWith("reasoning part ") || message.startsWith("text part ")) +} + +// Some OpenAI-compatible and Anthropic proxies stream text/reasoning deltas +// without ever emitting the matching *-start block. SessionProcessor requires a +// start before it will create a part (`if (!ctx.currentText) return` for text, +// `if (!(value.id in ctx.reasoningMap)) return` for reasoning), so those orphan +// deltas - and any providerMetadata riding on them, including Anthropic thinking +// signatures - are dropped entirely. Losing the signature breaks thinking-block +// replay on subsequent requests. +// +// The adapter is the protocol-normalization layer (it already synthesizes block +// ids via currentTextID/currentReasoningID), so repair it here: emit the missing +// start before the delta and hand the processor a well-formed stream. Streams +// that do emit starts are unaffected. +function synthesizeStart( + started: Record, + id: string, + make: () => LLMEvent, +): LLMEvent[] { + if (started[id]) return [] + started[id] = true + return [make()] +} + export function toLLMEvents( state: ReturnType, event: AISDKEvent, @@ -126,6 +164,7 @@ export function toLLMEvents( case "text-start": return Effect.sync(() => { state.currentTextID = currentTextID(state, event.id) + state.startedText[state.currentTextID] = true return [ LLMEvent.textStart({ id: state.currentTextID, @@ -135,19 +174,28 @@ export function toLLMEvents( }) case "text-delta": - return Effect.succeed([ - LLMEvent.textDelta({ - id: currentTextID(state, event.id), - text: event.text, - providerMetadata: providerMetadata(event.providerMetadata), - }), - ]) + return Effect.sync(() => { + const id = currentTextID(state, event.id) + return [ + ...synthesizeStart(state.startedText, id, () => + LLMEvent.textStart({ id, providerMetadata: providerMetadata(event.providerMetadata) }), + ), + LLMEvent.textDelta({ + id, + text: event.text, + providerMetadata: providerMetadata(event.providerMetadata), + }), + ] + }) case "text-end": return Effect.sync(() => { const id = currentTextID(state, event.id) state.currentTextID = undefined return [ + ...synthesizeStart(state.startedText, id, () => + LLMEvent.textStart({ id, providerMetadata: providerMetadata(event.providerMetadata) }), + ), LLMEvent.textEnd({ id, providerMetadata: providerMetadata(event.providerMetadata), @@ -158,6 +206,7 @@ export function toLLMEvents( case "reasoning-start": return Effect.sync(() => { state.currentReasoningID = currentReasoningID(state, event.id) + state.startedReasoning[state.currentReasoningID] = true return [ LLMEvent.reasoningStart({ id: state.currentReasoningID, @@ -167,19 +216,28 @@ export function toLLMEvents( }) case "reasoning-delta": - return Effect.succeed([ - LLMEvent.reasoningDelta({ - id: currentReasoningID(state, event.id), - text: event.text, - providerMetadata: providerMetadata(event.providerMetadata), - }), - ]) + return Effect.sync(() => { + const id = currentReasoningID(state, event.id) + return [ + ...synthesizeStart(state.startedReasoning, id, () => + LLMEvent.reasoningStart({ id, providerMetadata: providerMetadata(event.providerMetadata) }), + ), + LLMEvent.reasoningDelta({ + id, + text: event.text, + providerMetadata: providerMetadata(event.providerMetadata), + }), + ] + }) case "reasoning-end": return Effect.sync(() => { const id = currentReasoningID(state, event.id) state.currentReasoningID = undefined return [ + ...synthesizeStart(state.startedReasoning, id, () => + LLMEvent.reasoningStart({ id, providerMetadata: providerMetadata(event.providerMetadata) }), + ), LLMEvent.reasoningEnd({ id, providerMetadata: providerMetadata(event.providerMetadata), @@ -262,6 +320,10 @@ export function toLLMEvents( }) case "error": + if (isOrphanStreamStateError(event.error)) + return Effect.logDebug("dropping orphan reasoning/text stream part", { + detail: errorMessage(event.error), + }).pipe(Effect.as([])) return Effect.fail(event.error) case "abort": diff --git a/packages/opencode/test/session/ai-sdk.test.ts b/packages/opencode/test/session/ai-sdk.test.ts new file mode 100644 index 000000000000..55e6864913c9 --- /dev/null +++ b/packages/opencode/test/session/ai-sdk.test.ts @@ -0,0 +1,108 @@ +import { describe, expect, test } from "bun:test" +import { Effect, Exit } from "effect" +import { LLMAISDK } from "@/session/llm/ai-sdk" + +type ErrorEvent = Extract[1], { type: "error" }> + +const errorEvent = (error: unknown): ErrorEvent => ({ type: "error", error }) + +const run = (error: unknown) => Effect.runPromiseExit(LLMAISDK.toLLMEvents(LLMAISDK.adapterState(), errorEvent(error))) + +type AnyEvent = Parameters[1] + +// Feed a whole event sequence through one adapter state, as a real stream would. +const runStream = async (events: AnyEvent[]) => { + const state = LLMAISDK.adapterState() + const out: any[] = [] + for (const event of events) { + const exit = await Effect.runPromiseExit(LLMAISDK.toLLMEvents(state, event)) + expect(Exit.isSuccess(exit)).toBe(true) + if (Exit.isSuccess(exit)) out.push(...exit.value) + } + return out +} + +describe("LLMAISDK.toLLMEvents error handling", () => { + test("drops orphan reasoning/text stream-state errors without failing the turn", async () => { + // vercel/ai stream-text.ts enqueues these exact strings as non-fatal in-band + // error parts when a reasoning/text delta (or its end) has no preceding + // *-start block. They must not abort the assistant turn. + for (const message of [ + "reasoning part 0 not found", + "reasoning part 7 not found", + "reasoning part rs_abc:0 not found", + "text part 0 not found", + "text part X not found", + ]) { + const exit = await run(message) + expect(Exit.isSuccess(exit)).toBe(true) + if (Exit.isSuccess(exit)) expect(exit.value).toEqual([]) + } + }) + + test("tolerates surrounding whitespace via trim", async () => { + const exit = await run(" reasoning part 0 not found ") + expect(Exit.isSuccess(exit)).toBe(true) + }) + + test("still fails on genuine provider errors", async () => { + for (const error of ["rate limit exceeded", 'Tool "foo" not found', "context length exceeded", new Error("boom")]) { + const exit = await run(error) + expect(Exit.isFailure(exit)).toBe(true) + } + }) +}) + +describe("LLMAISDK.toLLMEvents orphan delta repair", () => { + test("synthesizes a reasoning-start for an orphan reasoning-delta", async () => { + const events = await runStream([ + { type: "reasoning-delta", id: "rs_1", text: "thinking" } as AnyEvent, + { type: "reasoning-delta", id: "rs_1", text: " more" } as AnyEvent, + ]) + // Start is synthesized exactly once, then both deltas flow through. + expect(events.map((e) => e.type)).toEqual(["reasoning-start", "reasoning-delta", "reasoning-delta"]) + expect(events[0].id).toBe("rs_1") + expect(events[1].text).toBe("thinking") + expect(events[2].text).toBe(" more") + }) + + test("carries providerMetadata onto the synthesized start so signatures survive", async () => { + const metadata = { anthropic: { signature: "sig-abc" } } + const events = await runStream([ + { type: "reasoning-delta", id: "rs_2", text: "x", providerMetadata: metadata } as AnyEvent, + ]) + expect(events[0].type).toBe("reasoning-start") + expect(events[0].providerMetadata).toEqual(metadata) + }) + + test("synthesizes a text-start for an orphan text-delta", async () => { + const events = await runStream([{ type: "text-delta", id: "t_1", text: "hello" } as AnyEvent]) + expect(events.map((e) => e.type)).toEqual(["text-start", "text-delta"]) + expect(events[0].id).toBe("t_1") + }) + + test("synthesizes a start for an orphan reasoning-end", async () => { + const events = await runStream([{ type: "reasoning-end", id: "rs_3" } as AnyEvent]) + expect(events.map((e) => e.type)).toEqual(["reasoning-start", "reasoning-end"]) + }) + + test("does not duplicate the start when the provider emits one", async () => { + const events = await runStream([ + { type: "reasoning-start", id: "rs_4" } as AnyEvent, + { type: "reasoning-delta", id: "rs_4", text: "a" } as AnyEvent, + { type: "reasoning-end", id: "rs_4" } as AnyEvent, + { type: "text-start", id: "t_2" } as AnyEvent, + { type: "text-delta", id: "t_2", text: "b" } as AnyEvent, + { type: "text-end", id: "t_2" } as AnyEvent, + ]) + // Well-formed streams must be passed through completely unchanged. + expect(events.map((e) => e.type)).toEqual([ + "reasoning-start", + "reasoning-delta", + "reasoning-end", + "text-start", + "text-delta", + "text-end", + ]) + }) +})