fix(ai): addressed anthropic stream truncation errors for tool call recovery
- Added `isStreamEnvelopeErrorText` to packages/ai/src/error/flags.ts to recognize stream envelope truncation errors. - Updated `streamAnthropicOnce` in packages/ai/src/providers/anthropic.ts to throw an envelope error when streams die mid-generation without a terminal frame. - Updated `recoverTransientErrorToolTurn` in packages/agent/src/agent-loop.ts to recognize Anthropic stream envelope truncation errors for tool call salvage.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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)
|
||||
)
|
||||
|
||||
@@ -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<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: "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" });
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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) ||
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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"}}';
|
||||
|
||||
Reference in New Issue
Block a user