From bbc3ae4fe1d42c732d5192c9c5e0845ed9c78828 Mon Sep 17 00:00:00 2001 From: soshymking Date: Tue, 26 May 2026 18:55:34 +0900 Subject: [PATCH 1/6] dzianisv:issue-17717-retry-unexpected-aborts --- packages/opencode/src/session/processor.ts | 7 + packages/opencode/src/session/retry.ts | 18 +- .../test/session/processor-effect.test.ts | 160 +++++++++++++++++- packages/opencode/test/session/retry.test.ts | 36 ++++ 4 files changed, 216 insertions(+), 5 deletions(-) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index a287c3b00680..b7c535248679 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -777,6 +777,9 @@ export const layer = Layer.effect( yield* status.set(ctx.sessionID, { type: "idle" }) }) + const assistantOutputEmpty = () => + MessageV2.parts(ctx.assistantMessage.id).every((part) => part.type === "step-start") + const process = Effect.fn("SessionProcessor.process")(function* (streamInput: LLM.StreamInput) { slog.info("process") ctx.needsCompaction = false @@ -811,6 +814,10 @@ export const layer = Layer.effect( SessionRetry.policy({ provider: input.model.providerID, parse, + context: () => ({ + aborted, + empty: assistantOutputEmpty(), + }), set: (info) => { // TODO(v2): Temporary dual-write while migrating session messages to v2 events. const event = flags.experimentalEventSystem diff --git a/packages/opencode/src/session/retry.ts b/packages/opencode/src/session/retry.ts index 463bc27a95db..90febfeba76b 100644 --- a/packages/opencode/src/session/retry.ts +++ b/packages/opencode/src/session/retry.ts @@ -10,6 +10,12 @@ 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 +} + export type Retryable = { message: string action?: { @@ -26,6 +32,7 @@ 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 function cap(ms: number) { return Math.min(ms, RETRY_MAX_DELAY) @@ -64,9 +71,15 @@ export function delay(attempt: number, error?: MessageV2.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 (MessageV2.ContextOverflowError.isInstance(error)) return undefined + if (MessageV2.AbortedError.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.APIError.isInstance(error)) { const status = error.data.statusCode // 5xx errors are transient server failures and should always be retried, @@ -176,11 +189,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, MessageV2.APIError.isInstance(error) ? error : undefined) diff --git a/packages/opencode/test/session/processor-effect.test.ts b/packages/opencode/test/session/processor-effect.test.ts index ede122297a17..3a9f003af697 100644 --- a/packages/opencode/test/session/processor-effect.test.ts +++ b/packages/opencode/test/session/processor-effect.test.ts @@ -2,6 +2,7 @@ import { NodeFileSystem } from "@effect/platform-node" import { expect } from "bun:test" import { tool } from "ai" 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" @@ -23,12 +24,13 @@ import { SessionSummary } from "../../src/session/summary" import { Snapshot } from "../../src/snapshot" import * as Log from "@opencode-ai/core/util/log" import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner" -import { provideTmpdirServer } from "../fixture/fixture" +import { provideTmpdirInstance, provideTmpdirServer } from "../fixture/fixture" import { testEffect } from "../lib/effect" import { raw, reply, TestLLMServer } from "../lib/llm-server" import { SyncEvent } from "@/sync" import { RuntimeFlags } from "@/effect/runtime-flags" import { EventV2Bridge } from "@/event-v2-bridge" +import { LLMEvent, Usage } from "@opencode-ai/llm" void Log.init({ print: false }) @@ -172,19 +174,19 @@ const assistant = Effect.fn("TestSession.assistant")(function* ( const status = SessionStatus.layer.pipe(Layer.provideMerge(Bus.layer)) const infra = Layer.mergeAll(NodeFileSystem.layer, CrossSpawnSpawner.defaultLayer) -const deps = Layer.mergeAll( +const depsWithoutLLM = Layer.mergeAll( Session.defaultLayer, Snapshot.defaultLayer, AgentSvc.defaultLayer, Permission.defaultLayer, Plugin.defaultLayer, Config.defaultLayer, - LLM.defaultLayer, Provider.defaultLayer, status, SyncEvent.defaultLayer, EventV2Bridge.defaultLayer, ).pipe(Layer.provideMerge(infra)) +const deps = Layer.mergeAll(LLM.defaultLayer, depsWithoutLLM) const env = Layer.mergeAll( TestLLMServer.layer, SessionProcessor.layer.pipe( @@ -204,6 +206,57 @@ 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 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 SessionProcessor.layer.pipe( + Layer.provide(summary), + Layer.provide(Image.defaultLayer), + Layer.provide(RuntimeFlags.layer({ experimentalEventSystem: true })), + + Layer.provideMerge(depsWithoutLLM), + Layer.provide(llmLayer), + ) +} // --------------------------------------------------------------------------- // Tests // --------------------------------------------------------------------------- @@ -568,6 +621,107 @@ 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 = 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 = 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("MessageAbortedError") + }), + { 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 22ff6cde811d..0f9d1b309360 100644 --- a/packages/opencode/test/session/retry.test.ts +++ b/packages/opencode/test/session/retry.test.ts @@ -30,6 +30,10 @@ function wrap(message: unknown): ReturnType { return { name: "", data: { message } } } +function abortedError(message = "The operation was aborted."): ReturnType { + return new MessageV2.AbortedError({ message }).toObject() +} + describe("session.retry.delay", () => { test("caps delay at 30 seconds when headers missing", () => { const error = apiError() @@ -172,6 +176,38 @@ 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 after output starts", () => { + const error = abortedError() + + expect(SessionRetry.retryable(error, retryProvider, { aborted: false, empty: false })).toBeUndefined() + }) + test("retries 500 errors even when isRetryable is false", () => { const error = Schema.decodeUnknownSync(MessageV2.APIError.Schema)( new MessageV2.APIError({ From ee4e0637201712b2a435a6194f11e70c041a5c6a Mon Sep 17 00:00:00 2001 From: soshymking Date: Wed, 27 May 2026 07:41:28 +0900 Subject: [PATCH 2/6] no response handling --- packages/opencode/src/session/llm.ts | 4 + packages/opencode/src/session/message-v2.ts | 99 ++++++++++- packages/opencode/src/session/processor.ts | 162 +++++++++++++++++- packages/opencode/src/session/prompt.ts | 13 ++ packages/opencode/src/session/retry.ts | 21 +++ .../opencode/test/session/message-v2.test.ts | 76 ++++++++ .../test/session/processor-effect.test.ts | 157 ++++++++++++++++- packages/opencode/test/session/retry.test.ts | 59 +++++++ 8 files changed, 581 insertions(+), 10 deletions(-) diff --git a/packages/opencode/src/session/llm.ts b/packages/opencode/src/session/llm.ts index ea2efc99d007..9ca58347b478 100644 --- a/packages/opencode/src/session/llm.ts +++ b/packages/opencode/src/session/llm.ts @@ -43,6 +43,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 2745ff4f45d7..c13e653cb7de 100644 --- a/packages/opencode/src/session/message-v2.ts +++ b/packages/opencode/src/session/message-v2.ts @@ -38,7 +38,65 @@ interface FetchDecompressionError extends Error { export const SYNTHETIC_ATTACHMENT_PROMPT = "Attached media from tool result:" export { isMedia } -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 + +const NoResponseDiagnostics = Schema.Struct({ + providerID: Schema.optional(ProviderID), + modelID: Schema.optional(ModelID), + sessionID: Schema.optional(SessionID), + 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, @@ -380,6 +438,10 @@ export type Part = const AssistantErrorSchema = Schema.Union([ ...MessageError.Shared, AbortedError.EffectSchema, + UnexpectedProviderAbortError.EffectSchema, + PostToolContinuationTimeoutError.EffectSchema, + EmptyAssistantResponseError.EffectSchema, + NoResponseError.EffectSchema, StructuredOutputError.EffectSchema, ContextOverflowError.EffectSchema, APIError.EffectSchema, @@ -746,7 +808,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") ) ) { @@ -1095,12 +1157,37 @@ export function latest(msgs: WithParts[]) { export function fromError( e: unknown, - ctx: { providerID: ProviderID; aborted?: boolean }, + ctx: { + providerID: ProviderID + 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, }, @@ -1130,7 +1217,7 @@ 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 b7c535248679..4b2725e3d8fc 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -30,6 +30,7 @@ import { RuntimeFlags } from "@/effect/runtime-flags" import { Usage, type LLMEvent } from "@opencode-ai/llm" const DOOM_LOOP_THRESHOLD = 3 +export const POST_TOOL_FIRST_EVENT_TIMEOUT_MS = 10_000 const log = Log.create({ service: "session.processor" }) export type Result = "compact" | "stop" | "continue" @@ -82,6 +83,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( @@ -120,12 +133,43 @@ export const layer = Layer.effect( reasoningMap: {}, } 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, + } const slog = log.clone().tag("session.id", input.sessionID).tag("messageID", input.assistantMessage.id) + 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) { @@ -302,10 +346,21 @@ 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++ + } + + const tokenTotal = (tokens: MessageV2.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 (flags.experimentalEventSystem) { yield* events.publish(SessionEvent.Reasoning.Started, { @@ -329,6 +384,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 yield* session.updatePartDelta({ @@ -351,6 +407,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 @@ -378,6 +435,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 = toolInput(value.input) if (!toolCall.call.inputEnded) { @@ -450,6 +508,7 @@ export const layer = Layer.effect( } case "tool-result": { + markVisiblePart(false) const toolCall = yield* readToolCall(value.id) const rawOutput = toolResultOutput(value) const normalized = yield* Effect.forEach(rawOutput.attachments ?? [], (attachment) => @@ -502,6 +561,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 (flags.experimentalEventSystem) { @@ -526,6 +586,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. @@ -576,6 +637,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, @@ -617,6 +680,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 (flags.experimentalEventSystem) { @@ -640,6 +704,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 yield* session.updatePartDelta({ @@ -653,6 +718,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( @@ -751,6 +817,16 @@ export const layer = Layer.effect( const halt = Effect.fn("SessionProcessor.halt")(function* (e: unknown) { slog.error("process", { error: errorMessage(e), stack: e instanceof Error ? e.stack : undefined }) const error = parse(e) + if (MessageV2.UnexpectedProviderAbortError.isInstance(error)) { + slog.warn("model.abort.unexpected_provider_abort", { + ...(error.data.diagnostics ?? diagnostics()), + abortSource: error.data.abortSource, + phase: error.data.phase, + }) + } + if (MessageV2.AbortedError.isInstance(error)) { + slog.info("model.abort.user_cancel", { ...diagnostics(), abortSource: error.data.abortSource ?? "user_cancel" }) + } if (MessageV2.ContextOverflowError.isInstance(error)) { ctx.needsCompaction = true yield* bus.publish(Session.Event.Error, { sessionID: ctx.sessionID, error }) @@ -784,19 +860,88 @@ export const layer = Layer.effect( slog.info("process") 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.reasoningMap = {} yield* status.set(ctx.sessionID, { type: "busy" }) 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.sync(() => { + slog.warn("model.no_response.post_tool_first_event_timeout", { + ...diagnostics(), + abortSource: "post_tool_first_event_timeout", + phase: "post_tool_continuation", + }) + return new MessageV2.PostToolContinuationTimeoutError({ + message: `No stream event within ${firstEventTimeoutMs}ms after tool continuation`, + abortSource: "post_tool_first_event_timeout", + phase: "post_tool_continuation", + retryable: true, + diagnostics: diagnostics(), + }) + }).pipe(Effect.flatMap((error) => Effect.fail(error))), + }), + Effect.andThen(Effect.never), + ) + : Effect.never + + yield* drain.pipe(Effect.raceFirst(postToolFirstEventTimeout)) + + if (ctx.assistantMessage.role === "assistant") { + const parts = MessageV2.parts(ctx.assistantMessage.id) + const noParts = parts.length === 0 || parts.every((part) => part.type === "step-start") + const noTokens = activity.tokenCount === 0 && ctx.assistantMessage.tokens.output === 0 + if (noParts && noTokens) { + slog.warn("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* () { @@ -817,8 +962,19 @@ export const layer = Layer.effect( context: () => ({ aborted, empty: assistantOutputEmpty(), + postToolContinuation: activity.isPostToolContinuation, }), set: (info) => { + retryAttempt = info.attempt + activity.retryAttempt = info.attempt + if (activity.isPostToolContinuation) { + slog.warn("model.no_response.retrying_continuation", { + ...diagnostics(), + abortSource: "post_tool_first_event_timeout", + phase: "post_tool_continuation", + attempt: info.attempt, + }) + } // TODO(v2): Temporary dual-write while migrating session messages to v2 events. const event = flags.experimentalEventSystem ? events.publish(SessionEvent.Retried, { diff --git a/packages/opencode/src/session/prompt.ts b/packages/opencode/src/session/prompt.ts index 22fe4d81cd40..2dee488dc098 100644 --- a/packages/opencode/src/session/prompt.ts +++ b/packages/opencode/src/session/prompt.ts @@ -87,6 +87,15 @@ function isOrphanedInterruptedTool(part: MessageV2.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 @@ -1441,6 +1450,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, @@ -1452,6 +1462,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 90febfeba76b..f45e6fec89c7 100644 --- a/packages/opencode/src/session/retry.ts +++ b/packages/opencode/src/session/retry.ts @@ -14,6 +14,7 @@ export type RetryContext = { aborted?: boolean empty?: boolean attempt?: number + postToolContinuation?: boolean } export type Retryable = { @@ -33,6 +34,7 @@ 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 function cap(ms: number) { return Math.min(ms, RETRY_MAX_DELAY) @@ -75,11 +77,30 @@ export function retryable(error: Err, provider: string, context?: RetryContext) // context overflow errors should not be retried if (MessageV2.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 + if (MessageV2.PostToolContinuationTimeoutError.isInstance(error) && context?.postToolContinuation !== true) + return undefined + if ((context?.attempt ?? 1) > RETRY_NO_RESPONSE_LIMIT) return undefined + return { message: error.data.message } + } if (MessageV2.APIError.isInstance(error)) { const status = error.data.statusCode // 5xx errors are transient server failures and should always be retried, diff --git a/packages/opencode/test/session/message-v2.test.ts b/packages/opencode/test/session/message-v2.test.ts index 82bed0e9cc6f..d8a8ad7538fc 100644 --- a/packages/opencode/test/session/message-v2.test.ts +++ b/packages/opencode/test/session/message-v2.test.ts @@ -1036,6 +1036,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 = { @@ -1546,6 +1585,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 3a9f003af697..09364ede5c07 100644 --- a/packages/opencode/test/session/processor-effect.test.ts +++ b/packages/opencode/test/session/processor-effect.test.ts @@ -231,6 +231,18 @@ function abortStream(...events: LLMEvent[]) { ) } +function delayedFirstEventStream(delayMs: number) { + return Stream.fromAsyncIterable( + { + async *[Symbol.asyncIterator]() { + await Bun.sleep(delayMs) + yield LLMEvent.stepStart({ index: 0 }) + }, + }, + (error) => error, + ) +} + function llmMock(...streams: Stream.Stream[]) { let calls = 0 return { @@ -716,7 +728,150 @@ partialAbortIt.live("session.processor effect tests do not retry unexpected abor 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("MessageAbortedError") + 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 = 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), delayedFirstEventStream(50)) +const postToolTimeoutExhaustIt = testEffect(processorEnv(postToolTimeoutExhaustLLM.layer)) + +postToolTimeoutExhaustIt.live("session.processor effect tests stop with typed timeout after post-tool retry exhaustion", () => + 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) + 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(2) + expect(handle.message.error?.name).toBe("PostToolContinuationTimeoutError") + }), + { config: cfg }, + ), +) + +const emptyAssistantLLM = llmMock(Stream.empty, Stream.empty) +const emptyAssistantIt = testEffect(processorEnv(emptyAssistantLLM.layer)) + +emptyAssistantIt.live("session.processor effect tests classify 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: {}, + }) + + expect(value).toBe("stop") + expect(emptyAssistantLLM.calls()).toBe(2) + expect(MessageV2.parts(msg.id)).toHaveLength(0) + expect(handle.message.error?.name).toBe("EmptyAssistantResponseError") }), { config: cfg }, ), diff --git a/packages/opencode/test/session/retry.test.ts b/packages/opencode/test/session/retry.test.ts index 0f9d1b309360..6b3082cf1f7f 100644 --- a/packages/opencode/test/session/retry.test.ts +++ b/packages/opencode/test/session/retry.test.ts @@ -34,6 +34,15 @@ function abortedError(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() @@ -202,12 +211,62 @@ describe("session.retry.retryable", () => { 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("retries post-tool continuation timeout once", () => { + 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: true, attempt: 1 }), + ).toEqual({ message: "No stream event" }) + expect( + SessionRetry.retryable(error, retryProvider, { aborted: false, empty: true, postToolContinuation: 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(MessageV2.APIError.Schema)( new MessageV2.APIError({ From da395f2eeb5a599ee3af65e878b7d16d51c9ad10 Mon Sep 17 00:00:00 2001 From: soshymking Date: Wed, 27 May 2026 16:29:40 +0900 Subject: [PATCH 3/6] fix empty SSE stream --- packages/opencode/src/session/processor.ts | 18 +- .../test/session/processor-effect.test.ts | 55 +++ packages/sdk/js/src/v2/gen/types.gen.ts | 164 +++++++ packages/sdk/openapi.json | 436 ++++++++++++++++++ 4 files changed, 669 insertions(+), 4 deletions(-) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 4b2725e3d8fc..de536854da97 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -30,7 +30,7 @@ import { RuntimeFlags } from "@/effect/runtime-flags" import { Usage, type LLMEvent } from "@opencode-ai/llm" const DOOM_LOOP_THRESHOLD = 3 -export const POST_TOOL_FIRST_EVENT_TIMEOUT_MS = 10_000 +export const POST_TOOL_FIRST_EVENT_TIMEOUT_MS = 100_000 const log = Log.create({ service: "session.processor" }) export type Result = "compact" | "stop" | "continue" @@ -853,8 +853,7 @@ export const layer = Layer.effect( yield* status.set(ctx.sessionID, { type: "idle" }) }) - const assistantOutputEmpty = () => - MessageV2.parts(ctx.assistantMessage.id).every((part) => part.type === "step-start") + const assistantOutputEmpty = () => MessageV2.parts(ctx.assistantMessage.id).every(isStructuralAssistantPart) const process = Effect.fn("SessionProcessor.process")(function* (streamInput: LLM.StreamInput) { slog.info("process") @@ -923,7 +922,7 @@ export const layer = Layer.effect( if (ctx.assistantMessage.role === "assistant") { const parts = MessageV2.parts(ctx.assistantMessage.id) - const noParts = parts.length === 0 || parts.every((part) => part.type === "step-start") + const noParts = isEmptyAssistantOutput(parts) const noTokens = activity.tokenCount === 0 && ctx.assistantMessage.tokens.output === 0 if (noParts && noTokens) { slog.warn("model.no_response.zero_part_assistant_turn", { @@ -1043,4 +1042,15 @@ 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 * as SessionProcessor from "./processor" diff --git a/packages/opencode/test/session/processor-effect.test.ts b/packages/opencode/test/session/processor-effect.test.ts index 09364ede5c07..5d386cd23cd9 100644 --- a/packages/opencode/test/session/processor-effect.test.ts +++ b/packages/opencode/test/session/processor-effect.test.ts @@ -243,6 +243,14 @@ function delayedFirstEventStream(delayMs: number) { ) } +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 { @@ -877,6 +885,53 @@ emptyAssistantIt.live("session.processor effect tests classify zero-part assista ), ) +const structuralEmptyAssistantLLM = llmMock(structuralEmptyAssistantStream(), structuralEmptyAssistantStream()) +const structuralEmptyAssistantIt = testEffect(processorEnv(structuralEmptyAssistantLLM.layer)) + +structuralEmptyAssistantIt.live("session.processor effect tests classify 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 = MessageV2.parts(msg.id) + + expect(value).toBe("stop") + expect(structuralEmptyAssistantLLM.calls()).toBe(2) + expect(parts.every((part) => part.type === "step-start" || part.type === "step-finish")).toBe(true) + expect(handle.message.error?.name).toBe("EmptyAssistantResponseError") + }), + { config: cfg }, + ), +) + it.live("session.processor effect tests publish retry status updates", () => provideTmpdirServer( ({ dir, llm }) => diff --git a/packages/sdk/js/src/v2/gen/types.gen.ts b/packages/sdk/js/src/v2/gen/types.gen.ts index 89cfc6559014..aebf1b0f4fda 100644 --- a/packages/sdk/js/src/v2/gen/types.gen.ts +++ b/packages/sdk/js/src/v2/gen/types.gen.ts @@ -224,6 +224,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 + } } } @@ -440,6 +596,10 @@ export type AssistantMessage = { | UnknownError | MessageOutputLengthError | MessageAbortedError + | UnexpectedProviderAbortError + | PostToolContinuationTimeoutError + | EmptyAssistantResponseError + | NoResponseError | StructuredOutputError | ContextOverflowError | ApiError @@ -2604,6 +2764,10 @@ export type EventSessionError = { | UnknownError | MessageOutputLengthError | MessageAbortedError + | UnexpectedProviderAbortError + | PostToolContinuationTimeoutError + | EmptyAssistantResponseError + | NoResponseError | StructuredOutputError | ContextOverflowError | ApiError diff --git a/packages/sdk/openapi.json b/packages/sdk/openapi.json index ebef9aea8db6..27760f4fdc51 100644 --- a/packages/sdk/openapi.json +++ b/packages/sdk/openapi.json @@ -11165,6 +11165,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"], @@ -11174,6 +11190,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": { @@ -11752,6 +12164,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" }, @@ -18406,6 +18830,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" }, From 2a3abb3d6901b8ba5f1c6ac78a34526f942c6464 Mon Sep 17 00:00:00 2001 From: soshymking Date: Thu, 28 May 2026 04:10:17 +0900 Subject: [PATCH 4/6] change main retry to 10 --- packages/opencode/src/session/processor.ts | 3 +- packages/opencode/src/session/retry.ts | 9 +- .../test/session/processor-effect.test.ts | 132 ++++++++++++++++-- packages/opencode/test/session/retry.test.ts | 96 +++++++++++-- 4 files changed, 211 insertions(+), 29 deletions(-) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index de536854da97..8a96a3ef74af 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -906,7 +906,7 @@ export const layer = Layer.effect( phase: "post_tool_continuation", }) return new MessageV2.PostToolContinuationTimeoutError({ - message: `No stream event within ${firstEventTimeoutMs}ms after tool continuation`, + 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, @@ -962,6 +962,7 @@ export const layer = Layer.effect( aborted, empty: assistantOutputEmpty(), postToolContinuation: activity.isPostToolContinuation, + subagent: streamInput.parentSessionID !== undefined, }), set: (info) => { retryAttempt = info.attempt diff --git a/packages/opencode/src/session/retry.ts b/packages/opencode/src/session/retry.ts index f45e6fec89c7..b187ed7d60f4 100644 --- a/packages/opencode/src/session/retry.ts +++ b/packages/opencode/src/session/retry.ts @@ -15,6 +15,7 @@ export type RetryContext = { empty?: boolean attempt?: number postToolContinuation?: boolean + subagent?: boolean } export type Retryable = { @@ -35,6 +36,7 @@ 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) @@ -96,9 +98,10 @@ export function retryable(error: Err, provider: string, context?: RetryContext) ) { if (context?.aborted) return undefined if (context?.empty === false) return undefined - if (MessageV2.PostToolContinuationTimeoutError.isInstance(error) && context?.postToolContinuation !== true) - return undefined - if ((context?.attempt ?? 1) > RETRY_NO_RESPONSE_LIMIT) 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 (MessageV2.APIError.isInstance(error)) { diff --git a/packages/opencode/test/session/processor-effect.test.ts b/packages/opencode/test/session/processor-effect.test.ts index 5d386cd23cd9..b032e84f92a4 100644 --- a/packages/opencode/test/session/processor-effect.test.ts +++ b/packages/opencode/test/session/processor-effect.test.ts @@ -792,10 +792,10 @@ postToolTimeoutRetryIt.live("session.processor effect tests retry stalled post-t ), ) -const postToolTimeoutExhaustLLM = llmMock(delayedFirstEventStream(50), delayedFirstEventStream(50)) +const postToolTimeoutExhaustLLM = llmMock(delayedFirstEventStream(50)) const postToolTimeoutExhaustIt = testEffect(processorEnv(postToolTimeoutExhaustLLM.layer)) -postToolTimeoutExhaustIt.live("session.processor effect tests stop with typed timeout after post-tool retry exhaustion", () => +postToolTimeoutExhaustIt.live("session.processor effect tests include retry attempt in post-tool timeout errors", () => provideTmpdirInstance( (dir) => Effect.gen(function* () { @@ -805,6 +805,13 @@ postToolTimeoutExhaustIt.live("session.processor effect tests stop with typed ti 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, @@ -833,17 +840,21 @@ postToolTimeoutExhaustIt.live("session.processor effect tests stop with typed ti }) expect(value).toBe("stop") - expect(postToolTimeoutExhaustLLM.calls()).toBe(2) - expect(handle.message.error?.name).toBe("PostToolContinuationTimeoutError") + 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, Stream.empty) +const emptyAssistantLLM = llmMock(Stream.empty, assistantSuccessStream("after-empty")) const emptyAssistantIt = testEffect(processorEnv(emptyAssistantLLM.layer)) -emptyAssistantIt.live("session.processor effect tests classify zero-part assistant finalization", () => +emptyAssistantIt.live("session.processor effect tests retry zero-part assistant finalization", () => provideTmpdirInstance( (dir) => Effect.gen(function* () { @@ -876,8 +887,56 @@ emptyAssistantIt.live("session.processor effect tests classify zero-part assista tools: {}, }) - expect(value).toBe("stop") + const parts = 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(MessageV2.parts(msg.id)).toHaveLength(0) expect(handle.message.error?.name).toBe("EmptyAssistantResponseError") }), @@ -885,10 +944,10 @@ emptyAssistantIt.live("session.processor effect tests classify zero-part assista ), ) -const structuralEmptyAssistantLLM = llmMock(structuralEmptyAssistantStream(), structuralEmptyAssistantStream()) +const structuralEmptyAssistantLLM = llmMock(structuralEmptyAssistantStream(), assistantSuccessStream("after-structural")) const structuralEmptyAssistantIt = testEffect(processorEnv(structuralEmptyAssistantLLM.layer)) -structuralEmptyAssistantIt.live("session.processor effect tests classify structural-only assistant finalization", () => +structuralEmptyAssistantIt.live("session.processor effect tests retry structural-only assistant finalization", () => provideTmpdirInstance( (dir) => Effect.gen(function* () { @@ -923,10 +982,59 @@ structuralEmptyAssistantIt.live("session.processor effect tests classify structu const parts = MessageV2.parts(msg.id) - expect(value).toBe("stop") + expect(value).toBe("continue") expect(structuralEmptyAssistantLLM.calls()).toBe(2) - expect(parts.every((part) => part.type === "step-start" || part.type === "step-finish")).toBe(true) - expect(handle.message.error?.name).toBe("EmptyAssistantResponseError") + 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 }, ), diff --git a/packages/opencode/test/session/retry.test.ts b/packages/opencode/test/session/retry.test.ts index 6b3082cf1f7f..9f5ec1a08b83 100644 --- a/packages/opencode/test/session/retry.test.ts +++ b/packages/opencode/test/session/retry.test.ts @@ -240,20 +240,90 @@ describe("session.retry.retryable", () => { expect(SessionRetry.retryable(error, retryProvider, { aborted: true, empty: true })).toBeUndefined() }) - test("retries post-tool continuation timeout once", () => { - const error = new MessageV2.PostToolContinuationTimeoutError({ - message: "No stream event", - abortSource: "post_tool_first_event_timeout", - phase: "post_tool_continuation", - retryable: true, - }).toObject() + 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() + } + }) - expect( - SessionRetry.retryable(error, retryProvider, { aborted: false, empty: true, postToolContinuation: true, attempt: 1 }), - ).toEqual({ message: "No stream event" }) - expect( - SessionRetry.retryable(error, retryProvider, { aborted: false, empty: true, postToolContinuation: true, attempt: 2 }), - ).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", () => { From 4a5b031f7dfc3931b3786df456b6da2a711b1393 Mon Sep 17 00:00:00 2001 From: soshymking Date: Wed, 10 Jun 2026 18:27:05 +0900 Subject: [PATCH 5/6] fix merge conflict --- packages/core/src/v1/session.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/core/src/v1/session.ts b/packages/core/src/v1/session.ts index edfd2a598406..a8934f7282da 100644 --- a/packages/core/src/v1/session.ts +++ b/packages/core/src/v1/session.ts @@ -59,7 +59,7 @@ export type RequestPhase = Schema.Schema.Type export const NoResponseDiagnostics = Schema.Struct({ providerID: Schema.optional(ProviderV2.ID), - modelID: Schema.optional(ProviderV2.ModelID), + modelID: Schema.optional(ModelV2.ID), sessionID: Schema.optional(SessionSchema.ID), messageID: Schema.optional(MessageID), elapsedMs: Schema.optional(NonNegativeInt), From 6f8213d941d7cede950fe985d3010bf84a343db2 Mon Sep 17 00:00:00 2001 From: soshymking Date: Wed, 10 Jun 2026 18:45:23 +0900 Subject: [PATCH 6/6] fix merge conflict --- packages/opencode/src/session/processor.ts | 18 ++++++++++-------- 1 file changed, 10 insertions(+), 8 deletions(-) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index aad1a7259296..70ccfd7c8e15 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -696,7 +696,7 @@ export const layer = Layer.effect( sessionID: ctx.sessionID, assistantMessageID, callID: value.id, - structured: toolOutput.structured, + structured: isRecord(toolOutput.structured) ? toolOutput.structured : { value: toolOutput.structured }, content: toolOutput.content, result: value.result, provider: { @@ -1105,13 +1105,15 @@ export const layer = Layer.effect( abortSource: "post_tool_first_event_timeout", phase: "post_tool_continuation", }) - return yield* 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(), - }) + 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),