fix(ai): honor model.compat.streamIdleTimeoutMs in the Anthropic idle watchdog
anthropic-messages requests dispatched through the direct Anthropic provider path ignored the catalog compat idle timeout. Slow relayed endpoints therefore remained limited by the default inter-event watchdog even when the model configured a wider or disabled timeout. Add streamIdleTimeoutMs to AnthropicCompat, preserve it through buildAnthropicCompat, and pass it as the fallback to getStreamIdleTimeoutMs in streamAnthropic. A value of 0 disables only the inter-event watchdog; the first-event watchdog is unchanged. Resolution order matches the existing timeout helper contract: caller option > PI_STREAM_IDLE_TIMEOUT_MS > PI_OPENAI_STREAM_IDLE_TIMEOUT_MS > model compat > 300s default The pi-native transport retains its existing transport-level timeout policy. Tested: anthropic stream timeout suite, env-contaminated regression cases, models config validation, and bun check.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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<T>(promise: Promise<T>, 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<StreamTimeoutEnvKey, string | undefined> = {
|
||||
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<typeof buildModel>[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;
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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<AnthropicCompat> & {
|
||||
export type ResolvedAnthropicCompat = Required<Omit<AnthropicCompat, "streamIdleTimeoutMs">> & {
|
||||
/**
|
||||
* 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
|
||||
|
||||
Reference in New Issue
Block a user