Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
70 changes: 65 additions & 5 deletions src/agent/retry-policy.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -104,8 +104,9 @@ describe("createCorbitsRetryPolicy", () => {
raw: { error: { message: "Too Many Requests" } },
},
});
// Remapped to retryable -> default backoff, not abort on moderate Retry-After.
expect(decision).toEqual({ kind: "retry", delayMs: 500 });
// Remapped to retryable -> paced retry honors the server's Retry-After,
// not abort on moderate Retry-After and not a capped 30s wait.
expect(decision).toEqual({ kind: "retry", delayMs: 45_000 });
});

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

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

// CL-6910: the harness only surfaces `inference.error` to the director
Expand Down Expand Up @@ -242,7 +243,7 @@ describe("createCorbitsRetryPolicy", () => {
raw: { error: { message: "Too Many Requests" } },
},
};
expect(await decide(bare429)).toEqual({ kind: "retry", delayMs: 500 });
expect(await decide(bare429)).toEqual({ kind: "retry", delayMs: 45_000 });
current = "openai";
expect(await decide(bare429)).toEqual({ kind: "abort" });
});
Expand Down Expand Up @@ -298,4 +299,63 @@ describe("createCorbitsRetryPolicy", () => {
});
expect(notes).toHaveLength(0);
});

test("retryable 429 honors Retry-After instead of the fixed 500/1000ms backoff", async () => {
const decide = policy({ providerId: "codex/abk-labs" });
const situation = (attempt: number) => ({
attempt,
elapsedMs: 0,
error: {
category: "retryable" as const,
message: "Too Many Requests",
statusCode: 429,
retryAfterMs: 5_000,
},
});
expect(await decide(situation(1))).toEqual({ kind: "retry", delayMs: 5_000 });
expect(await decide(situation(2))).toEqual({ kind: "retry", delayMs: 5_000 });
expect(await decide(situation(3))).toEqual({ kind: "abort" });
});

test("retryable 429 honors a Retry-After above the blind-wait ceiling", async () => {
const decide = policy({ providerId: "codex/abk-labs" });
const decision = await decide({
attempt: 1,
elapsedMs: 0,
error: {
category: "retryable" as const,
message: "Too Many Requests",
statusCode: 429,
retryAfterMs: 120_000,
},
});
expect(decision).toEqual({ kind: "retry", delayMs: 120_000 });
});

test("retryable 429 with a day-long Retry-After aborts instead of hanging", async () => {
const decide = policy({ providerId: "codex/abk-labs" });
const decision = await decide({
attempt: 1,
elapsedMs: 0,
error: {
category: "retryable" as const,
message: "Too Many Requests",
statusCode: 429,
retryAfterMs: 86_400_000,
},
});
expect(decision).toEqual({ kind: "abort" });
});

test("retryable 429 without Retry-After keeps the fixed backoff", async () => {
const decide = policy();
const situation = (attempt: number) => ({
attempt,
elapsedMs: 0,
error: { category: "retryable" as const, message: "boom", statusCode: 429 },
});
expect(await decide(situation(1))).toEqual({ kind: "retry", delayMs: 500 });
expect(await decide(situation(2))).toEqual({ kind: "retry", delayMs: 1000 });
expect(await decide(situation(3))).toEqual({ kind: "abort" });
});
});
17 changes: 17 additions & 0 deletions src/agent/retry-policy.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import { getProcessAdmissionQueue, type AdmissionQueue } from "../subagent/admis
// so the user can switch providers or decide when to retry manually.
export const MAX_BLIND_WAIT_MS = 30_000;
const DEFAULT_PRESSURE_PAUSE_MS = 1_000;
const RATE_LIMIT_HANG_MS = 86_400_000;

export interface CorbitsRetryPolicyOptions {
/**
Expand Down Expand Up @@ -49,6 +50,22 @@ export function createCorbitsRetryPolicy(options?: CorbitsRetryPolicyOptions): R
const pauseMs = Math.min(error.retryAfterMs ?? DEFAULT_PRESSURE_PAUSE_MS, MAX_BLIND_WAIT_MS);
const provider = withProvider.providerId ?? stampedProviderId ?? "unknown";
admission.notePressure(provider, now() + pauseMs);
// The vendored default retries `retryable` on a fixed 500/1000ms
// schedule and ignores Retry-After. A 429 carries the server's pacing
// instruction: honor the full window. Capping at MAX_BLIND_WAIT_MS and
// retrying early burns the attempt budget while the server is still
// closed (the 45s xAI/Codex fixtures). Days-long Retry-After is a hang
// — abort rather than park the session. Attempt abort comes from
// defaultPolicy so this path cannot drift from MAX_ATTEMPTS.
if (error.retryAfterMs !== undefined) {
if (error.retryAfterMs >= RATE_LIMIT_HANG_MS) return { kind: "abort" };
const retryAfterMs = error.retryAfterMs;
const honorRetryAfter = (decision: RetryDecision): RetryDecision =>
decision.kind === "retry" ? { kind: "retry", delayMs: retryAfterMs } : decision;
const decision = defaultPolicy({ ...situation, error });
if (decision instanceof Promise) return decision.then(honorRetryAfter);
return honorRetryAfter(decision);
}
}
if (
error.category === "quota_exhausted" &&
Expand Down
20 changes: 20 additions & 0 deletions src/inference-error-message.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { describe, expect, test } from "bun:test";

import { normalizeInferenceErrorForTerminal } from "./inference-gateway-error.js";
import {
inferenceErrorMessage,
terminalProviderFailureMessage,
Expand Down Expand Up @@ -211,6 +212,25 @@ describe("terminalProviderFailureMessage", () => {
);
});

test("terminal Codex short-429 failure does not claim to still be retrying", () => {
const normalized = normalizeInferenceErrorForTerminal(
{ category: "quota_exhausted", message: "Too Many Requests", statusCode: 429 },
"codex/default",
);
const message = terminalProviderFailureMessage("codex/default", normalized);
expect(message.toLowerCase()).toMatch(/rate limit/);
expect(message.toLowerCase()).not.toContain("retrying");
});

test("retryable 429 guidance asks the operator to wait before trying again", () => {
const message = terminalProviderFailureMessage("codex/default", {
category: "retryable",
message: "Rate limited",
statusCode: 429,
});
expect(message).toContain("Wait a moment and try again.");
});

test("uses a safe label when the provider id contains only control sequences", () => {
const message = terminalProviderFailureMessage("\u001b[31m\u001b[0m", {
category: "fatal",
Expand Down
5 changes: 5 additions & 0 deletions src/inference-error-message.ts
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,11 @@ export function terminalProviderFailureMessage(
function terminalProviderFailureGuidance(error: InferenceErrorLike, category: string): string {
if (category === "credential_failure") return CREDENTIAL_FAILURE_USER_MESSAGE;
if (category === "context_overflow") return "Try /clear to start fresh.";
// A 429 that survived the harness's paced retries is a wait-it-out rate
// limit, not a generic flake: say so instead of the bare "Try again."
if (category === "retryable" && error.statusCode === 429) {
return "Wait a moment and try again.";
}
if (
category === "retryable" ||
(error.statusCode !== undefined && error.statusCode >= 500 && error.statusCode <= 599)
Expand Down
8 changes: 6 additions & 2 deletions src/inference-gateway-error.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,8 +49,12 @@ const GATEWAY_OVERLOAD_TEXT_MARKERS = [
/** User-visible line while the harness retries a transient gateway overload. */
export const GATEWAY_OVERLOAD_USER_MESSAGE = "Inference gateway overloaded — retrying…";

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

/** Body markers that mean a real usage/quota window, not a short rate limit. */
const XAI_QUOTA_BODY_MARKERS = [
Expand Down
7 changes: 7 additions & 0 deletions src/provider/codex-responses-adapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,13 @@ describe("createCodexResponsesAdapter", () => {
const adapter = createCodexResponsesAdapter(source);
expect(adapter.isStreamTerminal).toBe(isResponsesStreamTerminal);
});

test("extracts Retry-After pacing from response headers", () => {
const adapter = createCodexResponsesAdapter(source);
expect(adapter.extractRetryAfterMs?.(new Headers({ "retry-after": "7" }))).toBe(7_000);
expect(adapter.extractRetryAfterMs?.(new Headers({ "retry-after-ms": "1500" }))).toBe(1_500);
expect(adapter.extractRetryAfterMs?.(new Headers({}))).toBeUndefined();
});
});

describe("createCodexResponsesAdapter usage parsing", () => {
Expand Down
22 changes: 22 additions & 0 deletions src/provider/codex-responses-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -635,6 +635,27 @@ export function isResponsesStreamTerminal(sseData: string): boolean {
return typeof eventType === "string" && RESPONSES_TERMINAL_EVENTS.has(eventType);
}

// Responses backends (Codex, Grok, OpenAI) signal 429 pacing with the same
// `retry-after` / `retry-after-ms` headers the Chat Completions adapter
// already reads. The shared Responses adapters never extracted them, so
// every 429 arrived with retryAfterMs undefined and the retry policy fell
// back to blind fixed backoff instead of waiting out the server's window.
export function extractResponsesRetryAfterMs(headers: Headers): number | undefined {
const retryMs = headers.get("retry-after-ms");
if (retryMs !== null) {
const ms = Number(retryMs);
if (Number.isFinite(ms) && ms > 0) return Math.ceil(ms);
}
const raw = headers.get("retry-after");
if (raw !== null) {
const seconds = Number(raw);
if (Number.isFinite(seconds) && seconds > 0) {
return Math.ceil(seconds * 1000);
}
}
return undefined;
}

export function createCodexResponsesAdapter(source: LastCycleSource): ProviderAdapter {
// Re-created per request in buildRequest, not just once here — otherwise
// block indices accumulate across every request the adapter instance ever
Expand All @@ -648,5 +669,6 @@ export function createCodexResponsesAdapter(source: LastCycleSource): ProviderAd
parseResponse: (sseData) => parseResponse(sseData, indexer, source),
parseJSONResponse,
isStreamTerminal: isResponsesStreamTerminal,
extractRetryAfterMs: extractResponsesRetryAfterMs,
};
}
7 changes: 7 additions & 0 deletions src/provider/grok-responses-adapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -162,4 +162,11 @@ describe("createGrokResponsesAdapter", () => {
expect(body.reasoning).toEqual({ summary: "detailed" });
expect(body.reasoning?.effort).toBeUndefined();
});

test("extracts Retry-After pacing from response headers", () => {
const adapter = createGrokResponsesAdapter(source);
expect(adapter.extractRetryAfterMs?.(new Headers({ "retry-after": "7" }))).toBe(7_000);
expect(adapter.extractRetryAfterMs?.(new Headers({ "retry-after-ms": "1500" }))).toBe(1_500);
expect(adapter.extractRetryAfterMs?.(new Headers({}))).toBeUndefined();
});
});
2 changes: 2 additions & 0 deletions src/provider/grok-responses-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import {
import {
RESPONSES_TOOL_NAME_LIMIT,
createResponsesBlockIndexer,
extractResponsesRetryAfterMs,
parseJSONResponse,
parseResponse,
signatureForModel,
Expand Down Expand Up @@ -251,5 +252,6 @@ export function createGrokResponsesAdapter(source: LastCycleSource): ProviderAda
},
parseResponse: (sseData) => parseResponse(sseData, indexer, source, GROK_RESPONSES_PROVIDER),
parseJSONResponse,
extractRetryAfterMs: extractResponsesRetryAfterMs,
};
}
2 changes: 2 additions & 0 deletions src/provider/openai-responses-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import type {
import {
RESPONSES_TOOL_NAME_LIMIT,
createResponsesBlockIndexer,
extractResponsesRetryAfterMs,
isResponsesStreamTerminal,
parseJSONResponse,
parseResponse,
Expand Down Expand Up @@ -233,5 +234,6 @@ export function createOpenAIResponsesAdapter(source: LastCycleSource): ProviderA
parseResponse: (sseData) => parseResponse(sseData, indexer, source, OPENAI_RESPONSES_PROVIDER),
parseJSONResponse,
isStreamTerminal: isResponsesStreamTerminal,
extractRetryAfterMs: extractResponsesRetryAfterMs,
};
}
2 changes: 1 addition & 1 deletion src/tui/stream-event-map.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -421,7 +421,7 @@ describe("inference.error text", () => {
mapProductionEvent({ type: "connector.reply", data: { content: "generic reply" } }, ctx),
).toContainEqual({
type: "assistant",
text: "Work Provider failed (retryable): Rate limited — retrying…. Try again.",
text: "Work Provider failed (retryable): Rate limited. Wait a moment and try again.",
});
});

Expand Down
9 changes: 9 additions & 0 deletions tests/unit/openai-responses-adapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -89,3 +89,12 @@ describe("openai-responses x-opencode-session header", () => {
expect(req.headers["x-opencode-session"]).toBeUndefined();
});
});

describe("openai-responses Retry-After extraction", () => {
test("extracts Retry-After pacing from response headers", () => {
const responses = adapter();
expect(responses.extractRetryAfterMs?.(new Headers({ "retry-after": "7" }))).toBe(7_000);
expect(responses.extractRetryAfterMs?.(new Headers({ "retry-after-ms": "1500" }))).toBe(1_500);
expect(responses.extractRetryAfterMs?.(new Headers({}))).toBeUndefined();
});
});
Loading