From cb4e68888f2c771a31f5df2eb93134bec22e4db3 Mon Sep 17 00:00:00 2001 From: ben Date: Sun, 28 Jun 2026 11:54:49 +0800 Subject: [PATCH] fix(agent): recover completed tools after stream read errors --- packages/agent/CHANGELOG.md | 3 + packages/agent/src/agent-loop.ts | 33 +++++++++- packages/agent/test/agent-loop.test.ts | 47 +++++++++++++++ packages/ai/CHANGELOG.md | 4 ++ packages/ai/src/error/flags.ts | 2 +- packages/ai/test/error-id.test.ts | 12 ++++ .../test/agent-session-retry-cap.test.ts | 60 +++++++++++++++++++ 7 files changed, 159 insertions(+), 2 deletions(-) diff --git a/packages/agent/CHANGELOG.md b/packages/agent/CHANGELOG.md index 546363a97..7b891b4d3 100644 --- a/packages/agent/CHANGELOG.md +++ b/packages/agent/CHANGELOG.md @@ -16,6 +16,9 @@ - Fixed an issue where assistant responses and encrypted reasoning were lost during local history trimming prior to remote compaction. - Improved reliability of remote compaction with transient error retries, configurable timeouts, and immediate termination upon user-initiated aborts. - Added title_change session metadata to the compaction entry type union to maintain type compatibility for hosts with title audit entries. +### Fixed + +- Fixed transient stream read failures after a completed tool call being treated as terminal errors; the agent now executes the completed tool call and continues the turn. ## [16.2.2] - 2026-06-27 diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index 889be34e7..fe92e5e2e 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -1355,7 +1355,10 @@ async function streamAssistantResponse( const event = next.value; if (event.type === "done" || event.type === "error") { - let finalMessage = retainCompletedToolCalls(await response.result(), completedToolCallIds); + let finalMessage = recoverTransientErrorToolTurn( + retainCompletedToolCalls(await response.result(), completedToolCallIds), + context.tools ?? [], + ); if (harmonyMitigationEnabled) { const detection = detectHarmonyLeakInAssistantMessage(finalMessage); if (detection) { @@ -1520,6 +1523,34 @@ function retainCompletedToolCalls( }; } +function recoverTransientErrorToolTurn( + message: AssistantMessage, + availableTools: ReadonlyArray>, +): AssistantMessage { + if (message.stopReason !== "error") return message; + const toolCalls = message.content.filter(block => block.type === "toolCall"); + if (toolCalls.length === 0) return message; + const availableToolNames = new Set(availableTools.map(tool => tool.name)); + if (!toolCalls.every(toolCall => availableToolNames.has(toolCall.name))) return message; + const errorId = AIError.classifyMessage(message); + if (!AIError.is(errorId, AIError.Flag.Transient)) return message; + return { + ...message, + stopReason: "toolUse", + stopDetails: + message.stopDetails?.type === STREAM_INTERRUPTED_AFTER_CONTENT_STOP_DETAIL + ? message.stopDetails + : { + type: STREAM_INTERRUPTED_AFTER_CONTENT_STOP_DETAIL, + category: message.stopDetails?.type ?? null, + explanation: message.stopDetails?.explanation ?? message.errorMessage ?? null, + }, + errorMessage: undefined, + errorId: undefined, + errorStatus: undefined, + }; +} + function emitDiscardedHarmonyPartial( partialMessage: AssistantMessage | null, stream: EventStream, diff --git a/packages/agent/test/agent-loop.test.ts b/packages/agent/test/agent-loop.test.ts index 1064134d6..511e5b628 100644 --- a/packages/agent/test/agent-loop.test.ts +++ b/packages/agent/test/agent-loop.test.ts @@ -565,6 +565,53 @@ describe("agentLoop with AgentMessage", () => { expect(toolStart.args.__parseError).toBeDefined(); // keeps __parseError for visibility of parse failure }); + it("runs completed tool calls after a transient stream_read_error", async () => { + const executedParams: Array<{ value: string }> = []; + const toolSchema = type({ value: "string" }); + const tool: AgentTool = { + name: "echo", + label: "Echo", + description: "Echo tool", + parameters: toolSchema, + async execute(_toolCallId, params) { + executedParams.push(params); + return { + content: [{ type: "text", text: `echoed: ${params.value}` }], + details: { value: params.value }, + }; + }, + }; + const context: AgentContext = { systemPrompt: [""], messages: [], tools: [tool] }; + const mock = createMockModel({ + responses: [ + { + content: [{ type: "toolCall", id: "tool-1", name: "echo", arguments: { value: "hello" } }], + stopReason: "error", + errorMessage: "Error Code stream_read_error: stream_read_error", + }, + { content: ["done after recovery"] }, + ], + }); + const config: AgentLoopConfig = { model: mock.model, convertToLlm: identityConverter }; + + const messages = await agentLoop( + [createUserMessage("run echo")], + context, + config, + undefined, + mock.stream, + ).result(); + + expect(executedParams).toEqual([{ value: "hello" }]); + expect(mock.calls).toHaveLength(2); + expect(messages.map(message => message.role)).toEqual(["user", "assistant", "toolResult", "assistant"]); + const recoveredTurn = messages[1] as AssistantMessage; + expect(recoveredTurn.stopReason).toBe("toolUse"); + expect(recoveredTurn.stopDetails?.type).toBe("stream_interrupted_after_content"); + const finalTurn = messages[3] as AssistantMessage; + expect(finalTurn.content).toContainEqual({ type: "text", text: "done after recovery" }); + }); + it("injects and strips intent when intent tracing is enabled", async () => { const toolSchema = type({ value: "string" }); const executedParams: Record[] = []; diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index acee75192..d4cfc4063 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -11,6 +11,10 @@ - Enabled freeform tool patch support for Azure OpenAI and Codex models - Fixed /usage show returning "No usage data available" when using a custom proxy base URL for Codex by routing usage and credit-reset requests to the canonical ChatGPT origin +### Fixed + +- Fixed OpenAI Responses `stream_read_error` provider events being classified as non-transient, which prevented the coding agent's auto-retry path from continuing after a recoverable stream read failure. + ## [16.2.2] - 2026-06-27 ### Added diff --git a/packages/ai/src/error/flags.ts b/packages/ai/src/error/flags.ts index 1044379bb..e13ffe688 100644 --- a/packages/ai/src/error/flags.ts +++ b/packages/ai/src/error/flags.ts @@ -86,7 +86,7 @@ const TIMEOUT_PATTERN = /\b(?:operation\s+)?timed?\s*out\b|\btimeout\b|\bstream const TRANSIENT_ENVELOPE_PATTERN = /anthropic stream envelope error:/i; const TRANSIENT_ENVELOPE_BEFORE_START_PATTERN = /before message_start/i; export const TRANSIENT_TRANSPORT_PATTERN = - /overloaded|provider.?returned.?error|rate.?limit|too many requests|429|500|502|503|504|service.?unavailable|server.?error|internal.?error|retry your request|network.?error|connection.?error|connection.?refused|other side closed|fetch failed|upstream.?connect|upstream.?request.?failed|reset before headers|socket hang up|timed? out|timeout|terminated|retry delay|stream stall|no error details in response|HTTP2(?:StreamReset|RefusedStream|EnhanceYourCalm)|malformed.?function.?call/i; + /overloaded|provider.?returned.?error|rate.?limit|too many requests|429|500|502|503|504|service.?unavailable|server.?error|internal.?error|retry your request|network.?error|connection.?error|connection.?refused|other side closed|fetch failed|upstream.?connect|upstream.?request.?failed|reset before headers|socket hang up|timed? out|timeout|terminated|retry delay|stream stall|stream[_ -]?read[_ -]?error|no error details in response|HTTP2(?:StreamReset|RefusedStream|EnhanceYourCalm)|malformed.?function.?call/i; const AUTH_FAILURE_PATTERN = /\b(?:401|403|unauthorized|forbidden|authentication|auth[_ ]?unavailable|no auth available|(?:invalid|no)[_ ]?api[_ ]?key)\b/i; const MALFORMED_FUNCTION_CALL_PATTERN = /\bmalformed.?function.?call\b/i; diff --git a/packages/ai/test/error-id.test.ts b/packages/ai/test/error-id.test.ts index e6afc1181..18b390baa 100644 --- a/packages/ai/test/error-id.test.ts +++ b/packages/ai/test/error-id.test.ts @@ -31,6 +31,18 @@ describe("error-id classification", () => { expect(AIError.is(id, AIError.Flag.Class)).toBe(true); }); + it("classifies OpenAI stream_read_error as transient", () => { + const assistant = message({ + api: "openai-responses", + provider: "openai", + model: "gpt-5", + errorMessage: "Error Code stream_read_error: stream_read_error", + }); + const id = AIError.classifyMessage(assistant); + expect(AIError.is(id, AIError.Flag.Transient)).toBe(true); + expect(AIError.retriable(id)).toBe(true); + }); + it("keeps raw status fallback unclassified", () => { const id = 503; expect(AIError.is(id, AIError.Flag.Class)).toBe(false); diff --git a/packages/coding-agent/test/agent-session-retry-cap.test.ts b/packages/coding-agent/test/agent-session-retry-cap.test.ts index 93cf9a3ce..60c54526a 100644 --- a/packages/coding-agent/test/agent-session-retry-cap.test.ts +++ b/packages/coding-agent/test/agent-session-retry-cap.test.ts @@ -141,6 +141,66 @@ describe("AgentSession retry delay cap", () => { expect(session.isRetrying).toBe(false); }); + it("auto-retries OpenAI Responses stream_read_error instead of stopping the conversation", async () => { + const model = getBundledModel("openai", "gpt-5"); + if (!model) { + throw new Error("Expected bundled OpenAI test model to exist"); + } + authStorage.setRuntimeApiKey("openai", "openai-test-key"); + + const mock = createMockModel({ + responses: [ + { throw: "Error Code stream_read_error: stream_read_error" }, + { content: ["recovered after stream read retry"], stopReason: "stop" }, + ], + }); + const agent = new Agent({ + getApiKey: requestedModel => `${requestedModel.provider}-test-key`, + initialState: { + model, + systemPrompt: ["Test"], + tools: [], + messages: [], + }, + streamFn: (requestedModel, context, options) => mock.stream(requestedModel, context, options), + }); + + const settings = Settings.isolated({ + "compaction.enabled": false, + "retry.baseDelayMs": 5, + "retry.maxDelayMs": 5_000, + "retry.maxRetries": 1, + "retry.modelFallback": false, + }); + settings.setModelRole("default", `${model.provider}/${model.id}`); + + session = new AgentSession({ + agent, + sessionManager: SessionManager.inMemory(), + settings, + modelRegistry, + }); + + vi.spyOn(scheduler, "wait").mockResolvedValue(undefined); + const retryStartEvents: AutoRetryStartEvent[] = []; + const retryEndEvents: AutoRetryEndEvent[] = []; + session.subscribe(event => { + if (event.type === "auto_retry_start") retryStartEvents.push(event); + if (event.type === "auto_retry_end") retryEndEvents.push(event); + }); + + await session.prompt("Trigger stream read retry"); + await session.waitForIdle(); + + expect(mock.calls).toHaveLength(2); + expect(retryStartEvents).toHaveLength(1); + expect(retryEndEvents).toHaveLength(1); + expect(retryEndEvents[0]).toMatchObject({ success: true }); + const last = lastAssistant(session); + expect(last.stopReason).toBe("stop"); + expect(last.content).toContainEqual({ type: "text", text: "recovered after stream read retry" }); + }); + it("switches credentials instead of failing the delay cap for account rate limits", async () => { const model = getBundledModel("anthropic", "claude-sonnet-4-5"); if (!model) {