Compare commits

..

3 Commits

Author SHA1 Message Date
rekram1-node 5a02790945 fix(ai): retry xAI capacity stream errors 2026-08-20 14:28:07 +00:00
Filip 838d747514 fix(core): revert SSE heartbeat handling (#43625) 2026-08-20 13:23:08 +02:00
Filip 879766aee7 fix(core): ignore SSE comment heartbeats (#43618) 2026-08-20 12:59:58 +02:00
6 changed files with 61 additions and 76 deletions
+8 -2
View File
@@ -1073,7 +1073,13 @@ const providerErrorMessage = (event: Event, fallback: string): string => {
}
export const providerFailure = (id: string, event: Event, fallback: string) => {
const code = event.code || event.error?.code || event.response?.error?.code || undefined
const code =
event.code ||
event.error?.code ||
event.error?.type ||
event.response?.error?.code ||
event.response?.error?.type ||
undefined
const message = providerErrorMessage(event, fallback)
const status =
typeof event.status === "number"
@@ -1084,7 +1090,7 @@ export const providerFailure = (id: string, event: Event, fallback: string) => {
return new AIError({
module: id,
method: "stream",
reason: classifyProviderFailure({ message, code, status }),
reason: classifyProviderFailure({ message, code, status, stream: true }),
})
}
+14
View File
@@ -74,9 +74,11 @@ const INVALID_REQUEST_CODES = new Set(["invalid_prompt", "invalid_request_error"
const RATE_LIMIT_TEXT = /rate increased too quickly|rate[-_\s]?limit|too[_\s]?many[_\s]?requests/i
const QUOTA_TEXT = /insufficient[-_\s]?quota|quota[-_\s]?exceeded/i
const CONTENT_POLICY_TEXT = /content[-_\s]?policy|content_filter|safety/i
const TRANSIENT_TEXT = /\btry again (?:later|in\b)|\b(?:currently|temporarily) at capacity\b|\bretry your request\b/i
export interface ProviderFailure {
readonly message: string
readonly stream?: boolean | undefined
readonly status?: number | undefined
readonly code?: string | undefined
readonly retryAfterMs?: number | undefined
@@ -133,6 +135,12 @@ export function classifyProviderFailure(input: ProviderFailure): AIError["reason
status: input.status,
retryAfterMs: input.retryAfterMs,
})
if (TRANSIENT_TEXT.test(text))
return new ProviderInternalReason({
...common,
status: input.status,
retryAfterMs: input.retryAfterMs,
})
if (input.status === 429) {
return new RateLimitReason({
...common,
@@ -149,6 +157,12 @@ export function classifyProviderFailure(input: ProviderFailure): AIError["reason
if (codes.some((code) => INVALID_REQUEST_CODES.has(code))) return new InvalidRequestReason(common)
if (input.status === 400 || input.status === 404 || input.status === 413 || input.status === 422)
return new InvalidRequestReason(common)
if (input.stream)
return new ProviderInternalReason({
...common,
status: input.status,
retryAfterMs: input.retryAfterMs,
})
return new UnknownProviderReason({ ...common, status: input.status })
}
+20
View File
@@ -75,6 +75,26 @@ describe("provider error classification", () => {
).toEqual(["ProviderInternal", "ProviderInternal", "ProviderInternal"])
})
test("classifies transient gateway message fallbacks as provider internal", () => {
expect(
[
"The model is currently at capacity due to high demand.",
"The service is temporarily at capacity.",
"Please try again in a few minutes.",
"Please retry your request shortly.",
].map((message) => classifyProviderFailure({ message })._tag),
).toEqual(["ProviderInternal", "ProviderInternal", "ProviderInternal", "ProviderInternal"])
})
test("classifies provider stream errors as provider internal", () => {
expect(
[
"The model is currently at capacity due to high demand. Please try again in a few minutes, or use a higher service tier for priority processing: https://docs.x.ai/developers/advanced-api-usage/priority-processing",
"The model is temporarily unavailable.",
].map((message) => classifyProviderFailure({ message, stream: true })._tag),
).toEqual(["ProviderInternal", "ProviderInternal"])
})
test("classifies transient client statuses as provider internal", () => {
expect([408, 409].map((status) => classifyProviderFailure({ message: `HTTP ${status}`, status })._tag)).toEqual([
"ProviderInternal",
@@ -809,7 +809,7 @@ describe("OpenAI Responses route", () => {
"ProviderInternal",
"RateLimit",
"ProviderInternal",
"UnknownProvider",
"ProviderInternal",
])
}),
)
@@ -2661,6 +2661,19 @@ describe("OpenAI Responses route", () => {
}),
)
it.effect("classifies code-less capacity errors as transient", () =>
Effect.gen(function* () {
const message =
"The model is currently at capacity due to high demand. Please try again in a few minutes, or use a higher service tier for priority processing: https://docs.x.ai/developers/advanced-api-usage/priority-processing"
const error = yield* LLMClient.generate(request).pipe(
Effect.provide(fixedResponse(sseEvents({ type: "error", message }))),
Effect.flip,
)
expect(error.reason).toMatchObject({ _tag: "ProviderInternal", message })
}),
)
it.effect("falls back to error code when no message is present", () =>
Effect.gen(function* () {
const error = yield* LLMClient.generate(request).pipe(
@@ -2801,7 +2814,7 @@ describe("OpenAI Responses route", () => {
Effect.flip,
)
expect(error.reason).toMatchObject({ _tag: "UnknownProvider", message: "Something went wrong" })
expect(error.reason).toMatchObject({ _tag: "ProviderInternal", message: "Something went wrong" })
}),
)
@@ -2812,7 +2825,7 @@ describe("OpenAI Responses route", () => {
Effect.flip,
)
expect(error.reason).toMatchObject({ _tag: "UnknownProvider", message: "OpenAI Responses stream error" })
expect(error.reason).toMatchObject({ _tag: "ProviderInternal", message: "OpenAI Responses stream error" })
}),
)
@@ -2823,7 +2836,7 @@ describe("OpenAI Responses route", () => {
Effect.flip,
)
expect(error.reason).toMatchObject({ _tag: "UnknownProvider", message: "OpenAI Responses stream error" })
expect(error.reason).toMatchObject({ _tag: "ProviderInternal", message: "OpenAI Responses stream error" })
}),
)
@@ -2834,7 +2847,7 @@ describe("OpenAI Responses route", () => {
Effect.flip,
)
expect(error.reason).toMatchObject({ _tag: "UnknownProvider", message: "OpenAI Responses response failed" })
expect(error.reason).toMatchObject({ _tag: "ProviderInternal", message: "OpenAI Responses response failed" })
}),
)
+1 -11
View File
@@ -34,7 +34,6 @@ import {
import { Auth, Endpoint, RequestExecutor, type AnyRoute } from "@opencode-ai/ai/route"
import { ProviderShared } from "@opencode-ai/ai/protocols/shared"
import { Cause, Context, Effect, Layer, Option, Schema, Scope, Stream } from "effect"
import { makeParser } from "effect/unstable/encoding/Sse"
import type { ID, Info } from "./model.js"
import { Provider } from "./provider.js"
import { State } from "./state.js"
@@ -66,23 +65,15 @@ function wrapSSE(res: Response, ms: number, ctl: AbortController) {
if (!res.headers.get("content-type")?.includes("text/event-stream")) return res
const reader = res.body.getReader()
const decoder = new TextDecoder()
let deadline: number | undefined
const parser = makeParser((event) => {
if (event._tag === "Event") deadline = Date.now() + ms
})
const body = new ReadableStream<Uint8Array>({
async pull(ctrl) {
const expires = deadline ?? Date.now() + ms
deadline = expires
const part = await new Promise<Awaited<ReturnType<typeof reader.read>>>((resolve, reject) => {
const remaining = Math.max(0, expires - Date.now())
const id = setTimeout(() => {
const err = new Error("SSE read timed out")
ctl.abort(err)
void reader.cancel(err)
reject(err)
}, remaining)
}, ms)
reader.read().then(
(part) => {
@@ -101,7 +92,6 @@ function wrapSSE(res: Response, ms: number, ctl: AbortController) {
return
}
parser.feed(decoder.decode(part.value, { stream: true }))
ctrl.enqueue(part.value)
},
async cancel(reason) {
-58
View File
@@ -1,7 +1,6 @@
import { APICallError } from "@ai-sdk/provider"
import type { LanguageModelV3, LanguageModelV3StreamPart } from "@ai-sdk/provider"
import { createMistral } from "@ai-sdk/mistral"
import { createOpenAICompatible } from "@ai-sdk/openai-compatible"
import { AISDK } from "@opencode-ai/core/aisdk"
import { SessionRunnerRetry } from "@opencode-ai/core/session/runner/retry"
import { toSessionError } from "@opencode-ai/core/session/to-session-error"
@@ -413,63 +412,6 @@ it.effect("moves a tool image through the real Mistral provider as a user messag
}),
)
it.effect("does not treat SSE comment heartbeats as model progress", () =>
Effect.gen(function* () {
const aisdk = yield* AISDK.Service
const encoder = new TextEncoder()
let heartbeat: ReturnType<typeof setInterval> | undefined
const customFetch = Object.assign(
async () =>
new Response(
new ReadableStream({
start(controller) {
controller.enqueue(
encoder.encode(
'data: {"id":"response-1","object":"chat.completion.chunk","created":0,"model":"api-model","choices":[{"index":0,"delta":{"content":"partial"},"finish_reason":null}]}\n\n',
),
)
heartbeat = setInterval(() => controller.enqueue(encoder.encode(": keepalive\n\n")), 5)
},
cancel() {
if (heartbeat) clearInterval(heartbeat)
},
}),
{ headers: { "content-type": "text/event-stream" } },
),
{ preconnect: fetch.preconnect },
)
yield* aisdk.hook.sdk((event) => {
event.sdk = createOpenAICompatible({
...event.options,
name: String(event.options.name),
baseURL: String(event.options.baseURL),
})
})
const resolved = yield* aisdk.model(
model("@ai-sdk/openai-compatible", {
apiKey: "test",
baseURL: "https://example.test/v1",
chunkTimeout: 25,
fetch: customFetch,
}),
)
const result = yield* LLMClient.generate(LLM.request({ model: resolved, prompt: "Hello" })).pipe(
Effect.provide(client),
Effect.result,
Effect.ensuring(
Effect.sync(() => {
if (heartbeat) clearInterval(heartbeat)
}),
),
)
expect(result).toMatchObject({
_tag: "Failure",
failure: { reason: { message: expect.stringContaining("SSE read timed out") } },
})
}),
)
it.effect("emits malformed AI SDK tool input without executing it", () =>
Effect.gen(function* () {
const aisdk = yield* AISDK.Service