diff --git a/packages/agent/CHANGELOG.md b/packages/agent/CHANGELOG.md index 8e3fb6a28..c42866b0f 100644 --- a/packages/agent/CHANGELOG.md +++ b/packages/agent/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- The error→toolUse salvage in the agent loop (`recoverTransientErrorToolTurn`) now recognizes Anthropic stream-envelope truncation errors, so a turn cut after streaming complete tool calls runs those calls instead of ending the run with an error. + ## [17.2.4] - 2026-08-01 ### Fixed diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index 409b76445..d92526325 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -1928,8 +1928,10 @@ function recoverTransientErrorToolTurn( if (tool.customWireName !== undefined) availableToolNames.add(tool.customWireName); } if (!toolCalls.every(toolCall => availableToolNames.has(toolCall.name))) return message; + const errorText = `${message.errorMessage ?? ""}\n${message.stopDetails?.explanation ?? ""}`; if ( - !AIError.isStreamReadErrorText(`${message.errorMessage ?? ""}\n${message.stopDetails?.explanation ?? ""}`) && + !AIError.isStreamReadErrorText(errorText) && + !AIError.isStreamEnvelopeErrorText(errorText) && !AIError.isTransientStreamParseError(message.errorMessage) && !AIError.isTransientStreamParseError(message.stopDetails?.explanation) ) diff --git a/packages/agent/test/agent-loop.test.ts b/packages/agent/test/agent-loop.test.ts index c6e574638..356d465cd 100644 --- a/packages/agent/test/agent-loop.test.ts +++ b/packages/agent/test/agent-loop.test.ts @@ -631,6 +631,51 @@ describe("agentLoop with AgentMessage", () => { expect(finalTurn.content).toContainEqual({ type: "text", text: "done after recovery" }); }); + it("runs completed tool calls after an Anthropic stream envelope truncation 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: "Anthropic stream envelope error: stream ended before message_stop", + }, + { 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"); + }); + it("runs completed tool calls after a transient stream JSON parse error", async () => { const executedParams: Array<{ value: string }> = []; const toolSchema = type({ value: "string" }); diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 7bcc538a7..5a5b88da9 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed Anthropic streams truncated mid-generation (connection closed with neither a `message_delta` stop_reason nor a `message_stop` frame) finalizing the partial message as a clean `stop`, which made the agent loop treat a truncated turn as complete and halt silently mid-sentence. Such streams raise the stream-envelope error again: transparently retried before replay-unsafe content streams; afterwards the turn surfaces as an error whose complete tool calls the agent loop still runs (`recoverTransientErrorToolTurn` now recognizes the envelope-error text after `retainCompletedToolCalls` drops half-streamed calls). Streams that delivered a `stop_reason` (or `message_stop`) keep degrading to best-effort content when the other terminal frame is missing. + ## [17.2.4] - 2026-08-01 ### Fixed diff --git a/packages/ai/src/error/flags.ts b/packages/ai/src/error/flags.ts index 01e4d71f9..09e54f17a 100644 --- a/packages/ai/src/error/flags.ts +++ b/packages/ai/src/error/flags.ts @@ -270,6 +270,14 @@ export function isStreamReadErrorText(text: string): boolean { return STREAM_READ_ERROR_PATTERN.test(text); } +/** Persisted-text form of {@link isStreamEnvelopeError}: recognizes the + * prefix-tagged envelope diagnostic on an aborted turn's `errorMessage` / + * `stopDetails.explanation` so loop-level salvage can classify it after the + * original `Error` instance is gone. */ +export function isStreamEnvelopeErrorText(text: string): boolean { + return text.includes(STREAM_ENVELOPE_ERROR_PREFIX); +} + function isTransientErrorText(text: string): boolean { return ( isUnexpectedSocketCloseMessage(text) || diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index d911a907e..4272f4224 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -2563,7 +2563,22 @@ const streamAnthropicOnce = ( if (!sawEvent || !sawMessageStart) { throw new AIError.AnthropicStreamEnvelopeError("stream ended before message_start"); } + if (!sawTerminalEnvelope) { + // Neither a message_delta stop_reason nor message_stop arrived: the + // connection died mid-generation. Finalizing the partial message as + // a clean "stop" would make the agent loop treat the truncated turn + // as complete (silent mid-sentence halt), so fail the turn. The + // envelope error is transparently retried before replay-unsafe + // content streams; afterwards it surfaces as an error turn whose + // complete tool calls the agent loop salvages + // (`recoverTransientErrorToolTurn` recognizes the envelope-error + // text and `retainCompletedToolCalls` drops half-streamed calls). + throw new AIError.AnthropicStreamEnvelopeError("stream ended before message_stop"); + } if (!sawMessageStop) { + // A stop_reason arrived via message_delta, so generation finished; + // only the trailing message_stop frame is missing (non-conforming + // gateway). Degrade to best-effort instead of discarding the turn. reportAnthropicEnvelopeAnomaly("stream ended before message_stop"); } if (openBlocks.size > 0) { diff --git a/packages/ai/test/anthropic-stream-envelope.test.ts b/packages/ai/test/anthropic-stream-envelope.test.ts index 1ac94995a..255d5a886 100644 --- a/packages/ai/test/anthropic-stream-envelope.test.ts +++ b/packages/ai/test/anthropic-stream-envelope.test.ts @@ -260,6 +260,12 @@ function createMalformedToolUseEvents(): MockAnthropicEvent[] { delta: { type: "input_json_delta", partial_json: '{"city":"Par' }, }, { type: "content_block_stop", index: 0 }, + { + type: "message_delta", + delta: { stop_reason: "end_turn" }, + usage: { input_tokens: 12, output_tokens: 4, cache_read_input_tokens: 0, cache_creation_input_tokens: 0 }, + }, + { type: "message_stop" }, ]; } @@ -288,6 +294,12 @@ function createGenuinelyMalformedToolUseEvents(): MockAnthropicEvent[] { delta: { type: "input_json_delta", partial_json: '{"city": Par' }, }, { type: "content_block_stop", index: 0 }, + { + type: "message_delta", + delta: { stop_reason: "end_turn" }, + usage: { input_tokens: 12, output_tokens: 4, cache_read_input_tokens: 0, cache_creation_input_tokens: 0 }, + }, + { type: "message_stop" }, ]; } @@ -1441,6 +1453,138 @@ describe("anthropic stream envelope handling", () => { expect(JSON.parse(JSON.stringify(result.content))).toEqual([{ type: "text", text: "partial" }]); }); + it("retries a stream cut before any stop_reason when no content streamed", async () => { + const successEvents = createTextSuccessEvents("recovered"); + let attempt = 0; + vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation(() => { + attempt += 1; + if (attempt === 1) { + // Connection died right after the envelope opened: only message_start + // arrived — no blocks, no message_delta stop_reason, no message_stop. + // (Any content_block_start already marks the stream replay-unsafe, + // which forbids the transparent retry.) + return createMockRequest([successEvents[0]]) as never; + } + return createMockRequest(createTextSuccessEvents("recovered")) as never; + }); + + const stream = streamAnthropic(model, context, { + apiKey: "sk-ant-test", + providerRetryWait: async () => {}, + }); + const events: AssistantMessageEvent[] = []; + for await (const event of stream) { + events.push(event); + } + const result = await stream.result(); + + expect(attempt).toBe(2); + expect(countEvents(events, "error")).toBe(0); + expect(countEvents(events, "done")).toBe(1); + expect(result.stopReason).toBe("stop"); + expect(JSON.parse(JSON.stringify(result.content))).toEqual([{ type: "text", text: "recovered" }]); + }); + + it("fails the turn instead of finalizing a clean stop when the stream dies mid-generation", async () => { + // message_start + content_block_start + text_delta, then the wire goes dead: + // no content_block_stop, no message_delta stop_reason, no message_stop. + const truncatedEvents = createTextSuccessEvents("partial").slice(0, 3); + let attempt = 0; + vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation(() => { + attempt += 1; + return createMockRequest(truncatedEvents) as never; + }); + + const stream = streamAnthropic(model, context, { + apiKey: "sk-ant-test", + providerRetryWait: async () => {}, + }); + const events: AssistantMessageEvent[] = []; + for await (const event of stream) { + events.push(event); + } + const result = await stream.result(); + + // Streamed text is replay-unsafe, so no transparent retry — the turn must + // surface as an error the agent loop can act on, never a silent "stop". + expect(attempt).toBe(1); + expect(countEvents(events, "done")).toBe(0); + expect(countEvents(events, "error")).toBe(1); + expect(result.stopReason).toBe("error"); + expect(result.errorMessage).toContain("message_stop"); + }); + + it("fails a truncated turn with mixed closed and half-streamed tool calls, keeping the closed call", async () => { + // One tool call closes cleanly, a second is cut mid-arguments, and the + // wire dies with no message_delta stop_reason and no message_stop. The + // turn must error (never finalize the half-streamed sibling into an + // executable call); the closed call stays in content so the agent loop's + // `retainCompletedToolCalls` + envelope-aware salvage can run it. + const truncatedEvents: MockAnthropicEvent[] = [ + { + type: "message_start", + message: { + id: "msg_mixed_truncated", + usage: { + input_tokens: 12, + output_tokens: 0, + cache_read_input_tokens: 0, + cache_creation_input_tokens: 0, + }, + }, + }, + { + type: "content_block_start", + index: 0, + content_block: { type: "tool_use", id: "tool_closed", name: "lookup_weather", input: {} }, + }, + { + type: "content_block_delta", + index: 0, + delta: { type: "input_json_delta", partial_json: '{"city":"Paris"}' }, + }, + { type: "content_block_stop", index: 0 }, + { + type: "content_block_start", + index: 1, + content_block: { type: "tool_use", id: "tool_open", name: "lookup_weather", input: {} }, + }, + { + type: "content_block_delta", + index: 1, + delta: { type: "input_json_delta", partial_json: '{"city":"Ber' }, + }, + ]; + let attempt = 0; + vi.spyOn(AnthropicMessages.prototype, "create").mockImplementation(() => { + attempt += 1; + return createMockRequest(truncatedEvents) as never; + }); + + const stream = streamAnthropic(model, context, { + apiKey: "sk-ant-test", + providerRetryWait: async () => {}, + }); + const events: AssistantMessageEvent[] = []; + for await (const event of stream) { + events.push(event); + } + const result = await stream.result(); + + expect(attempt).toBe(1); + expect(countEvents(events, "done")).toBe(0); + expect(countEvents(events, "error")).toBe(1); + expect(result.stopReason).toBe("error"); + expect(result.errorMessage).toContain("message_stop"); + const closedCall = result.content.find(block => block.type === "toolCall" && block.id === "tool_closed"); + expect(closedCall).toBeDefined(); + // The closed call emitted toolcall_end (loop-side completion marker); + // the half-streamed sibling did not. + const endedIds = events.filter(e => e.type === "toolcall_end").map(e => e.toolCall.id); + expect(endedIds).toContain("tool_closed"); + expect(endedIds).not.toContain("tool_open"); + }); + it("skips malformed raw SSE event frames and degrades to best-effort content", async () => { const malformedTextDelta = '{"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"line\\qbreak"}}';