Merge PR #8973: fix(ai): surface litellm concurrency-admission 429 immediately (@roboomp)

This commit is contained in:
can1357
2026-08-19 11:52:53 +02:00
3 changed files with 138 additions and 0 deletions
+1
View File
@@ -20,6 +20,7 @@
- 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 opencode-go (Console Go) rejecting Responses turns with `400 No tool output found for tool call …` (naming a random call of the batch on each retry) when a model streamed a trailing text/thinking block after its tool calls: `buildResponsesInput` emitted that block as an assistant `message` item wedged between the `function_call` batch and its `function_call_output` items. Such interleaved messages are now hoisted ahead of their call batch (canonical `message(s) → calls → outputs`), which the strict gateway validator accepts; content is unchanged ([#8789](https://github.com/can1357/oh-my-pi/issues/8789)).
- 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
+30
View File
@@ -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<string, string>;
@@ -69,6 +95,10 @@ export async function postOpenAIStream<TEvent>(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.
+107
View File
@@ -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);
});
});