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
8 changes: 8 additions & 0 deletions open-sse/services/quotaPreflight.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import { isCompatibleProviderConnectionId } from "@/shared/utils/compatibleProviderId";
import { isFeatureFlagEnabled } from "@/shared/utils/featureFlags";
import { isClaudeExtraUsageAllowed } from "@/lib/providers/claudeExtraUsage";
import { isQuotaHealthy } from "@/domain/quotaCache";
import { fetchNewApiAggregatorQuota } from "./newApiAggregatorQuotaFetcher.ts";
import {
isAntigravityQuotaProvider,
Expand All @@ -39,6 +40,8 @@ export interface QuotaCutoffScope {
provider?: string | null;
requestedModel?: string | null;
providerSpecificData?: unknown;
// #14359 — recent successful dispatch stands the cutoff down for this connection.
connectionId?: string | null;
}

export interface QuotaWindowInfo {
Expand Down Expand Up @@ -323,6 +326,10 @@ export function evaluateQuotaCutoff(
if (isClaudeExtraUsageAllowed(scope?.provider, scope?.providerSpecificData)) {
return { proceed: true, quotaPercent: quota.percentUsed };
}
// #14359 — same escape as the dispatch-time predicates: a recent success is not exhaustion.
if (scope?.connectionId && isQuotaHealthy(scope.connectionId)) {
return { proceed: true, quotaPercent: quota.percentUsed };
}

const windows = quota.windows;
if (windows && Object.keys(windows).length > 0) {
Expand Down Expand Up @@ -401,6 +408,7 @@ export async function preflightQuota(
provider,
requestedModel,
providerSpecificData: connection.providerSpecificData,
connectionId,
};
const windows = quota.windows;
if (windows && Object.keys(windows).length > 0) {
Expand Down
61 changes: 58 additions & 3 deletions src/domain/quotaCache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,9 @@ const EXHAUSTED_REFRESH_MS = 5 * 60 * 1000; // 5 minutes: recheck exhausted acco
const REFRESH_INTERVAL_MS = 60 * 1000; // Background tick every 1 minute
export const DEFAULT_QUOTA_THRESHOLD_PERCENT = 99;

// #14359 — a park is trusted at most EXHAUSTED_MAX_PARK_MS past its observation; a far weekly-window reset must not block for days.
export const EXHAUSTED_MAX_PARK_MS = 30 * 60 * 1000;

// ─── State ──────────────────────────────────────────────────────────────────
//
// #8065 — Next.js `output: "standalone"` builds can load this module from
Expand All @@ -105,6 +108,8 @@ export const DEFAULT_QUOTA_THRESHOLD_PERCENT = 99;
interface QuotaCacheState {
cache: Map<string, QuotaCacheEntry>;
refreshingSet: Set<string>;
// #14359 — connectionId → epoch ms the healthy override stands the predicates down; shared state per #8065.
healthyUntil: Map<string, number>;
refreshTimer: ReturnType<typeof setInterval> | null;
tickRunning: boolean;
}
Expand All @@ -118,6 +123,7 @@ function getState(): QuotaCacheState {
globalThis.__omnirouteQuotaCacheState = {
cache: new Map(),
refreshingSet: new Set(),
healthyUntil: new Map(),
refreshTimer: null,
tickRunning: false,
};
Expand All @@ -127,6 +133,47 @@ function getState(): QuotaCacheState {

const MAX_CONCURRENT_REFRESHES = 5;

// ─── #14359 park window + healthy override ──────────────────────────────────

// #14359 — cap the park deadline at `anchorMs + EXHAUSTED_MAX_PARK_MS`; unparseable passes through (fixed EXHAUSTED_TTL applies).
function capParkWindow(resetAt: string | null, anchorMs: number): string | null {
if (!resetAt) return resetAt;
const ms = parseDate(resetAt);
if (ms === null) return resetAt;
if (ms > anchorMs + EXHAUSTED_MAX_PARK_MS) {
return new Date(anchorMs + EXHAUSTED_MAX_PARK_MS).toISOString();
}
return resetAt;
}

// #14359 — keep the prior deadline while the exhausted streak continues, so the quota monitor's rewrites cannot re-anchor the park.
function preserveParkDeadline(
resetAt: string | null,
prior: QuotaCacheEntry | null | undefined
): string | null {
if (prior?.exhausted && prior.nextResetAt) {
const priorMs = parseDate(prior.nextResetAt);
if (priorMs !== null && priorMs > Date.now()) return prior.nextResetAt;
}
return capParkWindow(resetAt, Date.now());
}

// #14359 — arm the healthy override for one park window (chat success hook).
export function markQuotaHealthy(connectionId: string): void {
getState().healthyUntil.set(connectionId, Date.now() + EXHAUSTED_MAX_PARK_MS);
}

// #14359 — clear the healthy override (a genuine 429 re-parks immediately).
export function unmarkQuotaHealthy(connectionId: string): void {
getState().healthyUntil.delete(connectionId);
}

// #14359 — true while the healthy override for this connection is armed.
export function isQuotaHealthy(connectionId: string): boolean {
const until = getState().healthyUntil.get(connectionId);
return until !== undefined && until > Date.now();
}

// ─── Helpers ────────────────────────────────────────────────────────────────

function isExhausted(quotas: Record<string, QuotaInfo>): boolean {
Expand Down Expand Up @@ -394,7 +441,7 @@ export function hydrateCodexQuotaCacheForRequest(
}
}
entry.exhausted = isExhausted(entry.quotas);
if (exhaustedResetAt) entry.nextResetAt = exhaustedResetAt;
if (exhaustedResetAt) entry.nextResetAt = capParkWindow(exhaustedResetAt, Date.now());
cache.set(connection.id, entry);
}

Expand Down Expand Up @@ -431,6 +478,8 @@ export function isQuotaExhaustedForRequest(
requestedModel: string | null = null,
providerSpecificData?: unknown
): boolean {
// #14359 — a recent successful dispatch stands the predicates down for a park window.
if (isQuotaHealthy(connectionId)) return false;
if (isClaudeExtraUsageAllowed(provider, providerSpecificData)) return false;
const entry = getState().cache.get(connectionId) || hydrateQuotaCacheFromSnapshots(connectionId);
if (!entry) return false;
Expand Down Expand Up @@ -486,7 +535,8 @@ export function setQuotaCache(
quotas,
fetchedAt: Date.now(),
exhausted,
nextResetAt: exhausted ? earliestResetAt(quotas) : null,
// #14359 — cap the park; keep the prior deadline while the streak continues (monitor rewrites must not extend it).
nextResetAt: exhausted ? preserveParkDeadline(earliestResetAt(quotas), prior) : null,
};
getState().cache.set(connectionId, entry);

Expand Down Expand Up @@ -618,7 +668,8 @@ function hydrateQuotaCacheFromSnapshots(connectionId: string): QuotaCacheEntry |
quotas,
fetchedAt: fetchedAt || Date.now(),
exhausted,
nextResetAt: exhausted ? earliestResetAt(quotas) : null,
// #14359 — cap the hydrated park relative to the snapshot's observation time.
nextResetAt: exhausted ? capParkWindow(earliestResetAt(quotas), fetchedAt || Date.now()) : null,
windowDurationMs,
};
cache.set(connectionId, entry);
Expand All @@ -630,6 +681,8 @@ function hydrateQuotaCacheFromSnapshots(connectionId: string): QuotaCacheEntry |
* Returns false if no cache entry exists (unknown = assume available).
*/
export function isAccountQuotaExhausted(connectionId: string): boolean {
// #14359 — mirror of the request-time predicate: honour the healthy override.
if (isQuotaHealthy(connectionId)) return false;
const entry = getState().cache.get(connectionId) || hydrateQuotaCacheFromSnapshots(connectionId);
if (!entry) return false;
if (!entry.exhausted) return false;
Expand Down Expand Up @@ -767,6 +820,8 @@ export function getQuotaSnapshotFetchedAt(connectionId: string): number | null {
* Uses 5-minute fixed TTL since we don't know the actual resetAt.
*/
export function markAccountExhaustedFrom429(connectionId: string, provider: string) {
// #14359 — a real upstream 429 is authoritative: drop any healthy override.
unmarkQuotaHealthy(connectionId);
getState().cache.set(connectionId, {
connectionId,
provider,
Expand Down
4 changes: 3 additions & 1 deletion src/sse/handlers/chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ import { isFeatureFlagEnabled, isRotationAttributionEnabled } from "@/shared/uti
import * as agyLease from "../services/antigravityLeaseLifecycle";
import { shouldIsolateProbeFailures } from "@/shared/utils/probeOrigin";
import { getCircuitBreaker, isLocalStreamLifecycleError } from "../../shared/utils/circuitBreaker";
import { markAccountExhaustedFrom429 } from "../../domain/quotaCache";
import { markAccountExhaustedFrom429, markQuotaHealthy } from "../../domain/quotaCache";
import { resolveForcedConnectionForCredentialPool } from "../services/sessionAffinityPin.ts";
import { RequestTelemetry, recordTelemetry } from "../../shared/utils/requestTelemetry";
import { generateRequestId } from "../../shared/utils/requestId";
Expand Down Expand Up @@ -2075,6 +2075,8 @@ async function handleSingleModelChat(

if (result.success) {
clearModelLock(provider, credentials.connectionId, model);
// #14359 — a real upstream success is authoritative: arm the healthy override.
markQuotaHealthy(credentials.connectionId);
// #12254: exactly-once breaker accounting — combo successes are recorded by
// combo.ts (recordProviderSuccess); live combo tests never touch the breaker.
if (classifyProviderBreakerResult(result, isCombo, forceLiveComboTest) === "success") {
Expand Down
2 changes: 1 addition & 1 deletion src/sse/services/auth.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1902,7 +1902,7 @@ export async function getProviderCredentials(
allRateLimited: true,
retryAfter,
retryAfterHuman: formatRetryAfter(retryAfter),
lastError: `All ${provider} accounts have exhausted their quota`,
lastError: `All ${provider} accounts have exhausted their quota (cached quota state, no upstream attempt; earliest reset ${formatRetryAfter(retryAfter)})`,
lastErrorCode: 429,
};
}
Expand Down
88 changes: 88 additions & 0 deletions tests/unit/quota-healthy-override-14359.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
// #14359 — a successful upstream dispatch must arm a healthy override so the
// quota predicates stand down for the park window, and a genuine upstream 429
// must clear it again. This is the live-proven half of the fix: a connection
// that is actually serving traffic is, by definition, not quota-exhausted.
import test from "node:test";
import assert from "node:assert/strict";
import fs from "node:fs";
import os from "node:os";
import path from "node:path";

const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omni-quota-healthy-14359-"));
process.env.DATA_DIR = TEST_DATA_DIR;

// Dynamic imports are required: DATA_DIR must be set before the db modules
// open SQLite (same pattern as tests/unit/quota-cache-hydrate-5015.test.ts).
const coreDb = await import("../../src/lib/db/core.ts");
const quotaCache = await import("../../src/domain/quotaCache.ts");

// The park window under test (30 minutes — mirrors EXHAUSTED_MAX_PARK_MS).
const MAX_PARK_MS = 30 * 60 * 1000;

test.after(() => {
coreDb.resetDbInstance();
fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 });
});

test("#14359 markQuotaHealthy stands both predicates down on an exhausted entry", () => {
quotaCache.__clearForTests();
quotaCache.markAccountExhaustedFrom429("ho-14359", "qwen-cloud-token-plan");
assert.equal(quotaCache.isAccountQuotaExhausted("ho-14359"), true);
assert.equal(quotaCache.isQuotaExhaustedForRequest("ho-14359", "qwen-cloud-token-plan"), true);

// A successful dispatch arms the override (chat.ts success hook).
quotaCache.markQuotaHealthy("ho-14359");
assert.equal(
quotaCache.isAccountQuotaExhausted("ho-14359"),
false,
"healthy override must stand isAccountQuotaExhausted down"
);
assert.equal(
quotaCache.isQuotaExhaustedForRequest("ho-14359", "qwen-cloud-token-plan"),
false,
"healthy override must stand isQuotaExhaustedForRequest down"
);
});

test("#14359 a genuine 429 clears the healthy override", () => {
quotaCache.__clearForTests();
quotaCache.markQuotaHealthy("ho-429-14359");
assert.equal(quotaCache.isAccountQuotaExhausted("ho-429-14359"), false);

quotaCache.markAccountExhaustedFrom429("ho-429-14359", "zai");
assert.equal(
quotaCache.isAccountQuotaExhausted("ho-429-14359"),
true,
"a real 429 must clear the override and re-park immediately"
);
});

test("#14359 the healthy override expires after the park window", () => {
quotaCache.__clearForTests();
quotaCache.markAccountExhaustedFrom429("ho-exp-14359", "qwen-cloud-token-plan");
quotaCache.markQuotaHealthy("ho-exp-14359");
assert.equal(quotaCache.isAccountQuotaExhausted("ho-exp-14359"), false);

// The healthy map is a deliberate globalThis anchor (#8065 — chunk-shared
// singleton); reach it via a named reference so the test can simulate time.
interface QuotaCacheTestState {
healthyUntil: Map<string, number>;
}
const state: QuotaCacheTestState | undefined = (
globalThis as {
__omnirouteQuotaCacheState?: QuotaCacheTestState;
}
).__omnirouteQuotaCacheState;
assert.ok(state?.healthyUntil.has("ho-exp-14359"), "override armed in shared state");
assert.ok(
(state!.healthyUntil.get("ho-exp-14359") ?? 0) - Date.now() <= MAX_PARK_MS,
"override lifetime must equal the park window"
);
state!.healthyUntil.set("ho-exp-14359", Date.now() - 1);

assert.equal(
quotaCache.isAccountQuotaExhausted("ho-exp-14359"),
true,
"an expired override must not stand the predicates down any more"
);
});
Loading