diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index ff016809e..26e4b4c4c 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed direct Anthropic provider streams ignoring `model.compat.streamIdleTimeoutMs`. Requests dispatched through `streamAnthropic` can now widen the inter-event idle watchdog or set it to `0` to disable that watchdog; caller options and environment overrides retain precedence. Setting the compat value to `0` disables only the inter-event watchdog and leaves the first-event watchdog enabled; wider idle values continue to floor the first-event budget under the existing timeout contract. + ## [17.2.3] - 2026-08-01 ### Added diff --git a/packages/ai/src/providers/anthropic.ts b/packages/ai/src/providers/anthropic.ts index ca953db9c..d911a907e 100644 --- a/packages/ai/src/providers/anthropic.ts +++ b/packages/ai/src/providers/anthropic.ts @@ -2002,7 +2002,7 @@ const streamAnthropicOnce = ( | (AnthropicServerToolContent & { [kStreamingPartialJson]?: string }) | (ToolCall & { [kStreamingPartialJson]: string; [kStreamingLastParseLen]?: number }) ) & { [kStreamingBlockIndex]: number }; - const idleTimeoutMs = options?.streamIdleTimeoutMs ?? getStreamIdleTimeoutMs(); + const idleTimeoutMs = options?.streamIdleTimeoutMs ?? getStreamIdleTimeoutMs(model.compat.streamIdleTimeoutMs); const firstEventTimeoutMs = options?.streamFirstEventTimeoutMs ?? getStreamFirstEventTimeoutMs(idleTimeoutMs); const requestTimeoutMs = firstEventTimeoutMs !== undefined && firstEventTimeoutMs > 0 ? firstEventTimeoutMs : undefined; diff --git a/packages/ai/test/anthropic-stream-timeout.test.ts b/packages/ai/test/anthropic-stream-timeout.test.ts index c4a2d1e69..c5341ae75 100644 --- a/packages/ai/test/anthropic-stream-timeout.test.ts +++ b/packages/ai/test/anthropic-stream-timeout.test.ts @@ -1,4 +1,4 @@ -import { afterEach, describe, expect, it, vi } from "bun:test"; +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; import * as AIError from "@oh-my-pi/pi-ai/error"; import { streamAnthropic } from "@oh-my-pi/pi-ai/providers/anthropic"; import { AnthropicMessagesClient, type AnthropicMessagesClientLike } from "@oh-my-pi/pi-ai/providers/anthropic-client"; @@ -204,7 +204,37 @@ async function resolveAfterMicrotasks(promise: Promise, errorMessage: stri return outcome.value; } +const STREAM_TIMEOUT_ENV_KEYS = [ + "PI_STREAM_IDLE_TIMEOUT_MS", + "PI_OPENAI_STREAM_IDLE_TIMEOUT_MS", + "PI_STREAM_FIRST_EVENT_TIMEOUT_MS", +] as const; + +type StreamTimeoutEnvKey = (typeof STREAM_TIMEOUT_ENV_KEYS)[number]; + +const originalStreamTimeoutEnv: Record = { + PI_STREAM_IDLE_TIMEOUT_MS: undefined, + PI_OPENAI_STREAM_IDLE_TIMEOUT_MS: undefined, + PI_STREAM_FIRST_EVENT_TIMEOUT_MS: undefined, +}; + +beforeEach(() => { + for (const key of STREAM_TIMEOUT_ENV_KEYS) { + originalStreamTimeoutEnv[key] = Bun.env[key]; + delete Bun.env[key]; + } +}); + afterEach(() => { + for (const key of STREAM_TIMEOUT_ENV_KEYS) { + const previous = originalStreamTimeoutEnv[key]; + if (previous === undefined) { + delete Bun.env[key]; + } else { + Bun.env[key] = previous; + } + } + vi.useRealTimers(); vi.restoreAllMocks(); }); @@ -460,6 +490,114 @@ describe("anthropic first-event timeout retries", () => { }); }); +describe("anthropic model compat stream idle timeout floor", () => { + const baseModel = { + id: "claude-sonnet-4-5", + name: "Claude Sonnet 4.5", + api: "anthropic-messages" as const, + 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, + } satisfies Parameters[0]; + + function createStalledAfterFirstEventClient(onIteratorStart?: () => void): { + attempt: () => number; + client: AnthropicMessagesClientLike; + } { + let attempt = 0; + const create = ((_body: unknown, requestOptions?: { signal?: AbortSignal }) => { + attempt += 1; + return createAnthropicMockStream({ + signal: requestOptions?.signal, + events: [ + { + type: "message_start", + message: { + id: "msg_compat_stall", + usage: { + input_tokens: 12, + output_tokens: 0, + cache_read_input_tokens: 0, + cache_creation_input_tokens: 0, + }, + }, + }, + ], + hangAfterEvents: true, + onIteratorStart, + }) as never; + }) as unknown as AnthropicMessagesClientLike["messages"]["create"]; + return { attempt: () => attempt, client: { messages: { create } } as AnthropicMessagesClientLike }; + } + + it("uses model.compat.streamIdleTimeoutMs as the idle floor when no caller option is set", async () => { + const compatModel = buildModel({ ...baseModel, compat: { streamIdleTimeoutMs: 50 } }); + const { attempt, client } = createStalledAfterFirstEventClient(); + + const result = await streamAnthropic(compatModel, context, { + client, + streamFirstEventTimeoutMs: 5_000, + }).result(); + + expect(attempt()).toBe(1); + expect(result.stopReason).toBe("error"); + expect(result.errorMessage).toBe("Anthropic stream stalled while waiting for the next event"); + }); + + it("disables the idle watchdog when model.compat.streamIdleTimeoutMs is 0", async () => { + vi.useFakeTimers(); + const compatModel = buildModel({ ...baseModel, compat: { streamIdleTimeoutMs: 0 } }); + const controller = new AbortController(); + let iteratorStarted = false; + const { attempt, client } = createStalledAfterFirstEventClient(() => { + iteratorStarted = true; + }); + + let settled = false; + const resultPromise = streamAnthropic(compatModel, context, { + client, + signal: controller.signal, + // First-event watchdog stays out of this case's scope: it would need to + // be cleared by the mock's first event, and fake-timer advancement can + // outrun the microtask that consumes that event under filtered runs. + streamFirstEventTimeoutMs: 0, + }).result(); + void resultPromise.then( + () => { + settled = true; + }, + () => { + settled = true; + }, + ); + + await drainMicrotasksUntil( + () => iteratorStarted, + "Anthropic mock stream did not start for the compat-disabled watchdog test", + ); + // Well past the default 300s idle floor: the disabled watchdog must not + // classify the post-first-event silence as a stall. + vi.advanceTimersByTime(400_000); + await drainMicrotasksUntil( + () => vi.getTimerCount() === 0, + "Anthropic watchdog timer did not drain after advancing past the idle budget", + ); + expect(settled).toBe(false); + expect(attempt()).toBe(1); + + controller.abort(); + const result = await resolveAfterMicrotasks( + resultPromise, + "Anthropic compat-disabled stream did not settle after the caller aborted", + ); + expect(result.stopReason).toBe("aborted"); + }); +}); + describe("anthropic provider retry delays", () => { it("waits at least the server-suggested retry-after before retrying a retryable API error", async () => { let attempt = 0; diff --git a/packages/catalog/CHANGELOG.md b/packages/catalog/CHANGELOG.md index c329a78e5..e87f7e92e 100644 --- a/packages/catalog/CHANGELOG.md +++ b/packages/catalog/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Added + +- Added `AnthropicCompat.streamIdleTimeoutMs` and propagated it through `buildAnthropicCompat` so direct Anthropic provider streams can configure their inter-event idle watchdog. + ## [17.2.3] - 2026-08-01 ### Added diff --git a/packages/catalog/src/compat/anthropic.ts b/packages/catalog/src/compat/anthropic.ts index f5e1be87f..786b729c6 100644 --- a/packages/catalog/src/compat/anthropic.ts +++ b/packages/catalog/src/compat/anthropic.ts @@ -172,6 +172,7 @@ export function buildAnthropicCompat(spec: ModelSpec<"anthropic-messages">): Res // id or baseUrl marker. replayUnsignedThinking: !signingEndpoint && (Boolean(spec.reasoning) || modelMatchesHost(spec, "deepseekFamily")), escapeBuiltinToolNames: modelMatchesHost(spec, "umans"), + streamIdleTimeoutMs: spec.compat?.streamIdleTimeoutMs, }; applyCompatOverrides(compat, spec.compat); return compat; diff --git a/packages/catalog/src/types.ts b/packages/catalog/src/types.ts index 137e1af3e..4e289759e 100644 --- a/packages/catalog/src/types.ts +++ b/packages/catalog/src/types.ts @@ -401,6 +401,15 @@ export interface OpenAICompat { * that proxy gateways (Vertex AI, AWS Bedrock-style fronts, etc.) reject. */ export interface AnthropicCompat { + /** + * Stream-watchdog idle-timeout fallback in ms for slow reasoning hosts. + * Set to 0 to disable the inter-event idle watchdog entirely, matching + * `OpenAICompat.streamIdleTimeoutMs`. + * + * When unset, direct Anthropic streams use `PI_STREAM_IDLE_TIMEOUT_MS`, + * then the legacy `PI_OPENAI_STREAM_IDLE_TIMEOUT_MS` alias, then 300s. + */ + streamIdleTimeoutMs?: number; /** * Drop the top-level `strict: true` field on tool definitions. Vertex AI's * Anthropic-compatible endpoint rejects unknown tool fields with @@ -712,7 +721,13 @@ export interface ResolvedOpenAIResponsesCompat extends ResolvedOpenAISharedCompa export type ResolvedOpenRouterCompat = ResolvedOpenAICompat & ResolvedOpenAIResponsesCompat; /** Fully-resolved anthropic-messages compat view (same contract as `ResolvedOpenAICompat`). */ -export type ResolvedAnthropicCompat = Required & { +export type ResolvedAnthropicCompat = Required> & { + /** + * Stream-watchdog idle-timeout fallback in ms for slow reasoning hosts; 0 disables the idle watchdog. + * Undefined defers to `PI_STREAM_IDLE_TIMEOUT_MS`, then the legacy + * `PI_OPENAI_STREAM_IDLE_TIMEOUT_MS` alias, then 300s. + */ + streamIdleTimeoutMs?: number; /** * The configured endpoint is the official first-party Anthropic API * (https + exact `api.anthropic.com` host; a missing baseUrl counts as