diff --git a/packages/core/src/v1/session.ts b/packages/core/src/v1/session.ts index 181ba9807d05..4904b0c7bc8a 100644 --- a/packages/core/src/v1/session.ts +++ b/packages/core/src/v1/session.ts @@ -34,7 +34,70 @@ export const AuthError = NamedError.create("ProviderAuthError", { message: Schema.String, }) -export const AbortedError = NamedError.create("MessageAbortedError", { message: Schema.String }) +export const AbortSource = Schema.Literals([ + "user_cancel", + "session_cancel", + "provider_abort", + "network_abort", + "first_byte_timeout", + "stream_idle_timeout", + "post_tool_first_event_timeout", + "no_visible_part_timeout", + "server_restart", + "client_disconnect", + "unknown", +]) +export type AbortSource = Schema.Schema.Type + +export const RequestPhase = Schema.Literals([ + "model_stream", + "post_tool_continuation", + "message_finalization", + "unknown", +]) +export type RequestPhase = Schema.Schema.Type + +export const NoResponseDiagnostics = Schema.Struct({ + providerID: Schema.optional(ProviderV2.ID), + modelID: Schema.optional(ModelV2.ID), + sessionID: Schema.optional(SessionSchema.ID), + messageID: Schema.optional(MessageID), + elapsedMs: Schema.optional(NonNegativeInt), + isPostToolContinuation: Schema.optional(Schema.Boolean), + retryAttempt: Schema.optional(NonNegativeInt), + firstStreamEventAt: Schema.optional(NonNegativeInt), + lastStreamEventAt: Schema.optional(NonNegativeInt), + firstVisiblePartAt: Schema.optional(NonNegativeInt), + lastVisiblePartAt: Schema.optional(NonNegativeInt), + partCount: Schema.optional(NonNegativeInt), + tokenCount: Schema.optional(NonNegativeInt), +}) +export type NoResponseDiagnostics = Schema.Schema.Type + +const NoResponseErrorData = { + message: Schema.String, + abortSource: AbortSource, + phase: RequestPhase, + retryable: Schema.Boolean, + diagnostics: Schema.optional(NoResponseDiagnostics), +} + +export const AbortedError = NamedError.create("MessageAbortedError", { + message: Schema.String, + abortSource: Schema.optional(AbortSource), +}) +export const UnexpectedProviderAbortError = NamedError.create("UnexpectedProviderAbortError", { + ...NoResponseErrorData, +}) +export const PostToolContinuationTimeoutError = NamedError.create("PostToolContinuationTimeoutError", { + ...NoResponseErrorData, +}) +export const EmptyAssistantResponseError = NamedError.create("EmptyAssistantResponseError", { + ...NoResponseErrorData, +}) +export const NoResponseError = NamedError.create("NoResponseError", { + ...NoResponseErrorData, +}) export const StructuredOutputError = NamedError.create("StructuredOutputError", { message: Schema.String, retries: NonNegativeInt, @@ -389,6 +452,10 @@ const AssistantErrorSchema = Schema.Union([ NamedError.Unknown.EffectSchema, OutputLengthError.EffectSchema, AbortedError.EffectSchema, + UnexpectedProviderAbortError.EffectSchema, + PostToolContinuationTimeoutError.EffectSchema, + EmptyAssistantResponseError.EffectSchema, + NoResponseError.EffectSchema, StructuredOutputError.EffectSchema, ContextOverflowError.EffectSchema, ContentFilterError.EffectSchema, diff --git a/packages/opencode/src/session/llm.ts b/packages/opencode/src/session/llm.ts index adacfc431549..d2bd926dc0cb 100644 --- a/packages/opencode/src/session/llm.ts +++ b/packages/opencode/src/session/llm.ts @@ -45,6 +45,10 @@ export type StreamInput = { tools: Record retries?: number toolChoice?: "auto" | "required" | "none" + internal?: { + postToolContinuation?: boolean + postToolFirstEventTimeoutMs?: number + } } export type StreamRequest = StreamInput & { diff --git a/packages/opencode/src/session/message-v2.ts b/packages/opencode/src/session/message-v2.ts index 813fe49f325d..80f9cb214904 100644 --- a/packages/opencode/src/session/message-v2.ts +++ b/packages/opencode/src/session/message-v2.ts @@ -4,16 +4,23 @@ import { SessionV1 } from "@opencode-ai/core/v1/session" import { ProviderV2 } from "@opencode-ai/core/provider" import { APIError, + AbortSource, AbortedError, Assistant, AuthError, CompactionPart, ContextOverflowError, + EmptyAssistantResponseError, Info, + NoResponseDiagnostics, + NoResponseError, OutputLengthError, Part, + PostToolContinuationTimeoutError, + RequestPhase, StructuredOutputError, SubtaskPart, + UnexpectedProviderAbortError, User, WithParts, type ToolPart, @@ -50,6 +57,29 @@ interface FetchDecompressionError extends Error { export const SYNTHETIC_ATTACHMENT_PROMPT = "Attached media from tool result:" export { isMedia } +export { + APIError, + AbortSource, + AbortedError, + Assistant, + AuthError, + CompactionPart, + ContextOverflowError, + EmptyAssistantResponseError, + Info, + NoResponseDiagnostics, + NoResponseError, + OutputLengthError, + Part, + PostToolContinuationTimeoutError, + RequestPhase, + StructuredOutputError, + SubtaskPart, + UnexpectedProviderAbortError, + User, + WithParts, +} +export type { ToolPart } function truncateToolOutput(text: string, maxChars?: number) { if (!maxChars || text.length <= maxChars) return text @@ -262,7 +292,7 @@ export const toModelMessagesEffect = Effect.fnUntraced(function* ( if ( msg.info.error && !( - AbortedError.isInstance(msg.info.error) && + (AbortedError.isInstance(msg.info.error) || UnexpectedProviderAbortError.isInstance(msg.info.error)) && msg.parts.some((part) => part.type !== "step-start" && part.type !== "reasoning") ) ) { @@ -616,12 +646,39 @@ export function latest(msgs: WithParts[]) { export function fromError( e: unknown, - ctx: { providerID: ProviderV2.ID; aborted?: boolean }, + ctx: { + providerID: ProviderV2.ID + aborted?: boolean + abortSource?: AbortSource + phase?: RequestPhase + diagnostics?: NoResponseDiagnostics + }, ): NonNullable { switch (true) { + case UnexpectedProviderAbortError.isInstance(e): + case PostToolContinuationTimeoutError.isInstance(e): + case EmptyAssistantResponseError.isInstance(e): + case NoResponseError.isInstance(e): + return e instanceof NamedError + ? (e.toObject() as NonNullable) + : (e as NonNullable) case e instanceof DOMException && e.name === "AbortError": - return new AbortedError( - { message: e.message }, + if (ctx.aborted || ctx.abortSource === "user_cancel" || ctx.abortSource === "session_cancel") { + return new AbortedError( + { message: e.message, abortSource: ctx.abortSource ?? "user_cancel" }, + { + cause: e, + }, + ).toObject() + } + return new UnexpectedProviderAbortError( + { + message: e.message, + abortSource: ctx.abortSource ?? "provider_abort", + phase: ctx.phase ?? "model_stream", + retryable: true, + diagnostics: ctx.diagnostics, + }, { cause: e, }, @@ -651,7 +708,10 @@ export function fromError( ).toObject() case e instanceof Error && (e as FetchDecompressionError).code === "ZlibError": if (ctx.aborted) { - return new AbortedError({ message: e.message }, { cause: e }).toObject() + return new AbortedError( + { message: e.message, abortSource: ctx.abortSource ?? "user_cancel" }, + { cause: e }, + ).toObject() } return new APIError( { diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 2554315908ce..70ccfd7c8e15 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -33,6 +33,7 @@ import { RuntimeFlags } from "@/effect/runtime-flags" import { ToolOutput, Usage, type LLMEvent } from "@opencode-ai/llm" const DOOM_LOOP_THRESHOLD = 3 +export const POST_TOOL_FIRST_EVENT_TIMEOUT_MS = 100_000 export type Result = "compact" | "stop" | "continue" export interface Handle { @@ -87,6 +88,18 @@ interface ProcessorContext extends Input { type StreamEvent = LLMEvent +type StreamActivity = { + requestStartedAt: number + firstStreamEventAt: number | undefined + lastStreamEventAt: number | undefined + firstVisiblePartAt: number | undefined + lastVisiblePartAt: number | undefined + partCount: number + tokenCount: number + isPostToolContinuation: boolean + retryAttempt: number +} + export class Service extends Context.Service()("@opencode/SessionProcessor") {} export const layer = Layer.effect( @@ -128,11 +141,43 @@ export const layer = Layer.effect( } const mirrorAssistant = flags.experimentalEventSystem && !input.assistantMessage.summary let aborted = false + let retryAttempt = 0 + let activity: StreamActivity = { + requestStartedAt: Date.now(), + firstStreamEventAt: undefined, + lastStreamEventAt: undefined, + firstVisiblePartAt: undefined, + lastVisiblePartAt: undefined, + partCount: 0, + tokenCount: 0, + isPostToolContinuation: false, + retryAttempt, + } + let assistantOutputIsEmpty = true + + const diagnostics = (): MessageV2.NoResponseDiagnostics => ({ + providerID: input.model.providerID, + modelID: input.model.id, + sessionID: input.sessionID, + messageID: input.assistantMessage.id, + elapsedMs: Math.max(0, Date.now() - activity.requestStartedAt), + isPostToolContinuation: activity.isPostToolContinuation, + retryAttempt: activity.retryAttempt, + firstStreamEventAt: activity.firstStreamEventAt, + lastStreamEventAt: activity.lastStreamEventAt, + firstVisiblePartAt: activity.firstVisiblePartAt, + lastVisiblePartAt: activity.lastVisiblePartAt, + partCount: activity.partCount, + tokenCount: activity.tokenCount, + }) const parse = (e: unknown) => MessageV2.fromError(e, { providerID: input.model.providerID, aborted, + abortSource: aborted ? "user_cancel" : undefined, + phase: activity.isPostToolContinuation ? "post_tool_continuation" : "model_stream", + diagnostics: diagnostics(), }) const settleToolCall = Effect.fn("SessionProcessor.settleToolCall")(function* (toolCallID: string) { @@ -368,10 +413,24 @@ export const layer = Layer.effect( } } + const toolInput = (value: unknown): Record => (isRecord(value) ? value : { value }) + + const markVisiblePart = (increment = true) => { + const now = Date.now() + activity.firstVisiblePartAt ??= now + activity.lastVisiblePartAt = now + if (increment) activity.partCount++ + assistantOutputIsEmpty = false + } + + const tokenTotal = (tokens: SessionV1.Assistant["tokens"]) => + tokens.total ?? tokens.input + tokens.output + tokens.reasoning + tokens.cache.read + tokens.cache.write + const handleEvent = Effect.fnUntraced(function* (value: StreamEvent) { switch (value.type) { case "reasoning-start": if (value.id in ctx.reasoningMap) return + markVisiblePart() // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (mirrorAssistant) { yield* events.publish(SessionEvent.Reasoning.Started, { @@ -397,6 +456,7 @@ export const layer = Layer.effect( case "reasoning-delta": // Match dev: silently drop orphan deltas (no preceding reasoning-start). if (!(value.id in ctx.reasoningMap)) return + markVisiblePart(false) ctx.reasoningMap[value.id].text += value.text if (value.providerMetadata) ctx.reasoningMap[value.id].metadata = value.providerMetadata if (mirrorAssistant) { @@ -428,6 +488,7 @@ export const layer = Layer.effect( if (ctx.assistantMessage.summary) { throw new Error(`Tool call not allowed while generating summary: ${value.name}`) } + markVisiblePart() yield* ensureToolCall(value) return @@ -469,6 +530,7 @@ export const layer = Layer.effect( if (ctx.assistantMessage.summary) { throw new Error(`Tool call not allowed while generating summary: ${value.name}`) } + markVisiblePart(false) const toolCall = yield* ensureToolCall(value) const input = isRecord(value.input) ? value.input : { value: value.input } if (!toolCall.call.inputEnded) { @@ -547,6 +609,7 @@ export const layer = Layer.effect( } case "tool-result": { + markVisiblePart(false) const toolCall = yield* readToolCall(value.id) if (!toolCall && value.result.type === "error") return if (value.result.type === "error") { @@ -606,7 +669,8 @@ export const layer = Layer.effect( }) as const, ) ?? []), ] - const unsupported = content.find((item) => item.type === "file" && !item.uri.startsWith("data:")) + const toolOutput = ToolOutput.make(output.metadata, content) + const unsupported = toolOutput.content.find((item) => item.type === "file" && !item.uri.startsWith("data:")) if (unsupported?.type === "file") { const error = new Error( `Tool attachment URI "${unsupported.uri}" must be materialized before durable V2 settlement`, @@ -632,8 +696,8 @@ export const layer = Layer.effect( sessionID: ctx.sessionID, assistantMessageID, callID: value.id, - structured: output.metadata, - content, + structured: isRecord(toolOutput.structured) ? toolOutput.structured : { value: toolOutput.structured }, + content: toolOutput.content, result: value.result, provider: { executed: value.providerExecuted === true || toolCall?.part.metadata?.providerExecuted === true, @@ -647,6 +711,7 @@ export const layer = Layer.effect( } case "tool-error": { + markVisiblePart(false) const toolCall = yield* readToolCall(value.id) // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (mirrorAssistant) { @@ -674,6 +739,7 @@ export const layer = Layer.effect( throw new Error(value.message) case "step-start": + activity.partCount++ if (!ctx.snapshot) ctx.snapshot = yield* snapshot.track() if (!ctx.assistantMessage.summary) { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. @@ -716,6 +782,8 @@ export const layer = Layer.effect( ctx.assistantMessage.finish = value.reason ctx.assistantMessage.cost += usage.cost ctx.assistantMessage.tokens = usage.tokens + activity.tokenCount = tokenTotal(usage.tokens) + activity.partCount++ yield* session.updatePart({ id: PartID.ascending(), reason: value.reason, @@ -757,6 +825,7 @@ export const layer = Layer.effect( } case "text-start": + markVisiblePart() if (!ctx.assistantMessage.summary) { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. if (mirrorAssistant) { @@ -783,6 +852,7 @@ export const layer = Layer.effect( case "text-delta": if (!ctx.currentText) return + markVisiblePart(false) ctx.currentText.text += value.text if (value.providerMetadata) ctx.currentText.metadata = value.providerMetadata if (mirrorAssistant) { @@ -805,6 +875,7 @@ export const layer = Layer.effect( case "text-end": if (!ctx.currentText) return + markVisiblePart(false) // oxlint-disable-next-line no-self-assign -- reactivity trigger ctx.currentText.text = ctx.currentText.text ctx.currentText.text = (yield* plugin.trigger( @@ -922,6 +993,19 @@ export const layer = Layer.effect( stack: e instanceof Error ? e.stack : undefined, }) const error = parse(e) + if (MessageV2.UnexpectedProviderAbortError.isInstance(error)) { + yield* Effect.logWarning("model.abort.unexpected_provider_abort", { + ...(error.data.diagnostics ?? diagnostics()), + abortSource: error.data.abortSource, + phase: error.data.phase, + }) + } + if (MessageV2.AbortedError.isInstance(error)) { + yield* Effect.logInfo("model.abort.user_cancel", { + ...diagnostics(), + abortSource: error.data.abortSource ?? "user_cancel", + }) + } yield* flushV2Fragments() if (SessionV1.ContextOverflowError.isInstance(error)) { if ((yield* config.get()).compaction?.auto === false && !ctx.assistantMessage.summary) { @@ -957,6 +1041,14 @@ export const layer = Layer.effect( yield* status.set(ctx.sessionID, { type: "idle" }) }) + const readAssistantOutputIsEmpty = Effect.fnUntraced(function* () { + const parts = yield* MessageV2.parts(ctx.assistantMessage.id).pipe( + Effect.provideService(Database.Service, database), + ) + assistantOutputIsEmpty = isEmptyAssistantOutput(parts) + return assistantOutputIsEmpty + }) + const process = Effect.fn("SessionProcessor.process")(function* (streamInput: LLM.StreamInput) { yield* Effect.logInfo("process", { "session.id": input.sessionID, @@ -964,20 +1056,92 @@ export const layer = Layer.effect( }) ctx.needsCompaction = false ctx.shouldBreak = (yield* config.get()).experimental?.continue_loop_on_deny !== true - return yield* Effect.gen(function* () { yield* Effect.gen(function* () { + activity = { + requestStartedAt: Date.now(), + firstStreamEventAt: undefined, + lastStreamEventAt: undefined, + firstVisiblePartAt: undefined, + lastVisiblePartAt: undefined, + partCount: 0, + tokenCount: 0, + isPostToolContinuation: streamInput.internal?.postToolContinuation === true, + retryAttempt, + } ctx.currentText = undefined ctx.currentTextID = undefined ctx.reasoningMap = {} yield* status.set(ctx.sessionID, { type: "busy" }) + assistantOutputIsEmpty = yield* readAssistantOutputIsEmpty() const stream = llm.stream(streamInput) + const firstStreamEvent = yield* Deferred.make() + const firstEventTimeoutMs = + streamInput.internal?.postToolFirstEventTimeoutMs ?? POST_TOOL_FIRST_EVENT_TIMEOUT_MS - yield* stream.pipe( - Stream.tap((event) => handleEvent(event)), + const drain = stream.pipe( + Stream.tap((event) => + Effect.sync(() => { + const now = Date.now() + activity.firstStreamEventAt ??= now + activity.lastStreamEventAt = now + }).pipe( + Effect.andThen(Deferred.succeed(firstStreamEvent, undefined).pipe(Effect.ignore)), + Effect.andThen(handleEvent(event)), + ), + ), Stream.takeUntil(() => ctx.needsCompaction), Stream.runDrain, ) + + const postToolFirstEventTimeout = activity.isPostToolContinuation + ? Deferred.await(firstStreamEvent).pipe( + Effect.timeoutOrElse({ + duration: firstEventTimeoutMs, + orElse: () => + Effect.gen(function* () { + yield* Effect.logWarning("model.no_response.post_tool_first_event_timeout", { + ...diagnostics(), + abortSource: "post_tool_first_event_timeout", + phase: "post_tool_continuation", + }) + return yield* Effect.fail( + new MessageV2.PostToolContinuationTimeoutError({ + message: `No stream event within ${firstEventTimeoutMs}ms after tool continuation (retry attempt ${activity.retryAttempt})`, + abortSource: "post_tool_first_event_timeout", + phase: "post_tool_continuation", + retryable: true, + diagnostics: diagnostics(), + }), + ) + }), + }), + Effect.andThen(Effect.never), + ) + : Effect.never + + yield* drain.pipe(Effect.raceFirst(postToolFirstEventTimeout)) + + if (ctx.assistantMessage.role === "assistant") { + const noParts = yield* readAssistantOutputIsEmpty() + const noTokens = activity.tokenCount === 0 && ctx.assistantMessage.tokens.output === 0 + if (noParts && noTokens) { + yield* Effect.logWarning("model.no_response.zero_part_assistant_turn", { + ...diagnostics(), + abortSource: "unknown", + phase: "message_finalization", + }) + return yield* Effect.fail( + new MessageV2.EmptyAssistantResponseError({ + message: "Assistant stream ended without content", + abortSource: "unknown", + phase: "message_finalization", + retryable: true, + diagnostics: diagnostics(), + }), + ) + } + } }).pipe( Effect.onInterrupt(() => Effect.gen(function* () { @@ -995,7 +1159,15 @@ export const layer = Layer.effect( SessionRetry.policy({ provider: input.model.providerID, parse, + context: () => ({ + aborted, + empty: assistantOutputIsEmpty, + postToolContinuation: activity.isPostToolContinuation, + subagent: streamInput.parentSessionID !== undefined, + }), set: (info) => { + retryAttempt = info.attempt + activity.retryAttempt = info.attempt // TODO(v2): Temporary dual-write while migrating session messages to v2 events. const event = mirrorAssistant ? events.publish(SessionEvent.Retried, { @@ -1008,18 +1180,28 @@ export const layer = Layer.effect( timestamp: DateTime.makeUnsafe(Date.now()), }) : Effect.void - return flushV2Fragments().pipe( - Effect.andThen(event), - Effect.andThen( - status.set(ctx.sessionID, { - type: "retry", + return Effect.gen(function* () { + if (activity.isPostToolContinuation) { + yield* Effect.logWarning("model.no_response.retrying_continuation", { + ...diagnostics(), + abortSource: "post_tool_first_event_timeout", + phase: "post_tool_continuation", attempt: info.attempt, - message: info.message, - action: info.action, - next: info.next, - }), - ), - ) + }) + } + return yield* flushV2Fragments().pipe( + Effect.andThen(event), + Effect.andThen( + status.set(ctx.sessionID, { + type: "retry", + attempt: info.attempt, + message: info.message, + action: info.action, + next: info.next, + }), + ), + ) + }) }, }), ), @@ -1065,6 +1247,20 @@ export const defaultLayer = Layer.suspend(() => ), ) +function isEmptyAssistantOutput(parts: MessageV2.Part[]) { + if (parts.length === 0) return true + if (parts.every((part) => part.type === "step-start")) return true + if (!parts.every(isStructuralAssistantPart)) return false + return ( + parts.some((part) => part.type === "step-finish") && + parts.every((part) => part.type !== "step-finish" || part.reason === "unknown") + ) +} + +function isStructuralAssistantPart(part: MessageV2.Part) { + return part.type === "step-start" || part.type === "step-finish" +} + export const node = LayerNode.make(layer, [ Session.node, Config.node, diff --git a/packages/opencode/src/session/prompt.ts b/packages/opencode/src/session/prompt.ts index 92d371077fea..f50aa2555d4b 100644 --- a/packages/opencode/src/session/prompt.ts +++ b/packages/opencode/src/session/prompt.ts @@ -103,6 +103,15 @@ function isOrphanedInterruptedTool(part: SessionV1.ToolPart) { return part.state.status === "error" && part.state.metadata?.interrupted === true } +function isPostToolContinuation(lastAssistant: MessageV2.WithParts | undefined, finish: string | undefined) { + if (!lastAssistant || finish !== "tool-calls") return false + const tools = lastAssistant.parts.filter( + (part): part is MessageV2.ToolPart => part.type === "tool" && !isOrphanedInterruptedTool(part), + ) + if (tools.length === 0) return false + return tools.every((part) => part.state.status === "completed" || part.state.status === "error") +} + export interface Interface { readonly cancel: (sessionID: SessionID) => Effect.Effect readonly prompt: (input: PromptInput) => Effect.Effect @@ -1365,6 +1374,7 @@ export const layer = Layer.effect( const system = [...env, ...instructions, ...(skills ? [skills] : [])] const format = lastUser.format ?? { type: "text" as const } if (format.type === "json_schema") system.push(STRUCTURED_OUTPUT_SYSTEM_PROMPT) + const postToolContinuation = isPostToolContinuation(lastAssistantMsg, lastAssistant?.finish) const result = yield* handle.process({ user: lastUser, agent, @@ -1379,6 +1389,9 @@ export const layer = Layer.effect( tools, model, toolChoice: format.type === "json_schema" ? "required" : undefined, + internal: { + postToolContinuation, + }, }) if (structured !== undefined) { diff --git a/packages/opencode/src/session/retry.ts b/packages/opencode/src/session/retry.ts index 4139665bd2bd..000d9b7e84d8 100644 --- a/packages/opencode/src/session/retry.ts +++ b/packages/opencode/src/session/retry.ts @@ -11,6 +11,14 @@ export const GO_UPSELL_MESSAGE = "Free usage exceeded, subscribe to Go" export const GO_UPSELL_URL = "https://opencode.ai/go" export type RetryReason = "free_tier_limit" | "account_rate_limit" | (string & {}) +export type RetryContext = { + aborted?: boolean + empty?: boolean + attempt?: number + postToolContinuation?: boolean + subagent?: boolean +} + export type Retryable = { message: string action?: { @@ -27,6 +35,9 @@ export const RETRY_INITIAL_DELAY = 2000 export const RETRY_BACKOFF_FACTOR = 2 export const RETRY_MAX_DELAY_NO_HEADERS = 30_000 // 30 seconds export const RETRY_MAX_DELAY = 2_147_483_647 // max 32-bit signed integer for setTimeout +export const RETRY_UNEXPECTED_ABORT_LIMIT = 1 +export const RETRY_NO_RESPONSE_LIMIT = 1 +export const RETRY_MAIN_SESSION_NO_RESPONSE_LIMIT = 10 function cap(ms: number) { return Math.min(ms, RETRY_MAX_DELAY) @@ -65,9 +76,35 @@ export function delay(attempt: number, error?: SessionV1.APIError) { return cap(Math.min(RETRY_INITIAL_DELAY * Math.pow(RETRY_BACKOFF_FACTOR, attempt - 1), RETRY_MAX_DELAY_NO_HEADERS)) } -export function retryable(error: Err, provider: string) { +export function retryable(error: Err, provider: string, context?: RetryContext) { // context overflow errors should not be retried if (SessionV1.ContextOverflowError.isInstance(error)) return undefined + if (MessageV2.AbortedError.isInstance(error)) { + if (error.data.abortSource === "user_cancel" || error.data.abortSource === "session_cancel") return undefined + if (context?.aborted) return undefined + if (context?.empty === false) return undefined + if ((context?.attempt ?? 1) > RETRY_UNEXPECTED_ABORT_LIMIT) return undefined + return { message: error.data.message } + } + if (MessageV2.UnexpectedProviderAbortError.isInstance(error)) { + if (context?.aborted) return undefined + if (context?.empty === false) return undefined + if ((context?.attempt ?? 1) > RETRY_UNEXPECTED_ABORT_LIMIT) return undefined + return { message: error.data.message } + } + if ( + MessageV2.PostToolContinuationTimeoutError.isInstance(error) || + MessageV2.EmptyAssistantResponseError.isInstance(error) || + MessageV2.NoResponseError.isInstance(error) + ) { + if (context?.aborted) return undefined + if (context?.empty === false) return undefined + const isPostToolContinuationTimeout = MessageV2.PostToolContinuationTimeoutError.isInstance(error) + if (isPostToolContinuationTimeout && context?.postToolContinuation !== true) return undefined + const limit = context?.subagent === true ? RETRY_NO_RESPONSE_LIMIT : RETRY_MAIN_SESSION_NO_RESPONSE_LIMIT + if ((context?.attempt ?? 1) > limit) return undefined + return { message: error.data.message } + } if (SessionV1.APIError.isInstance(error)) { const status = error.data.statusCode // 5xx errors are transient server failures and should always be retried, @@ -177,11 +214,12 @@ export function policy(opts: { provider: string parse: (error: unknown) => Err set: (input: { attempt: number; message: string; action?: Retryable["action"]; next: number }) => Effect.Effect + context?: () => RetryContext }) { return Schedule.fromStepWithMetadata( Effect.succeed((meta: Schedule.InputMetadata) => { const error = opts.parse(meta.input) - const retry = retryable(error, opts.provider) + const retry = retryable(error, opts.provider, { ...opts.context?.(), attempt: meta.attempt }) if (!retry) return Cause.done(meta.attempt) return Effect.gen(function* () { const wait = delay(meta.attempt, SessionV1.APIError.isInstance(error) ? error : undefined) diff --git a/packages/opencode/test/session/message-v2.test.ts b/packages/opencode/test/session/message-v2.test.ts index 1de84c9dd95b..3e3d9bdc9894 100644 --- a/packages/opencode/test/session/message-v2.test.ts +++ b/packages/opencode/test/session/message-v2.test.ts @@ -1041,6 +1041,45 @@ describe("session.message-v2.toModelMessage", () => { ]) }) + test("includes unexpected provider abort messages when they have partial content", async () => { + const assistantID = "m-assistant" + const error = new MessageV2.UnexpectedProviderAbortError({ + message: "aborted", + abortSource: "provider_abort", + phase: "model_stream", + retryable: true, + }).toObject() as MessageV2.Assistant["error"] + + const input: MessageV2.WithParts[] = [ + { + info: assistantInfo(assistantID, "m-parent", error), + parts: [ + { + ...basePart(assistantID, "a1"), + type: "reasoning", + text: "thinking", + time: { start: 0 }, + }, + { + ...basePart(assistantID, "a2"), + type: "text", + text: "partial answer", + }, + ] as MessageV2.Part[], + }, + ] + + expect(await MessageV2.toModelMessages(input, model)).toStrictEqual([ + { + role: "assistant", + content: [ + { type: "reasoning", text: "thinking", providerOptions: undefined }, + { type: "text", text: "partial answer" }, + ], + }, + ]) + }) + test("preserves OpenRouter reasoning details through provider transform", async () => { const assistantID = "m-assistant" const openrouterModel: Provider.Model = { @@ -1551,6 +1590,43 @@ describe("session.message-v2.fromError", () => { expect(result.name).toBe("MessageAbortedError") }) + + test("classifies user AbortError as MessageAbortedError with provenance", () => { + const result = MessageV2.fromError(new DOMException("Aborted", "AbortError"), { + providerID, + aborted: true, + abortSource: "session_cancel", + }) + + expect(result).toStrictEqual({ + name: "MessageAbortedError", + data: { + message: "Aborted", + abortSource: "session_cancel", + }, + }) + }) + + test("classifies unexpected AbortError as provider abort", () => { + const result = MessageV2.fromError(new DOMException("Aborted", "AbortError"), { + providerID, + diagnostics: { + providerID, + elapsedMs: 25, + isPostToolContinuation: true, + retryAttempt: 0, + partCount: 0, + tokenCount: 0, + }, + }) + + expect(result.name).toBe("UnexpectedProviderAbortError") + if (!MessageV2.UnexpectedProviderAbortError.isInstance(result)) throw new Error("expected provider abort") + expect(result.data.abortSource).toBe("provider_abort") + expect(result.data.phase).toBe("model_stream") + expect(result.data.retryable).toBe(true) + expect(result.data.diagnostics?.isPostToolContinuation).toBe(true) + }) }) describe("session.message-v2.latest", () => { diff --git a/packages/opencode/test/session/processor-effect.test.ts b/packages/opencode/test/session/processor-effect.test.ts index c8f40d0de1f9..d72386b0021d 100644 --- a/packages/opencode/test/session/processor-effect.test.ts +++ b/packages/opencode/test/session/processor-effect.test.ts @@ -4,7 +4,8 @@ import { LayerNode } from "@opencode-ai/core/effect/layer-node" import { EventV2Bridge } from "@/event-v2-bridge" import { expect } from "bun:test" import { tool } from "ai" -import { Cause, Effect, Exit, Fiber, Layer, Stream } from "effect" +import { Cause, Effect, Exit, Fiber, Layer } from "effect" +import * as Stream from "effect/Stream" import path from "path" import z from "zod" import type { Agent } from "../../src/agent/agent" @@ -22,11 +23,11 @@ import { provideTmpdirInstance, provideTmpdirServer } from "../fixture/fixture" import { testEffect } from "../lib/effect" import { raw, reply, TestLLMServer } from "../lib/llm-server" import { RuntimeFlags } from "@/effect/runtime-flags" +import { LLMEvent, Usage } from "@opencode-ai/llm" import { ProviderV2 } from "@opencode-ai/core/provider" import { ModelV2 } from "@opencode-ai/core/model" import { SessionEvent } from "@opencode-ai/core/session/event" import { SessionProjector } from "@opencode-ai/core/session/projector" -import { LLMEvent } from "@opencode-ai/llm" const summary = Layer.succeed( SessionSummary.Service, @@ -235,6 +236,72 @@ const boot = Effect.fn("test.boot")(function* () { return { processors, session, provider } }) +const basicUsage = () => new Usage({ inputTokens: 1, outputTokens: 1, totalTokens: 2 }) + +function assistantSuccessStream(text: string) { + return Stream.make( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.textStart({ id: "txt-0" }), + LLMEvent.textDelta({ id: "txt-0", text }), + LLMEvent.textEnd({ id: "txt-0" }), + LLMEvent.stepFinish({ index: 0, reason: "stop", usage: basicUsage() }), + LLMEvent.finish({ reason: "stop", usage: basicUsage() }), + ) +} + +function abortStream(...events: LLMEvent[]) { + return Stream.fromAsyncIterable( + { + async *[Symbol.asyncIterator]() { + yield* events + throw new DOMException("The operation was aborted.", "AbortError") + }, + }, + (error) => error, + ) +} + +function delayedFirstEventStream(delayMs: number) { + return Stream.fromAsyncIterable( + { + async *[Symbol.asyncIterator]() { + await Bun.sleep(delayMs) + yield LLMEvent.stepStart({ index: 0 }) + }, + }, + (error) => error, + ) +} + +function structuralEmptyAssistantStream() { + return Stream.make( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.stepFinish({ index: 0, reason: "unknown", usage: new Usage({}) }), + LLMEvent.finish({ reason: "unknown", usage: new Usage({}) }), + ) +} + +function llmMock(...streams: Stream.Stream[]) { + let calls = 0 + return { + calls: () => calls, + layer: Layer.succeed( + LLM.Service, + LLM.Service.of({ + stream: () => { + calls += 1 + return streams.shift() ?? Stream.empty + }, + }), + ), + } +} + +function processorEnv(llmLayer: Layer.Layer) { + return LayerNode.buildLayer(root, { + replacements: [...replacements, LayerNode.replace(LLM.node, llmLayer)], + }) +} // --------------------------------------------------------------------------- // Tests // --------------------------------------------------------------------------- @@ -606,6 +673,410 @@ it.live("session.processor effect tests retry recognized structured json errors" ), ) +const unexpectedAbortLLM = llmMock(abortStream(LLMEvent.stepStart({ index: 0 })), assistantSuccessStream("after")) +const unexpectedAbortIt = testEffect(processorEnv(unexpectedAbortLLM.layer)) + +unexpectedAbortIt.live("session.processor effect tests retry unexpected aborts before assistant output", () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const { processors, session, provider } = yield* boot() + + const chat = yield* session.create({}) + const parent = yield* user(chat.id, "retry abort") + const msg = yield* assistant(chat.id, parent.id, path.resolve(dir)) + const mdl = yield* provider.getModel(ref.providerID, ref.modelID) + const handle = yield* processors.create({ + assistantMessage: msg, + sessionID: chat.id, + model: mdl, + }) + + const value = yield* handle.process({ + user: { + id: parent.id, + sessionID: chat.id, + role: "user", + time: parent.time, + agent: parent.agent, + model: { providerID: ref.providerID, modelID: ref.modelID }, + } satisfies MessageV2.User, + sessionID: chat.id, + model: mdl, + agent: agent(), + system: [], + messages: [{ role: "user", content: "retry abort" }], + tools: {}, + }) + + const parts = yield* MessageV2.parts(msg.id) + + expect(value).toBe("continue") + expect(unexpectedAbortLLM.calls()).toBe(2) + expect(parts.some((part) => part.type === "text" && part.text === "after")).toBe(true) + expect(handle.message.error).toBeUndefined() + }), + { config: cfg }, + ), +) + +const partialAbortLLM = llmMock( + abortStream( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.textStart({ id: "txt-0" }), + LLMEvent.textDelta({ id: "txt-0", text: "partial" }), + ), + assistantSuccessStream("after"), +) +const partialAbortIt = testEffect(processorEnv(partialAbortLLM.layer)) + +partialAbortIt.live("session.processor effect tests do not retry unexpected aborts after text output", () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const { processors, session, provider } = yield* boot() + + const chat = yield* session.create({}) + const parent = yield* user(chat.id, "partial abort") + const msg = yield* assistant(chat.id, parent.id, path.resolve(dir)) + const mdl = yield* provider.getModel(ref.providerID, ref.modelID) + const handle = yield* processors.create({ + assistantMessage: msg, + sessionID: chat.id, + model: mdl, + }) + + const value = yield* handle.process({ + user: { + id: parent.id, + sessionID: chat.id, + role: "user", + time: parent.time, + agent: parent.agent, + model: { providerID: ref.providerID, modelID: ref.modelID }, + } satisfies MessageV2.User, + sessionID: chat.id, + model: mdl, + agent: agent(), + system: [], + messages: [{ role: "user", content: "partial abort" }], + tools: {}, + }) + + const parts = yield* MessageV2.parts(msg.id) + + expect(value).toBe("stop") + expect(partialAbortLLM.calls()).toBe(1) + expect(parts.some((part) => part.type === "text" && part.text === "partial")).toBe(true) + expect(handle.message.error?.name).toBe("UnexpectedProviderAbortError") + }), + { config: cfg }, + ), +) + +const postToolTimeoutRetryLLM = llmMock(delayedFirstEventStream(50), assistantSuccessStream("after-timeout")) +const postToolTimeoutRetryIt = testEffect(processorEnv(postToolTimeoutRetryLLM.layer)) + +postToolTimeoutRetryIt.live("session.processor effect tests retry stalled post-tool continuation once", () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const { processors, session, provider } = yield* boot() + + const chat = yield* session.create({}) + const parent = yield* user(chat.id, "post-tool") + const msg = yield* assistant(chat.id, parent.id, path.resolve(dir)) + const mdl = yield* provider.getModel(ref.providerID, ref.modelID) + const handle = yield* processors.create({ + assistantMessage: msg, + sessionID: chat.id, + model: mdl, + }) + + const value = yield* handle.process({ + user: { + id: parent.id, + sessionID: chat.id, + role: "user", + time: parent.time, + agent: parent.agent, + model: { providerID: ref.providerID, modelID: ref.modelID }, + } satisfies MessageV2.User, + sessionID: chat.id, + model: mdl, + agent: agent(), + system: [], + messages: [{ role: "user", content: "post-tool" }], + tools: {}, + internal: { + postToolContinuation: true, + postToolFirstEventTimeoutMs: 10, + }, + }) + + const parts = yield* MessageV2.parts(msg.id) + + expect(value).toBe("continue") + expect(postToolTimeoutRetryLLM.calls()).toBe(2) + expect(parts.some((part) => part.type === "text" && part.text === "after-timeout")).toBe(true) + }), + { config: cfg }, + ), +) + +const postToolTimeoutExhaustLLM = llmMock(delayedFirstEventStream(50)) +const postToolTimeoutExhaustIt = testEffect(processorEnv(postToolTimeoutExhaustLLM.layer)) + +postToolTimeoutExhaustIt.live("session.processor effect tests include retry attempt in post-tool timeout errors", () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const { processors, session, provider } = yield* boot() + + const chat = yield* session.create({}) + const parent = yield* user(chat.id, "post-tool fail") + const msg = yield* assistant(chat.id, parent.id, path.resolve(dir)) + const mdl = yield* provider.getModel(ref.providerID, ref.modelID) + yield* session.updatePart({ + id: PartID.ascending(), + messageID: msg.id, + sessionID: chat.id, + type: "text", + text: "started", + }) + const handle = yield* processors.create({ + assistantMessage: msg, + sessionID: chat.id, + model: mdl, + }) + + const value = yield* handle.process({ + user: { + id: parent.id, + sessionID: chat.id, + role: "user", + time: parent.time, + agent: parent.agent, + model: { providerID: ref.providerID, modelID: ref.modelID }, + } satisfies MessageV2.User, + sessionID: chat.id, + model: mdl, + agent: agent(), + system: [], + messages: [{ role: "user", content: "post-tool fail" }], + tools: {}, + internal: { + postToolContinuation: true, + postToolFirstEventTimeoutMs: 10, + }, + }) + + expect(value).toBe("stop") + expect(postToolTimeoutExhaustLLM.calls()).toBe(1) + const error = handle.message.error + expect(MessageV2.PostToolContinuationTimeoutError.isInstance(error)).toBe(true) + if (!MessageV2.PostToolContinuationTimeoutError.isInstance(error)) return + expect(error.data.message).toContain("No stream event within 10ms after tool continuation") + expect(error.data.message).toContain("retry attempt 0") + }), + { config: cfg }, + ), +) + +const emptyAssistantLLM = llmMock(Stream.empty, assistantSuccessStream("after-empty")) +const emptyAssistantIt = testEffect(processorEnv(emptyAssistantLLM.layer)) + +emptyAssistantIt.live("session.processor effect tests retry zero-part assistant finalization", () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const { processors, session, provider } = yield* boot() + + const chat = yield* session.create({}) + const parent = yield* user(chat.id, "empty assistant") + const msg = yield* assistant(chat.id, parent.id, path.resolve(dir)) + const mdl = yield* provider.getModel(ref.providerID, ref.modelID) + const handle = yield* processors.create({ + assistantMessage: msg, + sessionID: chat.id, + model: mdl, + }) + + const value = yield* handle.process({ + user: { + id: parent.id, + sessionID: chat.id, + role: "user", + time: parent.time, + agent: parent.agent, + model: { providerID: ref.providerID, modelID: ref.modelID }, + } satisfies MessageV2.User, + sessionID: chat.id, + model: mdl, + agent: agent(), + system: [], + messages: [{ role: "user", content: "empty assistant" }], + tools: {}, + }) + + const parts = yield* MessageV2.parts(msg.id) + + expect(value).toBe("continue") + expect(emptyAssistantLLM.calls()).toBe(2) + expect(parts.some((part) => part.type === "text" && part.text === "after-empty")).toBe(true) + }), + { config: cfg }, + ), +) + +const subagentEmptyAssistantLLM = llmMock(Stream.empty, Stream.empty) +const subagentEmptyAssistantIt = testEffect(processorEnv(subagentEmptyAssistantLLM.layer)) + +subagentEmptyAssistantIt.live("session.processor effect tests limit child zero-part assistant retries", () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const { processors, session, provider } = yield* boot() + + const parentChat = yield* session.create({}) + const chat = yield* session.create({}) + const parent = yield* user(chat.id, "empty assistant") + const msg = yield* assistant(chat.id, parent.id, path.resolve(dir)) + const mdl = yield* provider.getModel(ref.providerID, ref.modelID) + const handle = yield* processors.create({ + assistantMessage: msg, + sessionID: chat.id, + model: mdl, + }) + + const value = yield* handle.process({ + user: { + id: parent.id, + sessionID: chat.id, + role: "user", + time: parent.time, + agent: parent.agent, + model: { providerID: ref.providerID, modelID: ref.modelID }, + } satisfies MessageV2.User, + sessionID: chat.id, + parentSessionID: parentChat.id, + model: mdl, + agent: agent(), + system: [], + messages: [{ role: "user", content: "empty assistant" }], + tools: {}, + }) + + expect(value).toBe("stop") + expect(subagentEmptyAssistantLLM.calls()).toBe(2) + expect(yield* MessageV2.parts(msg.id)).toHaveLength(0) + expect(handle.message.error?.name).toBe("EmptyAssistantResponseError") + }), + { config: cfg }, + ), +) + +const structuralEmptyAssistantLLM = llmMock( + structuralEmptyAssistantStream(), + assistantSuccessStream("after-structural"), +) +const structuralEmptyAssistantIt = testEffect(processorEnv(structuralEmptyAssistantLLM.layer)) + +structuralEmptyAssistantIt.live("session.processor effect tests retry structural-only assistant finalization", () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const { processors, session, provider } = yield* boot() + + const chat = yield* session.create({}) + const parent = yield* user(chat.id, "empty assistant") + const msg = yield* assistant(chat.id, parent.id, path.resolve(dir)) + const mdl = yield* provider.getModel(ref.providerID, ref.modelID) + const handle = yield* processors.create({ + assistantMessage: msg, + sessionID: chat.id, + model: mdl, + }) + + const value = yield* handle.process({ + user: { + id: parent.id, + sessionID: chat.id, + role: "user", + time: parent.time, + agent: parent.agent, + model: { providerID: ref.providerID, modelID: ref.modelID }, + } satisfies MessageV2.User, + sessionID: chat.id, + model: mdl, + agent: agent(), + system: [], + messages: [{ role: "user", content: "empty assistant" }], + tools: {}, + }) + + const parts = yield* MessageV2.parts(msg.id) + + expect(value).toBe("continue") + expect(structuralEmptyAssistantLLM.calls()).toBe(2) + expect(parts.some((part) => part.type === "text" && part.text === "after-structural")).toBe(true) + }), + { config: cfg }, + ), +) + +const subagentPostToolTimeoutLLM = llmMock(...Array.from({ length: 11 }, () => delayedFirstEventStream(50))) +const subagentPostToolTimeoutIt = testEffect(processorEnv(subagentPostToolTimeoutLLM.layer)) + +subagentPostToolTimeoutIt.live( + "session.processor effect tests limit child post-tool continuation timeout retries", + () => + provideTmpdirInstance( + (dir) => + Effect.gen(function* () { + const { processors, session, provider } = yield* boot() + + const parentChat = yield* session.create({}) + const chat = yield* session.create({}) + const parent = yield* user(chat.id, "post-tool child fail") + const msg = yield* assistant(chat.id, parent.id, path.resolve(dir)) + const mdl = yield* provider.getModel(ref.providerID, ref.modelID) + const handle = yield* processors.create({ + assistantMessage: msg, + sessionID: chat.id, + model: mdl, + }) + + const value = yield* handle.process({ + user: { + id: parent.id, + sessionID: chat.id, + role: "user", + time: parent.time, + agent: parent.agent, + model: { providerID: ref.providerID, modelID: ref.modelID }, + } satisfies MessageV2.User, + sessionID: chat.id, + parentSessionID: parentChat.id, + model: mdl, + agent: agent(), + system: [], + messages: [{ role: "user", content: "post-tool child fail" }], + tools: {}, + internal: { + postToolContinuation: true, + postToolFirstEventTimeoutMs: 10, + }, + }) + + expect(value).toBe("stop") + expect(subagentPostToolTimeoutLLM.calls()).toBe(2) + expect(handle.message.error?.name).toBe("PostToolContinuationTimeoutError") + }), + { config: cfg }, + ), +) + it.live("session.processor effect tests publish retry status updates", () => provideTmpdirServer( ({ dir, llm }) => diff --git a/packages/opencode/test/session/retry.test.ts b/packages/opencode/test/session/retry.test.ts index e53a6c1f18ee..aaa2ba21eb2d 100644 --- a/packages/opencode/test/session/retry.test.ts +++ b/packages/opencode/test/session/retry.test.ts @@ -32,6 +32,19 @@ function wrap(message: unknown): ReturnType { return { name: "", data: { message } } } +function abortedError(message = "The operation was aborted."): ReturnType { + return new MessageV2.AbortedError({ message }).toObject() +} + +function unexpectedProviderAbortError(message = "The operation was aborted."): ReturnType { + return new MessageV2.UnexpectedProviderAbortError({ + message, + abortSource: "provider_abort", + phase: "model_stream", + retryable: true, + }).toObject() +} + describe("session.retry.delay", () => { test("caps delay at 30 seconds when headers missing", () => { const error = apiError() @@ -191,6 +204,158 @@ describe("session.retry.retryable", () => { expect(SessionRetry.retryable(error, retryProvider)).toBeUndefined() }) + test("retries unexpected aborted errors before output starts", () => { + const error = abortedError() + + expect(SessionRetry.retryable(error, retryProvider, { aborted: false, empty: true })).toEqual({ + message: "The operation was aborted.", + }) + }) + + test("does not retry repeated unexpected aborted errors", () => { + const error = abortedError() + + expect( + SessionRetry.retryable(error, retryProvider, { + aborted: false, + empty: true, + attempt: SessionRetry.RETRY_UNEXPECTED_ABORT_LIMIT + 1, + }), + ).toBeUndefined() + }) + + test("does not retry intentional aborts", () => { + const error = abortedError() + + expect(SessionRetry.retryable(error, retryProvider, { aborted: true, empty: true })).toBeUndefined() + }) + + test("does not retry aborts with user provenance", () => { + const error = new MessageV2.AbortedError({ message: "Aborted", abortSource: "user_cancel" }).toObject() + + expect(SessionRetry.retryable(error, retryProvider, { aborted: false, empty: true })).toBeUndefined() + }) + + test("retries unexpected provider aborts before output starts", () => { + const error = unexpectedProviderAbortError() + + expect(SessionRetry.retryable(error, retryProvider, { aborted: false, empty: true })).toEqual({ + message: "The operation was aborted.", + }) + }) + + test("does not retry aborts after output starts", () => { + const error = abortedError() + + expect(SessionRetry.retryable(error, retryProvider, { aborted: false, empty: false })).toBeUndefined() + }) + + test("does not retry explicit user-cancel aborted errors", () => { + const error = new MessageV2.AbortedError({ + message: "Aborted", + abortSource: "user_cancel", + }).toObject() + + expect(SessionRetry.retryable(error, retryProvider, { aborted: true, empty: true })).toBeUndefined() + }) + + test("main sessions retry no-response errors up to ten attempts", () => { + const errors = [ + new MessageV2.PostToolContinuationTimeoutError({ + message: "No stream event", + abortSource: "post_tool_first_event_timeout", + phase: "post_tool_continuation", + retryable: true, + }).toObject(), + new MessageV2.EmptyAssistantResponseError({ + message: "Assistant stream ended without content", + abortSource: "unknown", + phase: "message_finalization", + retryable: true, + }).toObject(), + new MessageV2.NoResponseError({ + message: "Provider returned no response", + abortSource: "unknown", + phase: "message_finalization", + retryable: true, + }).toObject(), + ] + + for (const error of errors) { + expect( + SessionRetry.retryable(error, retryProvider, { + aborted: false, + empty: true, + postToolContinuation: true, + attempt: 10, + }), + ).toEqual({ message: error.data.message }) + expect( + SessionRetry.retryable(error, retryProvider, { + aborted: false, + empty: true, + postToolContinuation: true, + attempt: 11, + }), + ).toBeUndefined() + } + }) + + test("subagent sessions retry no-response errors once", () => { + const errors = [ + new MessageV2.PostToolContinuationTimeoutError({ + message: "No stream event", + abortSource: "post_tool_first_event_timeout", + phase: "post_tool_continuation", + retryable: true, + }).toObject(), + new MessageV2.EmptyAssistantResponseError({ + message: "Assistant stream ended without content", + abortSource: "unknown", + phase: "message_finalization", + retryable: true, + }).toObject(), + new MessageV2.NoResponseError({ + message: "Provider returned no response", + abortSource: "unknown", + phase: "message_finalization", + retryable: true, + }).toObject(), + ] + + for (const error of errors) { + expect( + SessionRetry.retryable(error, retryProvider, { + aborted: false, + empty: true, + postToolContinuation: true, + subagent: true, + attempt: 1, + }), + ).toEqual({ message: error.data.message }) + expect( + SessionRetry.retryable(error, retryProvider, { + aborted: false, + empty: true, + postToolContinuation: true, + subagent: true, + attempt: 2, + }), + ).toBeUndefined() + } + }) + + test("does not retry post-tool timeout outside post-tool continuation", () => { + const error = new MessageV2.PostToolContinuationTimeoutError({ + message: "No stream event", + abortSource: "post_tool_first_event_timeout", + phase: "post_tool_continuation", + retryable: true, + }).toObject() + + expect(SessionRetry.retryable(error, retryProvider, { aborted: false, empty: true, postToolContinuation: false })).toBeUndefined() + }) + test("retries 500 errors even when isRetryable is false", () => { const error = Schema.decodeUnknownSync(SessionV1.APIError.Schema)( new SessionV1.APIError({ diff --git a/packages/sdk/js/src/v2/gen/types.gen.ts b/packages/sdk/js/src/v2/gen/types.gen.ts index ef1c0142fb32..cd11fb1e1693 100644 --- a/packages/sdk/js/src/v2/gen/types.gen.ts +++ b/packages/sdk/js/src/v2/gen/types.gen.ts @@ -284,6 +284,162 @@ export type MessageAbortedError = { name: "MessageAbortedError" data: { message: string + abortSource?: + | "user_cancel" + | "session_cancel" + | "provider_abort" + | "network_abort" + | "first_byte_timeout" + | "stream_idle_timeout" + | "post_tool_first_event_timeout" + | "no_visible_part_timeout" + | "server_restart" + | "client_disconnect" + | "unknown" + } +} + +export type UnexpectedProviderAbortError = { + name: "UnexpectedProviderAbortError" + data: { + message: string + abortSource: + | "user_cancel" + | "session_cancel" + | "provider_abort" + | "network_abort" + | "first_byte_timeout" + | "stream_idle_timeout" + | "post_tool_first_event_timeout" + | "no_visible_part_timeout" + | "server_restart" + | "client_disconnect" + | "unknown" + phase: "model_stream" | "post_tool_continuation" | "message_finalization" | "unknown" + retryable: boolean + diagnostics?: { + providerID?: string + modelID?: string + sessionID?: string + messageID?: string + elapsedMs?: number + isPostToolContinuation?: boolean + retryAttempt?: number + firstStreamEventAt?: number + lastStreamEventAt?: number + firstVisiblePartAt?: number + lastVisiblePartAt?: number + partCount?: number + tokenCount?: number + } + } +} + +export type PostToolContinuationTimeoutError = { + name: "PostToolContinuationTimeoutError" + data: { + message: string + abortSource: + | "user_cancel" + | "session_cancel" + | "provider_abort" + | "network_abort" + | "first_byte_timeout" + | "stream_idle_timeout" + | "post_tool_first_event_timeout" + | "no_visible_part_timeout" + | "server_restart" + | "client_disconnect" + | "unknown" + phase: "model_stream" | "post_tool_continuation" | "message_finalization" | "unknown" + retryable: boolean + diagnostics?: { + providerID?: string + modelID?: string + sessionID?: string + messageID?: string + elapsedMs?: number + isPostToolContinuation?: boolean + retryAttempt?: number + firstStreamEventAt?: number + lastStreamEventAt?: number + firstVisiblePartAt?: number + lastVisiblePartAt?: number + partCount?: number + tokenCount?: number + } + } +} + +export type EmptyAssistantResponseError = { + name: "EmptyAssistantResponseError" + data: { + message: string + abortSource: + | "user_cancel" + | "session_cancel" + | "provider_abort" + | "network_abort" + | "first_byte_timeout" + | "stream_idle_timeout" + | "post_tool_first_event_timeout" + | "no_visible_part_timeout" + | "server_restart" + | "client_disconnect" + | "unknown" + phase: "model_stream" | "post_tool_continuation" | "message_finalization" | "unknown" + retryable: boolean + diagnostics?: { + providerID?: string + modelID?: string + sessionID?: string + messageID?: string + elapsedMs?: number + isPostToolContinuation?: boolean + retryAttempt?: number + firstStreamEventAt?: number + lastStreamEventAt?: number + firstVisiblePartAt?: number + lastVisiblePartAt?: number + partCount?: number + tokenCount?: number + } + } +} + +export type NoResponseError = { + name: "NoResponseError" + data: { + message: string + abortSource: + | "user_cancel" + | "session_cancel" + | "provider_abort" + | "network_abort" + | "first_byte_timeout" + | "stream_idle_timeout" + | "post_tool_first_event_timeout" + | "no_visible_part_timeout" + | "server_restart" + | "client_disconnect" + | "unknown" + phase: "model_stream" | "post_tool_continuation" | "message_finalization" | "unknown" + retryable: boolean + diagnostics?: { + providerID?: string + modelID?: string + sessionID?: string + messageID?: string + elapsedMs?: number + isPostToolContinuation?: boolean + retryAttempt?: number + firstStreamEventAt?: number + lastStreamEventAt?: number + firstVisiblePartAt?: number + lastVisiblePartAt?: number + partCount?: number + tokenCount?: number + } } } @@ -339,6 +495,10 @@ export type AssistantMessage = { | UnknownError | MessageOutputLengthError | MessageAbortedError + | UnexpectedProviderAbortError + | PostToolContinuationTimeoutError + | EmptyAssistantResponseError + | NoResponseError | StructuredOutputError | ContextOverflowError | ContentFilterError @@ -1209,6 +1369,10 @@ export type GlobalEvent = { | UnknownError | MessageOutputLengthError | MessageAbortedError + | UnexpectedProviderAbortError + | PostToolContinuationTimeoutError + | EmptyAssistantResponseError + | NoResponseError | StructuredOutputError | ContextOverflowError | ContentFilterError @@ -6487,6 +6651,10 @@ export type EventSessionError = { | UnknownError | MessageOutputLengthError | MessageAbortedError + | UnexpectedProviderAbortError + | PostToolContinuationTimeoutError + | EmptyAssistantResponseError + | NoResponseError | StructuredOutputError | ContextOverflowError | ContentFilterError diff --git a/packages/sdk/openapi.json b/packages/sdk/openapi.json index 932cb8a49549..e15e8b3acae9 100644 --- a/packages/sdk/openapi.json +++ b/packages/sdk/openapi.json @@ -15478,6 +15478,22 @@ "properties": { "message": { "type": "string" + }, + "abortSource": { + "type": "string", + "enum": [ + "user_cancel", + "session_cancel", + "provider_abort", + "network_abort", + "first_byte_timeout", + "stream_idle_timeout", + "post_tool_first_event_timeout", + "no_visible_part_timeout", + "server_restart", + "client_disconnect", + "unknown" + ] } }, "required": ["message"], @@ -15487,6 +15503,402 @@ "required": ["name", "data"], "additionalProperties": false }, + "UnexpectedProviderAbortError": { + "type": "object", + "properties": { + "name": { + "type": "string", + "enum": ["UnexpectedProviderAbortError"] + }, + "data": { + "type": "object", + "properties": { + "message": { + "type": "string" + }, + "abortSource": { + "type": "string", + "enum": [ + "user_cancel", + "session_cancel", + "provider_abort", + "network_abort", + "first_byte_timeout", + "stream_idle_timeout", + "post_tool_first_event_timeout", + "no_visible_part_timeout", + "server_restart", + "client_disconnect", + "unknown" + ] + }, + "phase": { + "type": "string", + "enum": ["model_stream", "post_tool_continuation", "message_finalization", "unknown"] + }, + "retryable": { + "type": "boolean" + }, + "diagnostics": { + "type": "object", + "properties": { + "providerID": { + "type": "string" + }, + "modelID": { + "type": "string" + }, + "sessionID": { + "type": "string", + "pattern": "^ses" + }, + "messageID": { + "type": "string", + "pattern": "^msg" + }, + "elapsedMs": { + "type": "integer", + "minimum": 0 + }, + "isPostToolContinuation": { + "type": "boolean" + }, + "retryAttempt": { + "type": "integer", + "minimum": 0 + }, + "firstStreamEventAt": { + "type": "integer", + "minimum": 0 + }, + "lastStreamEventAt": { + "type": "integer", + "minimum": 0 + }, + "firstVisiblePartAt": { + "type": "integer", + "minimum": 0 + }, + "lastVisiblePartAt": { + "type": "integer", + "minimum": 0 + }, + "partCount": { + "type": "integer", + "minimum": 0 + }, + "tokenCount": { + "type": "integer", + "minimum": 0 + } + }, + "additionalProperties": false + } + }, + "required": ["message", "abortSource", "phase", "retryable"], + "additionalProperties": false + } + }, + "required": ["name", "data"], + "additionalProperties": false + }, + "PostToolContinuationTimeoutError": { + "type": "object", + "properties": { + "name": { + "type": "string", + "enum": ["PostToolContinuationTimeoutError"] + }, + "data": { + "type": "object", + "properties": { + "message": { + "type": "string" + }, + "abortSource": { + "type": "string", + "enum": [ + "user_cancel", + "session_cancel", + "provider_abort", + "network_abort", + "first_byte_timeout", + "stream_idle_timeout", + "post_tool_first_event_timeout", + "no_visible_part_timeout", + "server_restart", + "client_disconnect", + "unknown" + ] + }, + "phase": { + "type": "string", + "enum": ["model_stream", "post_tool_continuation", "message_finalization", "unknown"] + }, + "retryable": { + "type": "boolean" + }, + "diagnostics": { + "type": "object", + "properties": { + "providerID": { + "type": "string" + }, + "modelID": { + "type": "string" + }, + "sessionID": { + "type": "string", + "pattern": "^ses" + }, + "messageID": { + "type": "string", + "pattern": "^msg" + }, + "elapsedMs": { + "type": "integer", + "minimum": 0 + }, + "isPostToolContinuation": { + "type": "boolean" + }, + "retryAttempt": { + "type": "integer", + "minimum": 0 + }, + "firstStreamEventAt": { + "type": "integer", + "minimum": 0 + }, + "lastStreamEventAt": { + "type": "integer", + "minimum": 0 + }, + "firstVisiblePartAt": { + "type": "integer", + "minimum": 0 + }, + "lastVisiblePartAt": { + "type": "integer", + "minimum": 0 + }, + "partCount": { + "type": "integer", + "minimum": 0 + }, + "tokenCount": { + "type": "integer", + "minimum": 0 + } + }, + "additionalProperties": false + } + }, + "required": ["message", "abortSource", "phase", "retryable"], + "additionalProperties": false + } + }, + "required": ["name", "data"], + "additionalProperties": false + }, + "EmptyAssistantResponseError": { + "type": "object", + "properties": { + "name": { + "type": "string", + "enum": ["EmptyAssistantResponseError"] + }, + "data": { + "type": "object", + "properties": { + "message": { + "type": "string" + }, + "abortSource": { + "type": "string", + "enum": [ + "user_cancel", + "session_cancel", + "provider_abort", + "network_abort", + "first_byte_timeout", + "stream_idle_timeout", + "post_tool_first_event_timeout", + "no_visible_part_timeout", + "server_restart", + "client_disconnect", + "unknown" + ] + }, + "phase": { + "type": "string", + "enum": ["model_stream", "post_tool_continuation", "message_finalization", "unknown"] + }, + "retryable": { + "type": "boolean" + }, + "diagnostics": { + "type": "object", + "properties": { + "providerID": { + "type": "string" + }, + "modelID": { + "type": "string" + }, + "sessionID": { + "type": "string", + "pattern": "^ses" + }, + "messageID": { + "type": "string", + "pattern": "^msg" + }, + "elapsedMs": { + "type": "integer", + "minimum": 0 + }, + "isPostToolContinuation": { + "type": "boolean" + }, + "retryAttempt": { + "type": "integer", + "minimum": 0 + }, + "firstStreamEventAt": { + "type": "integer", + "minimum": 0 + }, + "lastStreamEventAt": { + "type": "integer", + "minimum": 0 + }, + "firstVisiblePartAt": { + "type": "integer", + "minimum": 0 + }, + "lastVisiblePartAt": { + "type": "integer", + "minimum": 0 + }, + "partCount": { + "type": "integer", + "minimum": 0 + }, + "tokenCount": { + "type": "integer", + "minimum": 0 + } + }, + "additionalProperties": false + } + }, + "required": ["message", "abortSource", "phase", "retryable"], + "additionalProperties": false + } + }, + "required": ["name", "data"], + "additionalProperties": false + }, + "NoResponseError": { + "type": "object", + "properties": { + "name": { + "type": "string", + "enum": ["NoResponseError"] + }, + "data": { + "type": "object", + "properties": { + "message": { + "type": "string" + }, + "abortSource": { + "type": "string", + "enum": [ + "user_cancel", + "session_cancel", + "provider_abort", + "network_abort", + "first_byte_timeout", + "stream_idle_timeout", + "post_tool_first_event_timeout", + "no_visible_part_timeout", + "server_restart", + "client_disconnect", + "unknown" + ] + }, + "phase": { + "type": "string", + "enum": ["model_stream", "post_tool_continuation", "message_finalization", "unknown"] + }, + "retryable": { + "type": "boolean" + }, + "diagnostics": { + "type": "object", + "properties": { + "providerID": { + "type": "string" + }, + "modelID": { + "type": "string" + }, + "sessionID": { + "type": "string", + "pattern": "^ses" + }, + "messageID": { + "type": "string", + "pattern": "^msg" + }, + "elapsedMs": { + "type": "integer", + "minimum": 0 + }, + "isPostToolContinuation": { + "type": "boolean" + }, + "retryAttempt": { + "type": "integer", + "minimum": 0 + }, + "firstStreamEventAt": { + "type": "integer", + "minimum": 0 + }, + "lastStreamEventAt": { + "type": "integer", + "minimum": 0 + }, + "firstVisiblePartAt": { + "type": "integer", + "minimum": 0 + }, + "lastVisiblePartAt": { + "type": "integer", + "minimum": 0 + }, + "partCount": { + "type": "integer", + "minimum": 0 + }, + "tokenCount": { + "type": "integer", + "minimum": 0 + } + }, + "additionalProperties": false + } + }, + "required": ["message", "abortSource", "phase", "retryable"], + "additionalProperties": false + } + }, + "required": ["name", "data"], + "additionalProperties": false + }, "StructuredOutputError": { "type": "object", "properties": { @@ -15644,6 +16056,18 @@ { "$ref": "#/components/schemas/MessageAbortedError" }, + { + "$ref": "#/components/schemas/UnexpectedProviderAbortError" + }, + { + "$ref": "#/components/schemas/PostToolContinuationTimeoutError" + }, + { + "$ref": "#/components/schemas/EmptyAssistantResponseError" + }, + { + "$ref": "#/components/schemas/NoResponseError" + }, { "$ref": "#/components/schemas/StructuredOutputError" }, @@ -18421,6 +18845,18 @@ { "$ref": "#/components/schemas/MessageAbortedError" }, + { + "$ref": "#/components/schemas/UnexpectedProviderAbortError" + }, + { + "$ref": "#/components/schemas/PostToolContinuationTimeoutError" + }, + { + "$ref": "#/components/schemas/EmptyAssistantResponseError" + }, + { + "$ref": "#/components/schemas/NoResponseError" + }, { "$ref": "#/components/schemas/StructuredOutputError" }, @@ -34274,6 +34710,18 @@ { "$ref": "#/components/schemas/MessageAbortedError" }, + { + "$ref": "#/components/schemas/UnexpectedProviderAbortError" + }, + { + "$ref": "#/components/schemas/PostToolContinuationTimeoutError" + }, + { + "$ref": "#/components/schemas/EmptyAssistantResponseError" + }, + { + "$ref": "#/components/schemas/NoResponseError" + }, { "$ref": "#/components/schemas/StructuredOutputError" },