fix(ai): bounded anthropic ping keepalive extension of idle watchdog
A wedged Anthropic stream that kept emitting SSE ping keepalives reset the idle deadline on every ping, so an active tool-call stream (notably long write calls on Opus 4.8 high/xhigh) could hang forever with no timeout, no error, and no retry path. Pings now count as liveness only within a bounded window (3x the idle timeout) since the last semantic stream event; past it the idle watchdog fires and the turn surfaces a terminal stream-stall error that session-level auto-retry can recover. Pings before message_start still never consume the first-event watchdog, and pings bridging legitimate generation gaps within the window still keep slow streams alive. Fixes #4900
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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;
|
||||
},
|
||||
});
|
||||
|
||||
@@ -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<string, unknown>;
|
||||
|
||||
/** `{ 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<void> {
|
||||
for (let i = 0; i < count; i++) {
|
||||
await Promise.resolve();
|
||||
}
|
||||
}
|
||||
|
||||
async function drainMicrotasksUntil(predicate: () => boolean, errorMessage: string): Promise<void> {
|
||||
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" },
|
||||
},
|
||||
]);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user