fix(agent): recover completed tools after stream read errors

This commit is contained in:
ben
2026-06-28 11:54:49 +08:00
parent 4373dde48a
commit cb4e68888f
7 changed files with 159 additions and 2 deletions
+3
View File
@@ -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
+32 -1
View File
@@ -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<Pick<AgentTool, "name">>,
): 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<AgentEvent, AgentMessage[]>,
+47
View File
@@ -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<typeof toolSchema, { value: string }> = {
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<string, unknown>[] = [];
+4
View File
@@ -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
+1 -1
View File
@@ -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;
+12
View File
@@ -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);
@@ -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) {