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
1 change: 1 addition & 0 deletions changelog.d/features/claude-low-priority-mode.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **feat(sse):** Claude OAuth connections can opt in (per account, Edit connection → Claude section) to Claude Code's lower-priority lane and once-a-week session-limit reset. After the first 5-hour usage-wall 429 carrying `anthropic-ratelimit-unified-slow-offer: treatment`, OmniRoute retries the same account with `anthropic-usage-limit: slow` and keeps sending it until the window resets — the account keeps serving past the limit instead of being cooled down (slot_busy/529 wait the server's `slow-retry-after`, bounded by `slow-max-wait`). With auto-reset on, the wall first tries `POST /api/organizations/{org}/reset_rate_limits` (`juniper_tide`) and retries at full speed when the server grants it. Both default off; nothing is sent before the limit is hit.
60 changes: 60 additions & 0 deletions docs/architecture/RESILIENCE_GUIDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,66 @@ These persist until credentials change or an operator resets them. Do not overwr

**Lazy recovery:** when `rateLimitedUntil` is past, connection becomes eligible again. On successful use, `clearAccountError()` clears all error fields.

### Claude OAuth usage wall: lower-priority lane + session-limit reset

**Scope:** one Claude subscription (OAuth) connection. Both features are **opt-in per
connection** (Edit connection → Claude section → `lowPriorityMode` / `autoLimitReset` in
`providerSpecificData`, both default off) and mirror Claude Code's `/low-priority` and
`/limit-reset` commands (wire contract captured from Claude Code 2.1.263).

**Implementation:**

- State machine + response classification: `open-sse/services/claudeLowPriority.ts`
- Reset status/claim client: `open-sse/services/claudeLimitReset.ts`
- Executor hook (header injection + same-account retry): `open-sse/executors/base.ts::execute()`
- Opt-in persistence: `src/lib/providers/requestDefaults.ts::normalizeProviderSpecificData()`

**Trigger:** the 5-hour usage wall — a `429` whose headers carry
`anthropic-ratelimit-unified-status: rejected` and, when the account is eligible,
`anthropic-ratelimit-unified-slow-offer: treatment`. Nothing is sent before that first wall
429; a burst 429 without unified headers goes through the normal cooldown path.

**Lower-priority lane** (`lowPriorityMode`):

- On the wall 429 the executor accepts the offer and immediately retries the **same**
account with `anthropic-usage-limit: slow`; the lane stays active until the announced
`anthropic-ratelimit-unified-reset` (+60s grace) and every request in that window carries
the header. The intercepted 429 never reaches `handleChatCore`, so the connection is
**not** put in cooldown and is not rotated away.
- `anthropic-ratelimit-unified-slow-status` on later responses: `active` / `not_needed`
keep the lane; `slot_busy` (429) or a `529` wait the server's
`anthropic-ratelimit-unified-slow-retry-after` (default 20s, clamp 5–600s, ±30% jitter)
and retry, bounded by `anthropic-ratelimit-unified-slow-max-wait` (default 20 min, clamp
1 min–6 h) — past that the lane ends and a 10-minute cool-off blocks re-acceptance. The
wait is additionally capped by what is left of the request's own upstream-start timeout
(`resolveFetchStartTimeout`, 10 min by default) minus a 5 s margin: without that cap the
20-minute default max-wait would outlive the request and the sleep would be aborted
mid-wait, surfacing a `TimeoutError` instead of the graceful `max_wait` end + cool-off.
- `weekly_limit` / `budget_exhausted` / `off` / `ineligible`, a 5h-window rollover, or
`ineligible` + `anthropic-ratelimit-unified-overage-in-use: true` (which ends it as
`extra_usage` on any status, since paid overage now covers the wall) end the lane; the
response then flows to the normal cooldown path. `budget_exhausted` is remembered until
the announced budget reset (≤ 8 days).
- The wall check runs after the executor's own 400-driven intra-attempt retries (context
editing, thinking/effort clamps, param auto-learn), so a wall 429 that only surfaces on
one of those retries is still intercepted instead of reaching the cooldown path.
- State is in-memory per connection (a restart costs one extra wall 429 to re-accept).

**Session-limit reset** (`autoLimitReset`, tried before the lane when both are on):

- `GET https://api.anthropic.com/api/oauth/usage?at_wall=1&skip_spend=1` → `juniper_tide`
block; when `arm: "reset"` and `available: true`,
`POST https://api.anthropic.com/api/organizations/{orgUUID}/reset_rate_limits` with
`{ "program": "juniper_tide" }` (organization UUID from
`providerSpecificData.organizationUUID`, bootstrap fallback).
- `result: reset|not_limited` → the request is retried at full speed (no slow header).
`already_used` / `not_offered` memoise `next_available_at` (default one week); any
failure backs off 15 minutes. The reset is once a week and still counts toward the
weekly limit.

Regression guards: `tests/unit/claude-low-priority-mode.test.ts`,
`tests/unit/claude-limit-reset.test.ts`, `tests/unit/claude-low-priority-executor.test.ts`.

### Session affinity (#7274)

**Scope:** one client session (`X-Session-Id` / `x-codex-session-id` / `x-omniroute-session` header) pinned to one connection, for **any** provider.
Expand Down
34 changes: 34 additions & 0 deletions open-sse/config/anthropicHeaders.ts
Original file line number Diff line number Diff line change
Expand Up @@ -182,3 +182,37 @@ export const CLAUDE_CLI_USER_AGENT = getClaudeCodeUserAgent("cli");
export { getClaudeCodeUserAgent };
export const CLAUDE_CLI_STAINLESS_PACKAGE_VERSION = CLAUDE_CODE_SDK_PACKAGE_VERSION;
export const CLAUDE_CLI_STAINLESS_RUNTIME_VERSION = CLAUDE_CODE_RUNTIME_VERSION;

/**
* Merge a Claude-Code-shaped header set over the outbound headers, dropping any
* case variant of the same name first — undici would otherwise concatenate the two
* into a single rejected value (issue #1454).
*/
export function mergeCcHeaders(
headers: Record<string, string>,
ccHeaders: Record<string, string>
): void {
const ccKeysLower = new Set(Object.keys(ccHeaders).map((k) => k.toLowerCase()));
for (const key of Object.keys(headers)) {
if (ccKeysLower.has(key.toLowerCase())) delete headers[key];
}
Object.assign(headers, ccHeaders);
}

/**
* Stainless SDK metadata for the Claude wire image. OS/arch follow the host running
* the signed binary; the runtime version is pinned to the captured CLI, not OmniRoute's
* Node. Mutates `headers`.
*/
export function applyStainlessHeaders(
headers: Record<string, string>,
parts: { arch: string; os: string }
): void {
headers["X-Stainless-Arch"] = parts.arch;
headers["X-Stainless-Lang"] = "js";
headers["X-Stainless-OS"] = parts.os;
headers["X-Stainless-Runtime"] = "node";
headers["X-Stainless-Runtime-Version"] = CLAUDE_CLI_STAINLESS_RUNTIME_VERSION;
headers["X-Stainless-Retry-Count"] = "0";
delete headers["X-Stainless-Os"];
}
56 changes: 26 additions & 30 deletions open-sse/executors/base.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,9 @@ import {
type AlternateFormat,
} from "../config/providers/alternateFormats.ts";
import {
CLAUDE_CLI_STAINLESS_RUNTIME_VERSION,
applyStainlessHeaders,
getClaudeCliBillingVersion,
mergeCcHeaders,
mergeClientAnthropicBeta,
normalizeAnthropicHeaderVariants,
} from "../config/anthropicHeaders.ts";
Expand Down Expand Up @@ -40,6 +41,7 @@ import {
isFreeVariantModel,
} from "../services/openrouterFreeWindow.ts";
import { gateOutboundRequest } from "../services/wafRateLimit.ts";
import { ClaudeUsageLimitGuard } from "./claudeUsageLimit.ts";
import type { PoolConfig } from "../services/sessionPool/types.ts";
import type { Session } from "../services/sessionPool/session.ts";
import { SessionPool } from "../services/sessionPool/sessionPool.ts";
Expand Down Expand Up @@ -97,6 +99,7 @@ import {
selectBetaFlags,
stainlessArch,
stainlessOS,
stripClaudeSystemPrefixBlocks,
stripProxyToolPrefix,
} from "./claudeIdentity.ts";
import { withForcedResponsesUpstream } from "./forceResponsesUpstream.ts";
Expand Down Expand Up @@ -728,6 +731,8 @@ export class BaseExecutor {
let activeCredentials = credentials;
// Track per-URL intra-retry attempts to avoid infinite loops
const retryAttemptsByUrl: Record<number, number> = {};
// Claude OAuth usage wall (opt-in per connection): see ./claudeUsageLimit.ts.
const claudeUsageLimit = new ClaudeUsageLimitGuard(this.provider, log);

// Probe-origin dispatches must not consume a refresh-token rotation —
// routing state untouched; the reactive 401/403 path is probe-guarded
Expand Down Expand Up @@ -1177,18 +1182,7 @@ export class BaseExecutor {
// Strip any pre-existing billing/sentinel before re-prepending — keeps
// retries idempotent and avoids stacking that breaks prompt-cache prefix
// matching (see issue #1712).
for (let i = sysBlocks.length - 1; i >= 0; i--) {
const t = sysBlocks[i]?.text;
if (typeof t === "string" && t.startsWith("x-anthropic-billing-header:")) {
sysBlocks.splice(i, 1);
}
}
for (let i = sysBlocks.length - 1; i >= 0; i--) {
const t = sysBlocks[i]?.text;
if (typeof t === "string" && t.startsWith(SENTINEL)) {
sysBlocks.splice(i, 1);
}
}
stripClaudeSystemPrefixBlocks(sysBlocks, SENTINEL);
sysBlocks.unshift({ type: "text", text: billingLine }, { type: "text", text: SENTINEL });
tb.system = sysBlocks;
normalizeCacheControlTtl(tb);
Expand Down Expand Up @@ -1270,29 +1264,14 @@ export class BaseExecutor {
"X-Claude-Code-Session-Id": sessionId,
};

// Drop case variants of the same header name before merging — undici
// would otherwise concatenate them (issue #1454).
const ccKeysLower = new Set(Object.keys(ccHeaders).map((k) => k.toLowerCase()));
for (const key of Object.keys(headers)) {
if (ccKeysLower.has(key.toLowerCase())) delete headers[key];
}
Object.assign(headers, ccHeaders);
mergeCcHeaders(headers, ccHeaders);
if (usesCcWireImage(this.provider) && usesClaudeCodeProtocol) {
delete headers["Authorization"];
headers["x-api-key"] =
activeCredentials?.apiKey || activeCredentials?.accessToken || "";
}
delete headers["X-Stainless-Helper-Method"];

// OS/arch follow the host running the signed binary. Runtime version
// is pinned to the captured CLI wire image, not OmniRoute's Node.
headers["X-Stainless-Arch"] = stainlessArch();
headers["X-Stainless-Lang"] = "js";
headers["X-Stainless-OS"] = stainlessOS();
headers["X-Stainless-Runtime"] = "node";
headers["X-Stainless-Runtime-Version"] = CLAUDE_CLI_STAINLESS_RUNTIME_VERSION;
headers["X-Stainless-Retry-Count"] = "0";
delete headers["X-Stainless-Os"];
applyStainlessHeaders(headers, { arch: stainlessArch(), os: stainlessOS() });
}
// selectBetaFlags() above always includes redact-thinking for an
// "opaque" client (no client-negotiated anthropic-beta) — correct
Expand Down Expand Up @@ -1424,6 +1403,8 @@ export class BaseExecutor {
// Enforce peer tracing after all configurable headers have been merged so
// operator/provider metadata cannot accidentally erase the loop guard.
applyPeerTraceHeader(finalHeaders, clientHeaders, url);
// Rides `anthropic-usage-limit: slow` once this account accepted the offer.
const claudeSentSlow = claudeUsageLimit.applyHeader(finalHeaders, activeCredentials);
const serializedBody = prl.parseBody(bodyString);
// #4307 — Preserve the non-enumerable tool-name cloak/remap reverse map
// (`_toolNameMap`, set on the live `transformedBody` by
Expand Down Expand Up @@ -1672,6 +1653,21 @@ export class BaseExecutor {
}
}

// Claude OAuth usage wall: accept the slow-lane offer / claim the weekly
// session-limit reset and retry the SAME account instead of surfacing the 429
// (which would cool the connection down). Runs AFTER every 400-driven retry
// above so it classifies the FINAL response of this attempt.
const claudeRetry = await claudeUsageLimit.shouldRetry(response, url, {
credentials: activeCredentials,
signal,
budgetMs: fetchStartTimeoutMs,
sentSlow: claudeSentSlow,
});
if (claudeRetry) {
urlIndex--; // re-run this urlIndex (header injection sees the new lane state)
continue;
}

// Intra-URL retry: agentrouter.org WAF returns 400 content-blocked
// intermittently (burst-sensitive, recovers after cooldown). Retry the
// same URL with exponential backoff before falling through to the
Expand Down
18 changes: 18 additions & 0 deletions open-sse/executors/claudeIdentity.ts
Original file line number Diff line number Diff line change
Expand Up @@ -469,3 +469,21 @@ export function stripProxyToolPrefix(body: Record<string, unknown>): void {
}
}
}

/**
* Drop any previously injected billing header / Claude Code sentinel block from a
* `system` array, so re-prepending them on a retry stays idempotent instead of stacking
* (issue #1712 — stacking breaks prompt-cache prefix matching). Mutates `sysBlocks`.
*/
export function stripClaudeSystemPrefixBlocks(
sysBlocks: Array<Record<string, unknown>>,
sentinel: string
): void {
for (let i = sysBlocks.length - 1; i >= 0; i--) {
const text = sysBlocks[i]?.text;
if (typeof text !== "string") continue;
if (text.startsWith("x-anthropic-billing-header:") || text.startsWith(sentinel)) {
sysBlocks.splice(i, 1);
}
}
}
146 changes: 146 additions & 0 deletions open-sse/executors/claudeUsageLimit.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
/**
* claudeUsageLimit.ts — executor-side glue for the Claude OAuth usage wall.
*
* Keeps `BaseExecutor.execute()` free of the lower-priority lane's mechanics: this guard
* owns the per-request wait accounting, the header injection, the abort-aware sleep and the
* decision logging. The decision itself is the pure state machine in
* `open-sse/services/claudeLowPriority.ts`; the weekly session-limit reset lives in
* `open-sse/services/claudeLimitReset.ts`.
*
* Both behaviors are opt-in per connection (`providerSpecificData.lowPriorityMode` /
* `autoLimitReset`) and only ever act on a real 5-hour usage wall — see the module header of
* claudeLowPriority.ts for the wire contract and the lifecycle.
*/

import { attemptClaudeLimitReset } from "../services/claudeLimitReset.ts";
import {
CLAUDE_USAGE_LIMIT_HEADER,
CLAUDE_USAGE_LIMIT_SLOW,
createClaudeLowPriorityWait,
handleClaudeUsageLimitResponse,
isClaudeLowPriorityActive,
readClaudeUsageLimitConfig,
resolveClaudeUsageLimitKey,
type ClaudeLowPriorityWait,
} from "../services/claudeLowPriority.ts";
import type { ExecutorLog, ProviderCredentials } from "./base.ts";

/** Safety margin kept between the last lane wait and the request's own upstream timeout. */
export const CLAUDE_USAGE_LIMIT_WAIT_MARGIN_MS = 5_000;

type GuardResponse = { status: number; headers: Headers };

export type ClaudeUsageLimitRetryInput = {
credentials?: ProviderCredentials | null;
signal?: AbortSignal | null;
/** The request's upstream-start timeout; 0/undefined means "unbounded". */
budgetMs?: number;
/** Whether THIS request went out carrying the slow header (see `sentSlow` in the service). */
sentSlow: boolean;
};

/** True for a native Claude connection authenticated with a subscription OAuth token. */
function isClaudeOAuth(provider: string, credentials?: ProviderCredentials | null): boolean {
return (
provider === "claude" &&
typeof credentials?.accessToken === "string" &&
credentials.accessToken.startsWith("sk-ant-oat") &&
!credentials?.apiKey
);
}

/** Sleep that rejects as soon as the request is aborted, so a lane wait never outlives it. */
function abortableSleep(delayMs: number, signal?: AbortSignal | null): Promise<void> {
return new Promise<void>((resolve, reject) => {
const timer = setTimeout(() => {
signal?.removeEventListener("abort", onAbort);
resolve();
}, delayMs);
const onAbort = () => {
clearTimeout(timer);
reject(signal?.reason ?? new DOMException("Aborted", "AbortError"));
};
if (signal?.aborted) return onAbort();
signal?.addEventListener("abort", onAbort, { once: true });
});
}

/** One instance per `execute()` call: the wait window spans the intra-URL retries. */
export class ClaudeUsageLimitGuard {
private readonly wait: ClaudeLowPriorityWait = createClaudeLowPriorityWait();
private readonly startedAtMs = Date.now();
private key: string | null = null;

constructor(
private readonly provider: string,
private readonly log?: ExecutorLog | null
) {}

/**
* Stamp `anthropic-usage-limit: slow` when this connection's lane is active. Returns
* whether the header went out, which the response side needs to tell a lane verdict from
* a header-less sibling's outcome.
*/
applyHeader(headers: Record<string, string>, credentials?: ProviderCredentials | null): boolean {
this.key = isClaudeOAuth(this.provider, credentials)
? resolveClaudeUsageLimitKey(credentials ?? {})
: null;
if (this.key === null || !isClaudeLowPriorityActive(this.key)) return false;
headers[CLAUDE_USAGE_LIMIT_HEADER] = CLAUDE_USAGE_LIMIT_SLOW;
return true;
}

/**
* Classify an upstream response. Resolves true when the caller must retry the SAME
* account (the sleep, if any, has already happened) instead of surfacing the response.
*/
async shouldRetry(
response: GuardResponse,
url: string,
input: ClaudeUsageLimitRetryInput
): Promise<boolean> {
if (this.key === null) return false;
const key = this.key;
const credentials = input.credentials;
const decision = await handleClaudeUsageLimitResponse({
key,
config: readClaudeUsageLimitConfig(credentials?.providerSpecificData),
response,
wait: this.wait,
sentSlow: input.sentSlow,
waitCeilingMs: this.waitCeilingMs(input.budgetMs),
claimLimitReset: () =>
attemptClaudeLimitReset({
key,
accessToken: credentials?.accessToken ?? "",
providerSpecificData: credentials?.providerSpecificData,
log: this.log,
}).then((attempt) => attempt.reset),
});

if (decision.kind === "ended") {
this.log?.info?.("CLAUDE_LOW_PRIORITY", `lane ended (${decision.reason}) on ${url}`);
return false;
}
if (decision.kind !== "retry") return false;

this.log?.info?.(
"CLAUDE_LOW_PRIORITY",
`${decision.via} on ${url} — retrying same account in ${decision.delayMs}ms`
);
if (decision.delayMs > 0) await abortableSleep(decision.delayMs, input.signal);
return true;
}

/**
* What is left of the request's upstream timeout, minus a safety margin. Without this the
* server-announced max-wait (20 min by default, up to 6 h) outlives the request and the
* sleep is aborted mid-wait, surfacing a TimeoutError instead of the graceful `max_wait`
* end plus its cool-off.
*/
private waitCeilingMs(budgetMs?: number): number | undefined {
if (!budgetMs || budgetMs <= 0) return undefined;
const elapsed = Date.now() - this.startedAtMs;
return Math.max(0, budgetMs - elapsed - CLAUDE_USAGE_LIMIT_WAIT_MARGIN_MS);
}
}
Loading
Loading