diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index bf7f3af47..45a58da7e 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -10,6 +10,7 @@ ### Fixed - Fixed OpenAI Codex/Responses reasoning streams so streamed thinking content is preserved when the final `output_item.done` reconstructs to an empty summary ([#4918](https://github.com/can1357/oh-my-pi/issues/4918)). +- Fixed Anthropic streams hanging forever when generation wedges mid-stream (notably long `write` tool calls on Opus 4.8 high/xhigh) while the server keeps sending `ping` keepalives: pings now extend the idle watchdog only within a bounded window (3x the idle timeout) since the last real stream event, so a stalled tool-call stream times out and recovers instead of hanging with no retry path ([#4900](https://github.com/can1357/oh-my-pi/issues/4900)). ## [16.3.12] - 2026-07-08 diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index 6b3570c23..a22302fd2 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -1462,6 +1462,16 @@ async function* observeDecodedAnthropicSdkEvents( const PROVIDER_MAX_RETRIES = 10; +/** + * How long `ping` keepalives may keep extending the idle deadline without any + * semantic stream progress, as a multiple of the idle timeout. Anthropic pings + * across legitimate generation gaps, so pings count as liveness — but a wedged + * upstream that pings forever while producing no events must eventually trip + * the idle watchdog instead of hanging an active tool-call stream without a + * recovery path (#4900). + */ +const PING_PROGRESS_MAX_IDLE_MULTIPLIER = 3; + /** * Log a malformed-stream-envelope anomaly without aborting the turn. The strict * parser would `throw new AnthropicStreamEnvelopeError(...)` here; we instead @@ -2007,11 +2017,20 @@ const streamAnthropicOnce = ( } >(); - // Pings keep the idle deadline alive once content is flowing, but a - // ping before message_start must not consume the first-event watchdog: - // it would flip the (retryable) pre-content stall classification into - // a terminal mid-stream idle timeout. + // Pings keep the idle deadline alive once content is flowing (Anthropic + // bridges legitimate generation gaps with keepalives), but only within a + // bounded window: a wedged upstream that pings forever while the model + // produces nothing must still trip the idle watchdog, otherwise an + // active tool-call stream hangs unrecoverably with no retry (#4900). + // A ping before message_start must not consume the first-event watchdog + // either: it would flip the (retryable) pre-content stall classification + // into a terminal mid-stream idle timeout. let sawNonPingEvent = false; + let lastNonPingProgressAtMs = 0; + const pingProgressCapMs = + idleTimeoutMs !== undefined && idleTimeoutMs > 0 + ? idleTimeoutMs * PING_PROGRESS_MAX_IDLE_MULTIPLIER + : undefined; const timedAnthropicStream = iterateWithIdleTimeout(anthropicStream, { idleTimeoutMs, firstItemTimeoutMs: firstEventTimeoutMs, @@ -2021,8 +2040,13 @@ const streamAnthropicOnce = ( onFirstItemTimeout: () => activeAbortTracker.abortLocally(firstEventTimeoutAbortError), abortSignal: options?.signal, isProgressItem: item => { - if ((item as AnthropicStreamEvent).type === "ping") return sawNonPingEvent; + if ((item as AnthropicStreamEvent).type === "ping") { + if (!sawNonPingEvent) return false; + if (pingProgressCapMs === undefined) return true; + return Date.now() - lastNonPingProgressAtMs < pingProgressCapMs; + } sawNonPingEvent = true; + lastNonPingProgressAtMs = Date.now(); return true; }, }); diff --git a/packages/ai/test/anthropic-ping-keepalive.test.ts b/packages/ai/test/anthropic-ping-keepalive.test.ts new file mode 100644 index 000000000..a7e309b9e --- /dev/null +++ b/packages/ai/test/anthropic-ping-keepalive.test.ts @@ -0,0 +1,233 @@ +import { afterEach, describe, expect, it, vi } from "bun:test"; +import { buildModel } from "@oh-my-pi/pi-catalog/build"; +import { streamAnthropic } from "../src/providers/anthropic"; +import type { AnthropicMessagesClientLike } from "../src/providers/anthropic-client"; +import type { Context, Model } from "../src/types"; +import { waitForDelayOrAbort } from "./helpers"; + +const model: Model<"anthropic-messages"> = buildModel({ + id: "claude-opus-4-8", + name: "Claude Opus 4.8", + api: "anthropic-messages", + provider: "anthropic", + baseUrl: "https://api.anthropic.com", + reasoning: true, + input: ["text"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 200_000, + maxTokens: 8_192, +}); + +const context: Context = { + messages: [{ role: "user", content: "write a file", timestamp: Date.now() }], +}; + +type MockAnthropicEvent = Record; + +/** `{ waitMs, event }` script step; `waitMs` elapses (fake clock) before the event is yielded. */ +type ScriptStep = { waitMs: number; event: MockAnthropicEvent | "hang-with-pings" }; + +const writeToolCallOpening: MockAnthropicEvent[] = [ + { + type: "message_start", + message: { + id: "msg_ping_keepalive", + usage: { input_tokens: 10, 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: "toolu_ping_keepalive", name: "write", input: {} }, + }, + { + type: "content_block_delta", + index: 0, + delta: { type: "input_json_delta", partial_json: '{"path":"notes.md",' }, + }, +]; + +const writeToolCallClosing: MockAnthropicEvent[] = [ + { + type: "content_block_delta", + index: 0, + delta: { type: "input_json_delta", partial_json: '"content":"hello world"}' }, + }, + { type: "content_block_stop", index: 0 }, + { + type: "message_delta", + delta: { stop_reason: "tool_use" }, + usage: { input_tokens: 10, output_tokens: 6, cache_read_input_tokens: 0, cache_creation_input_tokens: 0 }, + }, + { type: "message_stop" }, +]; + +function createScriptedClient( + script: ScriptStep[], + counters: { pings: number }, + onIteratorStart: () => void, +): AnthropicMessagesClientLike { + const create = ((_body: unknown, requestOptions?: { signal?: AbortSignal }) => { + const signal = requestOptions?.signal; + const response = new Response(null, { status: 200, headers: { "request-id": "req_ping_keepalive" } }); + const stream = { + async *[Symbol.asyncIterator]() { + onIteratorStart(); + for (const step of script) { + if (step.event === "hang-with-pings") { + // Wedged upstream: no semantic events ever again, but the edge + // keeps the SSE connection alive with keepalive pings. + while (true) { + await waitForDelayOrAbort(step.waitMs, signal); + counters.pings += 1; + yield { type: "ping" }; + } + } + if (step.waitMs > 0) { + await waitForDelayOrAbort(step.waitMs, signal); + } + if (step.event.type === "ping") counters.pings += 1; + yield step.event; + } + }, + }; + return { + async withResponse() { + return { data: stream, response, request_id: "req_ping_keepalive" }; + }, + } as never; + }) as unknown as AnthropicMessagesClientLike["messages"]["create"]; + return { messages: { create } } as AnthropicMessagesClientLike; +} + +async function drainMicrotasks(count: number): Promise { + for (let i = 0; i < count; i++) { + await Promise.resolve(); + } +} + +async function drainMicrotasksUntil(predicate: () => boolean, errorMessage: string): Promise { + for (let i = 0; i < 1000; i++) { + if (predicate()) return; + await Promise.resolve(); + } + throw new Error(errorMessage); +} + +afterEach(() => { + vi.useRealTimers(); + vi.restoreAllMocks(); +}); + +describe("anthropic ping keepalive idle cap", () => { + it("times out a stalled tool-call stream instead of letting pings extend it forever", async () => { + vi.useFakeTimers(); + const counters = { pings: 0 }; + let iteratorStarted = false; + const script: ScriptStep[] = [ + ...writeToolCallOpening.map(event => ({ waitMs: 0, event })), + { waitMs: 500, event: "hang-with-pings" as const }, + ]; + const client = createScriptedClient(script, counters, () => { + iteratorStarted = true; + }); + const providerRetryWait = vi.fn(async () => {}); + + let settled = false; + const resultPromise = streamAnthropic(model, context, { + client, + streamFirstEventTimeoutMs: 1_000, + streamIdleTimeoutMs: 1_000, + providerRetryWait, + }) + .result() + .then(message => { + settled = true; + return message; + }); + + await drainMicrotasksUntil(() => iteratorStarted, "Anthropic mock stream never started"); + await drainMicrotasks(30); + + // Pings arrive every 500 fake-ms while generation is wedged. Drive far + // past the bounded keepalive window (3x idle = 3_000ms) plus one idle + // budget; without the cap the idle deadline is reset by every ping and + // this loop ends with the result still pending (issue #4900's hang). + let stepsRun = 0; + for (let step = 0; step < 40 && !settled; step++) { + vi.advanceTimersByTime(500); + await drainMicrotasks(30); + stepsRun = step + 1; + } + + expect(settled).toBe(true); + // Cap (3_000ms) + idle budget (1_000ms) = fires at 3_500-4_000 fake ms. + expect(stepsRun).toBeLessThanOrEqual(9); + // Keepalives within the window were honored before the watchdog fired. + expect(counters.pings).toBeGreaterThanOrEqual(5); + + const result = await resultPromise; + expect(result.stopReason).toBe("error"); + expect(result.errorMessage).toBe("Anthropic stream stalled while waiting for the next event"); + // Mid-stream idle stalls are terminal for the provider loop (session-level + // auto-retry owns recovery); the provider must not silently re-request. + expect(providerRetryWait).not.toHaveBeenCalled(); + }); + + it("keeps a slow-but-alive stream open across ping-bridged gaps within the cap", async () => { + vi.useFakeTimers(); + const counters = { pings: 0 }; + let iteratorStarted = false; + // Silent generation gap of 1_800ms (> 1_000ms idle budget) bridged by + // pings at t=600 and t=1200, then semantic progress resumes and the + // tool call completes. Pings within the cap must count as liveness. + const script: ScriptStep[] = [ + ...writeToolCallOpening.map(event => ({ waitMs: 0, event })), + { waitMs: 600, event: { type: "ping" } }, + { waitMs: 600, event: { type: "ping" } }, + { waitMs: 600, event: writeToolCallClosing[0]! }, + ...writeToolCallClosing.slice(1).map(event => ({ waitMs: 0, event })), + ]; + const client = createScriptedClient(script, counters, () => { + iteratorStarted = true; + }); + const providerRetryWait = vi.fn(async () => {}); + + let settled = false; + const resultPromise = streamAnthropic(model, context, { + client, + streamFirstEventTimeoutMs: 1_000, + streamIdleTimeoutMs: 1_000, + providerRetryWait, + }) + .result() + .then(message => { + settled = true; + return message; + }); + + await drainMicrotasksUntil(() => iteratorStarted, "Anthropic mock stream never started"); + await drainMicrotasks(30); + + for (let step = 0; step < 30 && !settled; step++) { + vi.advanceTimersByTime(200); + await drainMicrotasks(30); + } + + expect(settled).toBe(true); + expect(counters.pings).toBe(2); + + const result = await resultPromise; + expect(result.errorMessage).toBeUndefined(); + expect(result.stopReason).toBe("toolUse"); + expect(providerRetryWait).not.toHaveBeenCalled(); + expect(JSON.parse(JSON.stringify(result.content))).toEqual([ + { + type: "toolCall", + id: "toolu_ping_keepalive", + name: "write", + arguments: { path: "notes.md", content: "hello world" }, + }, + ]); + }); +});