fix(ai/providers): handled Anthropic ping keepalives and excluded usage limits from retrying

- Anthropic streaming now yielded explicit `ping` events and propagated them through the stream event types.
- Ping keepalive markers reset idle liveness handling so long-running streams no longer stalled as idle.
- Usage and quota limit errors were treated as non-retryable, and tests were added to keep transient rate-limit retries unchanged.
This commit is contained in:
can1357
2026-06-08 02:52:45 +02:00
parent 733d1dabb1
commit c4b3f267f5
4 changed files with 54 additions and 7 deletions
+1
View File
@@ -15,6 +15,7 @@
### Fixed
- Fixed duplicate upstream `tool_call_id` values collapsing distinct tool calls during message transformation, preserving one call/result pairing per emitted tool call before provider replay and keeping generated duplicate IDs distinct after OpenAI/Mistral wire-length caps. ([#2055](https://github.com/can1357/oh-my-pi/issues/2055))
- Fixed the Anthropic provider retrying persistent account usage/quota limits (e.g. `429 "This request would exceed your account's rate limit"`, `usage_limit_reached`) as if they were transient. Because the error text contains "rate limit", `isProviderRetryableError` matched it and the stream retry loop looped through its 2s/4s/8s backoff (then the `streamSimple` a/b/c policy re-minted the credential and ran the whole thing again) before surfacing the failure — even though the server's `retry-after` parked the account for minutes-to-hours. These errors are now recognized via `isUsageLimitError` and surfaced immediately to the credential-rotation layer, so e.g. `omp dry-balance --bench` reports a rate-limited account as failed at once instead of appearing to hang.
## [15.10.1] - 2026-06-07
+34 -5
View File
@@ -18,6 +18,7 @@ import {
supportsMidConversationSystemMessages,
} from "../model-thinking";
import { calculateCost } from "../models";
import { isUsageLimitError } from "../rate-limit-utils";
import { getEnvApiKey, OUTPUT_FALLBACK_BUFFER } from "../stream";
import type {
Api,
@@ -1036,11 +1037,25 @@ const ANTHROPIC_MESSAGE_EVENTS: ReadonlySet<string> = new Set([
"content_block_stop",
]);
/**
* Anthropic keepalive `ping` events carry no message content, but they prove the
* upstream connection is alive during long server-side gaps (extended thinking,
* slow tool execution). They are normally dropped before reaching the consumer;
* we instead surface them as lightweight markers so the idle watchdog
* (`iterateWithIdleTimeout`) resets its deadline on every ping. Without this, a
* connection that is demonstrably still streaming pings still trips
* "Anthropic stream stalled while waiting for the next event". The message-event
* branches in `streamAnthropic` match none of these markers, so they are ignored.
*/
type RawMessagePingEvent = { type: "ping" };
type AnthropicStreamEvent = RawMessageStreamEvent | RawMessagePingEvent;
const ANTHROPIC_PING_EVENT: RawMessagePingEvent = { type: "ping" };
async function* iterateAnthropicEvents(
response: Response,
signal?: AbortSignal,
onSseEvent?: AnthropicOptions["onSseEvent"],
): AsyncGenerator<RawMessageStreamEvent> {
): AsyncGenerator<AnthropicStreamEvent> {
if (!response.body) {
throw new Error("Attempted to iterate over an Anthropic response with no body");
}
@@ -1054,6 +1069,12 @@ async function* iterateAnthropicEvents(
throw new Error(sse.data);
}
if (sse.event === "ping") {
// Surface keepalives so the idle watchdog treats them as liveness.
yield ANTHROPIC_PING_EVENT;
continue;
}
if (!ANTHROPIC_MESSAGE_EVENTS.has(sse.event ?? "")) {
continue;
}
@@ -1104,7 +1125,7 @@ async function getAnthropicStreamResponse(
signal?: AbortSignal,
onSseEvent?: AnthropicOptions["onSseEvent"],
): Promise<{
events: AsyncIterable<RawMessageStreamEvent>;
events: AsyncIterable<AnthropicStreamEvent>;
response: Response;
requestId: string | null;
recordsRawSseEvents: boolean;
@@ -1126,9 +1147,9 @@ async function getAnthropicStreamResponse(
}
async function* observeDecodedAnthropicSdkEvents(
events: AsyncIterable<RawMessageStreamEvent>,
events: AsyncIterable<AnthropicStreamEvent>,
observer: (event: RawSseEvent) => void,
): AsyncGenerator<RawMessageStreamEvent> {
): AsyncGenerator<AnthropicStreamEvent> {
for await (const event of events) {
const data = JSON.stringify(event);
// Reconstructed from decoded SDK event; not literal wire bytes.
@@ -1207,6 +1228,14 @@ function isProviderRetryableStreamEnvelopeError(error: unknown): boolean {
export function isProviderRetryableError(error: unknown, provider?: string): boolean {
if (!(error instanceof Error)) return false;
if (provider === "github-copilot" && isCopilotTransientModelError(error)) return true;
// Account-level usage/quota limits ("usage_limit_reached", "exceed your
// account's rate limit", "quota exceeded") are persistent — the server
// parks the credential for minutes-to-hours (see the long `retry-after`).
// Retrying the same key with the provider's seconds-scale backoff never
// helps; these are owned by the credential-rotation layer (auth-gateway /
// `streamSimple` a/b/c policy), so surface them immediately instead of
// burning the retry budget here.
if (isUsageLimitError(error.message)) return false;
const msg = error.message.toLowerCase();
if (
isUnexpectedSocketCloseMessage(msg) ||
@@ -1415,7 +1444,7 @@ export const streamAnthropic: StreamFunction<"anthropic-messages"> = (
requestTimeoutMs,
);
}
let anthropicStream: AsyncIterable<RawMessageStreamEvent>;
let anthropicStream: AsyncIterable<AnthropicStreamEvent>;
let response: Response;
let requestId: string | null;
let recordsRawSseEvents: boolean;
+17
View File
@@ -55,6 +55,23 @@ describe("isProviderRetryableError", () => {
expect(isProviderRetryableError(new Error("Bad request"))).toBe(false);
});
it("does not retry persistent account usage/quota limits despite rate-limit wording", () => {
// Account-level 429 that says "rate limit" but is really a parked
// credential (long retry-after). Must surface immediately so the
// credential-rotation layer takes over instead of looping on backoff.
expect(
isProviderRetryableError(
new Error(
'429 {"type":"error","error":{"type":"rate_limit_error","message":"This request would exceed your account\'s rate limit. Please try again later."}}',
),
),
).toBe(false);
expect(isProviderRetryableError(new Error("usage_limit_reached"))).toBe(false);
expect(isProviderRetryableError(new Error("You have hit your ChatGPT usage limit"))).toBe(false);
// A generic transient rate limit (no account/usage framing) still retries.
expect(isProviderRetryableError(new Error("Rate limit exceeded"))).toBe(true);
});
it("retries Copilot transient model_not_supported only for github-copilot provider", () => {
const err = new Error("400 The requested model is not supported.");
(err as unknown as { status: number; code: string }).status = 400;
@@ -558,8 +558,8 @@ describe("Duplicate Tool Results Regression", () => {
];
const context: Context = { messages };
const wireMessages = convertMessages(providerModel, context, detectCompat(providerModel));
const assistantIds = assistantWireMessages(wireMessages).flatMap(message =>
message.tool_calls?.map(toolCall => toolCall.id) ?? [],
const assistantIds = assistantWireMessages(wireMessages).flatMap(
message => message.tool_calls?.map(toolCall => toolCall.id) ?? [],
);
expect(assistantIds, providerModel.provider).toEqual([duplicateId, expectedDuplicateId]);