Skip to content

Commit 48ea82c

Browse files
Merge pull request #839 from corbitsdev/cl-7523-fix-rate-limit-retry-to-honor-retry-after-header-and-clarify
Honor Retry-After header in rate-limit retry policy and clarify terminal message
2 parents 6ea5969 + 39e88fa commit 48ea82c

12 files changed

Lines changed: 163 additions & 8 deletions

src/agent/retry-policy.test.ts

Lines changed: 65 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -104,8 +104,9 @@ describe("createCorbitsRetryPolicy", () => {
104104
raw: { error: { message: "Too Many Requests" } },
105105
},
106106
});
107-
// Remapped to retryable -> default backoff, not abort on moderate Retry-After.
108-
expect(decision).toEqual({ kind: "retry", delayMs: 500 });
107+
// Remapped to retryable -> paced retry honors the server's Retry-After,
108+
// not abort on moderate Retry-After and not a capped 30s wait.
109+
expect(decision).toEqual({ kind: "retry", delayMs: 45_000 });
109110
});
110111

111112
test("stamped Codex usage-limit 429 retries as retryable, not long-quota abort", async () => {
@@ -120,7 +121,7 @@ describe("createCorbitsRetryPolicy", () => {
120121
raw: "You have hit your ChatGPT usage limit",
121122
},
122123
});
123-
expect(decision).toEqual({ kind: "retry", delayMs: 500 });
124+
expect(decision).toEqual({ kind: "retry", delayMs: 45_000 });
124125
});
125126

126127
test("stamped xAI usage/quota body still aborts on long retryAfterMs", async () => {
@@ -174,7 +175,7 @@ describe("createCorbitsRetryPolicy", () => {
174175
};
175176
expect(await decide(bare429)).toEqual({ kind: "abort" });
176177
current = "xai/thegreataxios";
177-
expect(await decide(bare429)).toEqual({ kind: "retry", delayMs: 500 });
178+
expect(await decide(bare429)).toEqual({ kind: "retry", delayMs: 45_000 });
178179
});
179180

180181
// CL-6910: the harness only surfaces `inference.error` to the director
@@ -242,7 +243,7 @@ describe("createCorbitsRetryPolicy", () => {
242243
raw: { error: { message: "Too Many Requests" } },
243244
},
244245
};
245-
expect(await decide(bare429)).toEqual({ kind: "retry", delayMs: 500 });
246+
expect(await decide(bare429)).toEqual({ kind: "retry", delayMs: 45_000 });
246247
current = "openai";
247248
expect(await decide(bare429)).toEqual({ kind: "abort" });
248249
});
@@ -298,4 +299,63 @@ describe("createCorbitsRetryPolicy", () => {
298299
});
299300
expect(notes).toHaveLength(0);
300301
});
302+
303+
test("retryable 429 honors Retry-After instead of the fixed 500/1000ms backoff", async () => {
304+
const decide = policy({ providerId: "codex/abk-labs" });
305+
const situation = (attempt: number) => ({
306+
attempt,
307+
elapsedMs: 0,
308+
error: {
309+
category: "retryable" as const,
310+
message: "Too Many Requests",
311+
statusCode: 429,
312+
retryAfterMs: 5_000,
313+
},
314+
});
315+
expect(await decide(situation(1))).toEqual({ kind: "retry", delayMs: 5_000 });
316+
expect(await decide(situation(2))).toEqual({ kind: "retry", delayMs: 5_000 });
317+
expect(await decide(situation(3))).toEqual({ kind: "abort" });
318+
});
319+
320+
test("retryable 429 honors a Retry-After above the blind-wait ceiling", async () => {
321+
const decide = policy({ providerId: "codex/abk-labs" });
322+
const decision = await decide({
323+
attempt: 1,
324+
elapsedMs: 0,
325+
error: {
326+
category: "retryable" as const,
327+
message: "Too Many Requests",
328+
statusCode: 429,
329+
retryAfterMs: 120_000,
330+
},
331+
});
332+
expect(decision).toEqual({ kind: "retry", delayMs: 120_000 });
333+
});
334+
335+
test("retryable 429 with a day-long Retry-After aborts instead of hanging", async () => {
336+
const decide = policy({ providerId: "codex/abk-labs" });
337+
const decision = await decide({
338+
attempt: 1,
339+
elapsedMs: 0,
340+
error: {
341+
category: "retryable" as const,
342+
message: "Too Many Requests",
343+
statusCode: 429,
344+
retryAfterMs: 86_400_000,
345+
},
346+
});
347+
expect(decision).toEqual({ kind: "abort" });
348+
});
349+
350+
test("retryable 429 without Retry-After keeps the fixed backoff", async () => {
351+
const decide = policy();
352+
const situation = (attempt: number) => ({
353+
attempt,
354+
elapsedMs: 0,
355+
error: { category: "retryable" as const, message: "boom", statusCode: 429 },
356+
});
357+
expect(await decide(situation(1))).toEqual({ kind: "retry", delayMs: 500 });
358+
expect(await decide(situation(2))).toEqual({ kind: "retry", delayMs: 1000 });
359+
expect(await decide(situation(3))).toEqual({ kind: "abort" });
360+
});
301361
});

src/agent/retry-policy.ts

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import { getProcessAdmissionQueue, type AdmissionQueue } from "../subagent/admis
1313
// so the user can switch providers or decide when to retry manually.
1414
export const MAX_BLIND_WAIT_MS = 30_000;
1515
const DEFAULT_PRESSURE_PAUSE_MS = 1_000;
16+
const RATE_LIMIT_HANG_MS = 86_400_000;
1617

1718
export interface CorbitsRetryPolicyOptions {
1819
/**
@@ -49,6 +50,22 @@ export function createCorbitsRetryPolicy(options?: CorbitsRetryPolicyOptions): R
4950
const pauseMs = Math.min(error.retryAfterMs ?? DEFAULT_PRESSURE_PAUSE_MS, MAX_BLIND_WAIT_MS);
5051
const provider = withProvider.providerId ?? stampedProviderId ?? "unknown";
5152
admission.notePressure(provider, now() + pauseMs);
53+
// The vendored default retries `retryable` on a fixed 500/1000ms
54+
// schedule and ignores Retry-After. A 429 carries the server's pacing
55+
// instruction: honor the full window. Capping at MAX_BLIND_WAIT_MS and
56+
// retrying early burns the attempt budget while the server is still
57+
// closed (the 45s xAI/Codex fixtures). Days-long Retry-After is a hang
58+
// — abort rather than park the session. Attempt abort comes from
59+
// defaultPolicy so this path cannot drift from MAX_ATTEMPTS.
60+
if (error.retryAfterMs !== undefined) {
61+
if (error.retryAfterMs >= RATE_LIMIT_HANG_MS) return { kind: "abort" };
62+
const retryAfterMs = error.retryAfterMs;
63+
const honorRetryAfter = (decision: RetryDecision): RetryDecision =>
64+
decision.kind === "retry" ? { kind: "retry", delayMs: retryAfterMs } : decision;
65+
const decision = defaultPolicy({ ...situation, error });
66+
if (decision instanceof Promise) return decision.then(honorRetryAfter);
67+
return honorRetryAfter(decision);
68+
}
5269
}
5370
if (
5471
error.category === "quota_exhausted" &&

src/inference-error-message.test.ts

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import { describe, expect, test } from "bun:test";
22

3+
import { normalizeInferenceErrorForTerminal } from "./inference-gateway-error.js";
34
import {
45
inferenceErrorMessage,
56
terminalProviderFailureMessage,
@@ -211,6 +212,25 @@ describe("terminalProviderFailureMessage", () => {
211212
);
212213
});
213214

215+
test("terminal Codex short-429 failure does not claim to still be retrying", () => {
216+
const normalized = normalizeInferenceErrorForTerminal(
217+
{ category: "quota_exhausted", message: "Too Many Requests", statusCode: 429 },
218+
"codex/default",
219+
);
220+
const message = terminalProviderFailureMessage("codex/default", normalized);
221+
expect(message.toLowerCase()).toMatch(/rate limit/);
222+
expect(message.toLowerCase()).not.toContain("retrying");
223+
});
224+
225+
test("retryable 429 guidance asks the operator to wait before trying again", () => {
226+
const message = terminalProviderFailureMessage("codex/default", {
227+
category: "retryable",
228+
message: "Rate limited",
229+
statusCode: 429,
230+
});
231+
expect(message).toContain("Wait a moment and try again.");
232+
});
233+
214234
test("uses a safe label when the provider id contains only control sequences", () => {
215235
const message = terminalProviderFailureMessage("\u001b[31m\u001b[0m", {
216236
category: "fatal",

src/inference-error-message.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -130,6 +130,11 @@ export function terminalProviderFailureMessage(
130130
function terminalProviderFailureGuidance(error: InferenceErrorLike, category: string): string {
131131
if (category === "credential_failure") return CREDENTIAL_FAILURE_USER_MESSAGE;
132132
if (category === "context_overflow") return "Try /clear to start fresh.";
133+
// A 429 that survived the harness's paced retries is a wait-it-out rate
134+
// limit, not a generic flake: say so instead of the bare "Try again."
135+
if (category === "retryable" && error.statusCode === 429) {
136+
return "Wait a moment and try again.";
137+
}
133138
if (
134139
category === "retryable" ||
135140
(error.statusCode !== undefined && error.statusCode >= 500 && error.statusCode <= 599)

src/inference-gateway-error.ts

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -49,8 +49,12 @@ const GATEWAY_OVERLOAD_TEXT_MARKERS = [
4949
/** User-visible line while the harness retries a transient gateway overload. */
5050
export const GATEWAY_OVERLOAD_USER_MESSAGE = "Inference gateway overloaded — retrying…";
5151

52-
/** User-visible line while the harness retries a short known-provider HTTP 429. */
53-
export const RATE_LIMIT_USER_MESSAGE = "Rate limited — retrying…";
52+
/**
53+
* User-visible line for a short known-provider HTTP 429. Worded without
54+
* "retrying": this message also surfaces terminally after the harness has
55+
* exhausted its retries, where claiming an ongoing retry is wrong.
56+
*/
57+
export const RATE_LIMIT_USER_MESSAGE = "Rate limited";
5458

5559
/** Body markers that mean a real usage/quota window, not a short rate limit. */
5660
const XAI_QUOTA_BODY_MARKERS = [

src/provider/codex-responses-adapter.test.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,13 @@ describe("createCodexResponsesAdapter", () => {
6363
const adapter = createCodexResponsesAdapter(source);
6464
expect(adapter.isStreamTerminal).toBe(isResponsesStreamTerminal);
6565
});
66+
67+
test("extracts Retry-After pacing from response headers", () => {
68+
const adapter = createCodexResponsesAdapter(source);
69+
expect(adapter.extractRetryAfterMs?.(new Headers({ "retry-after": "7" }))).toBe(7_000);
70+
expect(adapter.extractRetryAfterMs?.(new Headers({ "retry-after-ms": "1500" }))).toBe(1_500);
71+
expect(adapter.extractRetryAfterMs?.(new Headers({}))).toBeUndefined();
72+
});
6673
});
6774

6875
describe("createCodexResponsesAdapter usage parsing", () => {

src/provider/codex-responses-adapter.ts

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -635,6 +635,27 @@ export function isResponsesStreamTerminal(sseData: string): boolean {
635635
return typeof eventType === "string" && RESPONSES_TERMINAL_EVENTS.has(eventType);
636636
}
637637

638+
// Responses backends (Codex, Grok, OpenAI) signal 429 pacing with the same
639+
// `retry-after` / `retry-after-ms` headers the Chat Completions adapter
640+
// already reads. The shared Responses adapters never extracted them, so
641+
// every 429 arrived with retryAfterMs undefined and the retry policy fell
642+
// back to blind fixed backoff instead of waiting out the server's window.
643+
export function extractResponsesRetryAfterMs(headers: Headers): number | undefined {
644+
const retryMs = headers.get("retry-after-ms");
645+
if (retryMs !== null) {
646+
const ms = Number(retryMs);
647+
if (Number.isFinite(ms) && ms > 0) return Math.ceil(ms);
648+
}
649+
const raw = headers.get("retry-after");
650+
if (raw !== null) {
651+
const seconds = Number(raw);
652+
if (Number.isFinite(seconds) && seconds > 0) {
653+
return Math.ceil(seconds * 1000);
654+
}
655+
}
656+
return undefined;
657+
}
658+
638659
export function createCodexResponsesAdapter(source: LastCycleSource): ProviderAdapter {
639660
// Re-created per request in buildRequest, not just once here — otherwise
640661
// block indices accumulate across every request the adapter instance ever
@@ -648,5 +669,6 @@ export function createCodexResponsesAdapter(source: LastCycleSource): ProviderAd
648669
parseResponse: (sseData) => parseResponse(sseData, indexer, source),
649670
parseJSONResponse,
650671
isStreamTerminal: isResponsesStreamTerminal,
672+
extractRetryAfterMs: extractResponsesRetryAfterMs,
651673
};
652674
}

src/provider/grok-responses-adapter.test.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -162,4 +162,11 @@ describe("createGrokResponsesAdapter", () => {
162162
expect(body.reasoning).toEqual({ summary: "detailed" });
163163
expect(body.reasoning?.effort).toBeUndefined();
164164
});
165+
166+
test("extracts Retry-After pacing from response headers", () => {
167+
const adapter = createGrokResponsesAdapter(source);
168+
expect(adapter.extractRetryAfterMs?.(new Headers({ "retry-after": "7" }))).toBe(7_000);
169+
expect(adapter.extractRetryAfterMs?.(new Headers({ "retry-after-ms": "1500" }))).toBe(1_500);
170+
expect(adapter.extractRetryAfterMs?.(new Headers({}))).toBeUndefined();
171+
});
165172
});

src/provider/grok-responses-adapter.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import {
1919
import {
2020
RESPONSES_TOOL_NAME_LIMIT,
2121
createResponsesBlockIndexer,
22+
extractResponsesRetryAfterMs,
2223
parseJSONResponse,
2324
parseResponse,
2425
signatureForModel,
@@ -251,5 +252,6 @@ export function createGrokResponsesAdapter(source: LastCycleSource): ProviderAda
251252
},
252253
parseResponse: (sseData) => parseResponse(sseData, indexer, source, GROK_RESPONSES_PROVIDER),
253254
parseJSONResponse,
255+
extractRetryAfterMs: extractResponsesRetryAfterMs,
254256
};
255257
}

src/provider/openai-responses-adapter.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import type {
1313
import {
1414
RESPONSES_TOOL_NAME_LIMIT,
1515
createResponsesBlockIndexer,
16+
extractResponsesRetryAfterMs,
1617
isResponsesStreamTerminal,
1718
parseJSONResponse,
1819
parseResponse,
@@ -233,5 +234,6 @@ export function createOpenAIResponsesAdapter(source: LastCycleSource): ProviderA
233234
parseResponse: (sseData) => parseResponse(sseData, indexer, source, OPENAI_RESPONSES_PROVIDER),
234235
parseJSONResponse,
235236
isStreamTerminal: isResponsesStreamTerminal,
237+
extractRetryAfterMs: extractResponsesRetryAfterMs,
236238
};
237239
}

0 commit comments

Comments
 (0)