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
90 changes: 76 additions & 14 deletions packages/opencode/src/session/llm/ai-sdk.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, boolean>,
startedReasoning: {} as Record<string, boolean>,
toolNames: {} as Record<string, string>,
copilotTotalNanoAiu: undefined as number | undefined,
}
Expand Down Expand Up @@ -73,6 +77,40 @@ function currentReasoningID(state: ReturnType<typeof adapterState>, 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<string, boolean>,
id: string,
make: () => LLMEvent,
): LLMEvent[] {
if (started[id]) return []
started[id] = true
return [make()]
}

export function toLLMEvents(
state: ReturnType<typeof adapterState>,
event: AISDKEvent,
Expand Down Expand Up @@ -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,
Expand All @@ -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),
Expand All @@ -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,
Expand All @@ -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),
Expand Down Expand Up @@ -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<LLMEvent[]>([]))
return Effect.fail(event.error)

case "abort":
Expand Down
108 changes: 108 additions & 0 deletions packages/opencode/test/session/ai-sdk.test.ts
Original file line number Diff line number Diff line change
@@ -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<Parameters<typeof LLMAISDK.toLLMEvents>[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<typeof LLMAISDK.toLLMEvents>[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",
])
})
})
Loading