From 3d844bf3b231a26074327de7a5efcec2a7a10d66 Mon Sep 17 00:00:00 2001 From: roboomp Date: Wed, 19 Aug 2026 08:55:17 +0000 Subject: [PATCH] fix(ai): surface litellm concurrency-admission 429 immediately The OpenAI-wire transport called fetchWithRetry with maxAttempts: 6 and the default 60s maxDelayMs cap, so a LiteLLM max_parallel_requests rejection (HTTP 429, Retry-After: 60) was slept-and-retried up to six times before TurnRecovery ever saw it. A 60s hint equals the cap, so fetchWithRetry never bailed early and one turn could stall ~300s, bypassing the user's retry.maxDelayMs/maxRetries and the session-level CONCURRENT_LIMIT backoff + model fallback. postOpenAIStream now opts out of transport-level retry for this concurrency-admission response class via fetchWithRetry's shouldRetryResponse gate, detecting the rate_limit_type=max_parallel_requests marker in the response header or structured body. The 429 surfaces on the first attempt so session recovery owns retry/fallback. Genuine RPM/quota 429s carry no such marker and keep honoring Retry-After. Fixes #8854 --- packages/ai/CHANGELOG.md | 1 + packages/ai/src/utils/openai-http.ts | 30 ++++++ packages/ai/test/issue-8854-repro.test.ts | 107 ++++++++++++++++++++++ 3 files changed, 138 insertions(+) create mode 100644 packages/ai/test/issue-8854-repro.test.ts diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index cb40dada2..07783186f 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -19,6 +19,7 @@ - Cloud Code Assist Gemini 3.6/3.7 Flash requests at `minimal` now send `thinkingLevel: LOW` on the aliased `-low` SKU instead of `MINIMAL`, which the API rejects with HTTP 400. - Answer Cursor `interaction_query` permission gates (hosted web search, Exa, unnamed field-9 WebFetch) so the Run RPC continues instead of sitting silent until the 300s idle watchdog. - Fixed provider tool calls arriving with flattened array argument paths (e.g. Gemini's `questions[0].id`) being stripped and rejected by argument validation; well-formed flattened paths are now rebuilt into the nested arrays the tool schema expects ([#8886](https://github.com/can1357/oh-my-pi/issues/8886)). +- Fixed the OpenAI-wire transport sleeping on a LiteLLM concurrency-admission 429 (`rate_limit_type: max_parallel_requests`, `Retry-After: 60`) and retrying it up to 6 times (~300s) before session recovery saw the error. Because a 60s hint equals the transport's `maxDelayMs` cap, `fetchWithRetry` kept sleeping and retrying; the request now surfaces on the first attempt so `TurnRecovery`'s concurrency backoff/model fallback runs promptly. Genuine RPM/quota 429s (no such marker) still honor `Retry-After` ([#8854](https://github.com/can1357/oh-my-pi/issues/8854)). ## [17.3.7] - 2026-08-17 diff --git a/packages/ai/src/utils/openai-http.ts b/packages/ai/src/utils/openai-http.ts index 0dab0231a..7757a17e8 100644 --- a/packages/ai/src/utils/openai-http.ts +++ b/packages/ai/src/utils/openai-http.ts @@ -35,6 +35,32 @@ const DEFAULT_MAX_ATTEMPTS = 6; /** Bound the `Error.message` allocation for proxy HTML error pages and the like. */ const MAX_DETAIL_CHARS = 4096; +/** + * LiteLLM (and compatible proxies) shed over-concurrency requests *before* the + * upstream call with an immediate HTTP 429 marked `rate_limit_type: + * max_parallel_requests` — as a response header and/or a structured body field. + * This is an admission failure, not an upstream rate/quota limit: the request + * never reached a model. Retrying it inside the transport (honoring the proxy's + * `Retry-After`, up to {@link DEFAULT_MAX_ATTEMPTS} times) duplicates — worse, + * at 60s per sleep instead of 5s — the concurrency backoff and model fallback + * that `TurnRecovery` already owns, stalling one turn for up to ~300s + * (issue #8854). {@link isConcurrencyAdmissionRejection} lets the transport + * surface it on the first attempt so session recovery runs promptly. Genuine + * RPM/quota 429s carry no such marker and keep honoring `Retry-After`. + */ +const CONCURRENCY_ADMISSION_LIMITER = "max_parallel_requests"; + +/** Body form of the marker: `"rate_limit_type": "max_parallel_requests"` (top level or under `error`). */ +const CONCURRENCY_ADMISSION_BODY_PATTERN = /"rate_limit_type"\s*:\s*"max_parallel_requests"/; + +/** `true` for a proxy concurrency-admission 429 that must bypass transport-level retry. */ +function isConcurrencyAdmissionRejection(response: Response, bodyText: string): boolean { + return ( + response.headers.get("rate_limit_type")?.trim() === CONCURRENCY_ADMISSION_LIMITER || + CONCURRENCY_ADMISSION_BODY_PATTERN.test(bodyText) + ); +} + export interface OpenAIStreamRequestInit { url: string; headers: Record; @@ -69,6 +95,10 @@ export async function postOpenAIStream(init: OpenAIStreamRequestInit): P signal: init.signal, fetch: init.fetch, maxAttempts: DEFAULT_MAX_ATTEMPTS, + // A proxy concurrency-admission 429 (`rate_limit_type: max_parallel_requests`) + // surfaces immediately instead of being slept-and-retried here; session + // recovery owns its backoff/fallback (issue #8854). + shouldRetryResponse: (response, bodyText) => !isConcurrencyAdmissionRejection(response, bodyText), // Bun's native fetch enforces a hard ~300s pre-response timeout (issue #2422). // Cold large-context streams legitimately exceed it; the caller's // `firstEventTimeoutMs`/`AbortSignal` already govern stuck requests. diff --git a/packages/ai/test/issue-8854-repro.test.ts b/packages/ai/test/issue-8854-repro.test.ts new file mode 100644 index 000000000..5b3070771 --- /dev/null +++ b/packages/ai/test/issue-8854-repro.test.ts @@ -0,0 +1,107 @@ +import { describe, expect, test } from "bun:test"; +import { OpenAIHttpError, postOpenAIStream } from "../src/utils/openai-http"; +import { mockFetch } from "./helpers/fetch-mock"; + +// LiteLLM (and compatible proxies) shed over-concurrency requests before the +// upstream call with an immediate HTTP 429 marked `rate_limit_type: +// max_parallel_requests` and `Retry-After: 60`. Because a 60s hint equals the +// transport's `maxDelayMs` cap, `fetchWithRetry` used to sleep and retry it up +// to 6 times (~300s) before the error ever reached `TurnRecovery`, which owns +// the real concurrency backoff + model fallback. The transport must surface +// this admission failure on the first attempt instead. Regression guard for #8854. +describe("OpenAI transport concurrency-admission 429 (#8854)", () => { + const concurrencyBody = JSON.stringify({ + error: { + message: "Max parallel request limit reached", + type: "rate_limit_error", + rate_limit_type: "max_parallel_requests", + }, + }); + + test("surfaces a header-marked limiter 429 on the first attempt", async () => { + let attempts = 0; + const fetch = mockFetch(() => { + attempts++; + return new Response(concurrencyBody, { + status: 429, + headers: { + "content-type": "application/json", + "retry-after": "60", + rate_limit_type: "max_parallel_requests", + }, + }); + }); + + const error = await postOpenAIStream({ + url: "https://litellm.local/v1/chat/completions", + headers: {}, + body: { model: "gpt-4o", messages: [] }, + signal: new AbortController().signal, + fetch, + }).then( + () => undefined, + (err: unknown) => err, + ); + + expect(attempts).toBe(1); + expect(error).toBeInstanceOf(OpenAIHttpError); + expect((error as OpenAIHttpError).status).toBe(429); + }); + + test("surfaces a body-marked limiter 429 even without the header", async () => { + let attempts = 0; + const fetch = mockFetch(() => { + attempts++; + return new Response(concurrencyBody, { + status: 429, + headers: { "content-type": "application/json", "retry-after": "60" }, + }); + }); + + const error = await postOpenAIStream({ + url: "https://litellm.local/v1/chat/completions", + headers: {}, + body: { model: "gpt-4o", messages: [] }, + signal: new AbortController().signal, + fetch, + }).then( + () => undefined, + (err: unknown) => err, + ); + + expect(attempts).toBe(1); + expect(error).toBeInstanceOf(OpenAIHttpError); + expect((error as OpenAIHttpError).status).toBe(429); + }); + + // Scope guard: the opt-out must not globally disable Retry-After. A generic + // RPM/quota 429 without the concurrency marker is still retried. + test("still retries a generic 429 that lacks the concurrency marker", async () => { + let attempts = 0; + const fetch = mockFetch(() => { + attempts++; + if (attempts === 1) { + return new Response(JSON.stringify({ error: { message: "Rate limit reached" } }), { + status: 429, + // Short delta hint so the retry sleep is negligible in-test. + headers: { "content-type": "application/json", "retry-after-ms": "5" }, + }); + } + return new Response("data: [DONE]\n\n", { + status: 200, + headers: { "content-type": "text/event-stream" }, + }); + }); + + const handle = await postOpenAIStream({ + url: "https://api.openai.com/v1/chat/completions", + headers: {}, + body: { model: "gpt-4o", messages: [] }, + signal: new AbortController().signal, + fetch, + }); + + expect(attempts).toBe(2); + expect(handle.response.status).toBe(200); + }); +});