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
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(combo):** quota-weighted routing stops drawing on an out-of-credit connection — a 402 now invalidates the stored quota snapshot instead of leaving its stale remaining percentage in place, and a snapshot older than 10 minutes no longer counts as confident headroom for the primary pool ([#12972](https://github.com/diegosouzapw/OmniRoute/pull/12972)) — thanks @HouMinXi
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **fix(combo):** quota-aware expansion drops banned, inactive, missing, and wrong-provider connections before quota fetch or model dispatch; pins and allowlists stay selectors, not a bypass. Antigravity automatic exhaustion now requires a reported zero remaining, so a positive balance below 1% stays eligible.
7 changes: 7 additions & 0 deletions open-sse/services/combo/executeTargetAttempt.ts
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ import {
isQuotaExhaustionResponse,
recordQuotaExhaustionClassification,
} from "./quotaExhaustion.ts";
import { markAccountExhaustedFromCredits } from "../../../src/domain/quotaCache.ts";
import { classifyComboOutcome, redactConnectionLabel } from "./comboErrorAggregation.ts";
import { readConnectionForCooldownGate } from "./executeTargetGates.ts";
import {
Expand Down Expand Up @@ -994,6 +995,12 @@ export async function executeTargetAttempt(opts: {

const quotaExhausted = await isQuotaExhaustionResponse(result, provider, rawModel, profile);
recordQuotaExhaustionClassification(result, quotaExhausted);
// Balance exhaustion is upstream truth about credits, and it outranks the
// stored snapshot — which can be hours stale and still claim headroom. Mark
// it so the next quota-weighted draw stops picking this connection.
if (quotaExhausted && result.status === 402 && targetWithConnection.connectionId && provider) {
markAccountExhaustedFromCredits(targetWithConnection.connectionId, provider);
}
state.observeFailure(quotaExhausted, target.executionKey);

// Check if this is a transient error worth retrying on same model.
Expand Down
94 changes: 56 additions & 38 deletions open-sse/services/combo/quotaStrategies.ts
Original file line number Diff line number Diff line change
@@ -1,18 +1,17 @@
/**
* Stateful + async reset-aware / reset-window quota strategies for combo routing.
*
* Holds the two mutable module-level caches that back reset-aware routing
* (`resetAwareConnectionCache` for per-provider active connections and
* `resetAwareQuotaCache` for per-connection quota snapshots), plus the helpers
* Holds the per-connection quota snapshot cache and helpers
* that read/write them and the strategy orderers. Extracted byte-identically
* from combo.ts (QG v2 Fase 9 T5 D7b) — the larger, stateful half of the
* reset-aware quota block. The pure scoring/window-math half lives in
* ./quotaScoring.ts and is imported here.
*
* State cohesion: `resetAwareConnectionCache`, `resetAwareQuotaCache`, and
* State cohesion: `resetAwareQuotaCache` and
* `MAX_RESET_AWARE_CACHE` MUST remain single instances defined once here,
* alongside their only readers/writers (getQuotaAwareConnectionsForTarget,
* fetchResetAwareQuotaWithCache) — never duplicate a Map.
* alongside their only readers/writers (`fetchResetAwareQuotaWithCache`).
* Connection lists go through `getCachedProviderConnections` (5s TTL,
* invalidated on connection writes). Do not add a second connection cache.
*
* Cross-module state: the tie-band round-robin in orderTargetsByResetAwareQuota
* and orderTargetsByResetWindow shares the same rrCounters Map from ./rrState.ts
Expand Down Expand Up @@ -50,18 +49,28 @@ import { rankByHeadroom, type HeadroomSaturation } from "./headroomRanking.ts";
import { getInflight, incrementInflight } from "./quotaShareInflight.ts";
import { preferAntigravityConnectionsWithStoredProject } from "../antigravityProjectPersist.ts";
import { getQuotaFetchScope } from "../antigravityQuotaFamily.ts";
import { isQuotaExhaustedForRequest } from "../../../src/domain/quotaCache.ts";
import {
getQuotaSnapshotFetchedAt,
getQuotaWeightedRemainingPercent,
isQuotaExhaustedForRequest,
} from "../../../src/domain/quotaCache.ts";

/**
* How long a stored quota snapshot stays good enough to be counted as confident
* headroom by the quota-weighted A pool.
*
* Matches the background refresh cadence for active accounts (quotaCache's
* ACTIVE_TTL_MS), doubled to absorb one missed refresh tick. Past that the
* snapshot says "unknown", not "empty": the connection drops to the B pool and
* is still routed to when nothing fresher has room.
*/
export const QUOTA_WEIGHTED_MAX_SNAPSHOT_AGE_MS = 10 * 60 * 1000;

const RESET_AWARE_CONNECTION_CACHE_TTL_MS = 30_000;
const RESET_AWARE_QUOTA_FETCH_CONCURRENCY = 5;
const HEADROOM_SATURATION_FETCH_CONCURRENCY = 5;

const MAX_RESET_AWARE_CACHE = 200;

const resetAwareConnectionCache = new Map<
string,
{ fetchedAt: number; connections: Array<Record<string, unknown>> }
>();
const resetAwareQuotaCache = new Map<
string,
{ fetchedAt: number; quota: unknown; refreshPromise: Promise<unknown> | null }
Expand All @@ -77,35 +86,24 @@ async function getQuotaAwareConnectionsForTarget(
const provider = getResetAwareProvider(target);
if (!provider || !getQuotaFetcher(provider)) return [];
if (!connectionCache.has(provider)) {
const cached = resetAwareConnectionCache.get(provider);
if (cached && Date.now() - cached.fetchedAt < RESET_AWARE_CONNECTION_CACHE_TTL_MS) {
connectionCache.set(provider, cached.connections);
return cached.connections;
}

if (!connectionLoadPromises.has(provider)) {
connectionLoadPromises.set(
provider,
(async () => {
try {
const connections = await getCachedProviderConnections({ provider, isActive: true });
let activeConnections = Array.isArray(connections)
? (connections as Array<Record<string, unknown>>)
? (connections as Array<Record<string, unknown>>).filter(
(connection) =>
connection.isActive !== false &&
String(connection.testStatus || "")
.trim()
.toLowerCase() !== "banned"
)
: [];
if (provider === "antigravity" || provider === "agy") {
activeConnections = preferAntigravityConnectionsWithStoredProject(activeConnections);
}
if (
!resetAwareConnectionCache.has(provider) &&
resetAwareConnectionCache.size >= MAX_RESET_AWARE_CACHE
) {
const oldest = resetAwareConnectionCache.keys().next().value;
if (oldest !== undefined) resetAwareConnectionCache.delete(oldest);
}
resetAwareConnectionCache.set(provider, {
connections: activeConnections,
fetchedAt: Date.now(),
});
return activeConnections;
} catch (error) {
log.warn?.("COMBO", "Reset-aware failed to load quota-aware connections.", {
Expand Down Expand Up @@ -212,6 +210,8 @@ export async function expandTargetsByQuotaAwareConnections(
apiKeyAllowedConnectionIds
);
if (connectionIds.length === 0) {
const provider = getResetAwareProvider(target);
if (provider && getQuotaFetcher(provider)) continue;
if (
unrestrictedConnectionIds.length > 0 &&
normalizeConnectionIds(apiKeyAllowedConnectionIds)
Expand All @@ -225,6 +225,7 @@ export async function expandTargetsByQuotaAwareConnections(
for (const connectionId of connectionIds) {
const provider = getResetAwareProvider(target);
const connection = connectionById.get(connectionId);
if (provider && getQuotaFetcher(provider) && connection?.provider !== provider) continue;
if (
connection &&
typeof connection.rateLimitedUntil === "string" &&
Expand Down Expand Up @@ -750,14 +751,15 @@ function sortByScoreThenIndex(a: QuotaWeightedScored, b: QuotaWeightedScored): n
return a.index - b.index;
}

function resolveQuotaWeightedFloor(configSource: Record<string, unknown> | null | undefined): number {
function resolveQuotaWeightedFloor(
configSource: Record<string, unknown> | null | undefined
): number {
// Number(null) and Number("") are both 0, so an unset or blank key would
// switch the floor off instead of taking the default. Only a value that is
// actually a number, or a non-empty numeric string, gets to move it.
const configured = configSource?.quotaWeightedFloorPercent;
const raw =
typeof configured === "number" ||
(typeof configured === "string" && configured.trim() !== "")
typeof configured === "number" || (typeof configured === "string" && configured.trim() !== "")
? Number(configured)
: Number.NaN;
return Number.isFinite(raw) ? Math.max(0, Math.min(100, raw)) : 1;
Expand Down Expand Up @@ -798,12 +800,28 @@ export async function orderTargetsByQuotaWeighted(
}),
});

const eligible = scoredTargets.filter((entry) => entry.remainingPercent > 0);
// The live snapshot outranks the freshly-scored fetch on two counts: a 402
// recorded against this connection zeroes it, and an observation older than
// the staleness bound is not confident enough to sit in the A pool.
const now = Date.now();
const withSnapshot = scoredTargets.map((entry) => {
const connectionId = entry.target.connectionId ?? "";
const marked = connectionId ? getQuotaWeightedRemainingPercent(connectionId) : null;
const fetchedAt = connectionId ? getQuotaSnapshotFetchedAt(connectionId) : null;
return {
...entry,
remainingPercent: marked === 0 ? 0 : entry.remainingPercent,
stale: fetchedAt !== null && now - fetchedAt > QUOTA_WEIGHTED_MAX_SNAPSHOT_AGE_MS,
};
});

const eligible = withSnapshot.filter((entry) => entry.remainingPercent > 0);
const floor = resolveQuotaWeightedFloor(configSource);
const poolA =
floor === 0 ? eligible : eligible.filter((entry) => entry.remainingPercent > floor);
const poolB =
floor === 0 ? [] : eligible.filter((entry) => entry.remainingPercent > 0 && entry.remainingPercent <= floor);
const hasRoom = (entry: (typeof eligible)[number]) =>
floor === 0 ? true : entry.remainingPercent > floor;
// A holds only connections we both believe have room AND observed recently.
const poolA = eligible.filter((entry) => hasRoom(entry) && !entry.stale);
const poolB = eligible.filter((entry) => !hasRoom(entry) || entry.stale);
const selected = poolA.length > 0 ? poolA : poolB;
if (selected.length === 0) return [];

Expand Down
43 changes: 41 additions & 2 deletions src/domain/quotaCache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -304,8 +304,8 @@ function isAntigravityQuotaExhausted(
matchingWindows.length > 0 &&
matchingWindows.every(
(windowName) =>
getQuotaWindowStatus(connectionId, windowName, DEFAULT_QUOTA_THRESHOLD_PERCENT)
?.reachedThreshold
// Automatic exhaustion is not the operator's optional usage cutoff.
getQuotaWindowStatus(connectionId, windowName, 100)?.reachedThreshold
)
);
}
Expand Down Expand Up @@ -683,6 +683,45 @@ export function getQuotaWindowObservation(
};
}

/**
* Mark an account as out of credits from a 402-class response.
*
* Upstream refusing the request for balance is authoritative: it outranks
* whatever remaining percentage the last snapshot happened to hold, which may
* be hours old. Without this, a connection that answered 402 keeps its stale
* non-zero remaining and the next quota-weighted draw can pick it again.
*
* The entry is kept (never deactivated or deleted) — credits come back, and a
* later successful refresh or window reset clears the flag through the same
* paths that clear a 429 mark.
*/
export function markAccountExhaustedFromCredits(connectionId: string, provider: string) {
markAccountExhaustedFrom429(connectionId, provider);
}

/**
* Remaining headroom the quota-weighted strategy should credit this connection
* with, as a percentage. Returns 0 once the connection is known exhausted so a
* 402-marked account cannot be weighted back into the draw.
*/
export function getQuotaWeightedRemainingPercent(connectionId: string): number | null {
const entry = getState().cache.get(connectionId) || hydrateQuotaCacheFromSnapshots(connectionId);
if (!entry) return null;
if (isAccountQuotaExhausted(connectionId)) return 0;

const remaining = Object.values(entry.quotas)
.filter((quota) => quota.fractionReported !== false)
.map((quota) => clampPercent(quota.remainingPercentage));
if (remaining.length === 0) return null;
return Math.min(...remaining);
}

/** Epoch-ms of the observation backing this connection's snapshot, if any. */
export function getQuotaSnapshotFetchedAt(connectionId: string): number | null {
const entry = getState().cache.get(connectionId) || hydrateQuotaCacheFromSnapshots(connectionId);
return entry ? entry.fetchedAt : null;
}

/**
* Mark an account as quota-exhausted from a 429 response (no quota data available).
* Uses 5-minute fixed TTL since we don't know the actual resetAt.
Expand Down
3 changes: 3 additions & 0 deletions stryker.conf.json
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@
"tests/unit/alibaba-free-tier-exhaustion.test.ts",
"tests/unit/anthropic-thinking-signature-recovery.test.ts",
"tests/unit/agy-family-not-connection-cooldown.test.ts",
"tests/unit/agy-quota-exhaustion-threshold.test.ts",
"tests/unit/antigravity-429-quota-cooldown.test.ts",
"tests/unit/antigravity-429-quota-tdd.test.ts",
"tests/unit/antigravity-prefer-stored-project.test.ts",
Expand Down Expand Up @@ -230,6 +231,8 @@
"tests/unit/combo/combo-exhausted-skip.test.ts",
"tests/unit/combo/combo-target-timeout-standards.test.ts",
"tests/unit/combo/effective-max-concurrency.test.ts",
"tests/unit/combo/quota-connection-eligibility.test.ts",
"tests/unit/combo/quota-weighted-stale-402.test.ts",
"tests/unit/combo/quota-weighted-strategy.test.ts",
"tests/unit/combo/recovery-hint.test.ts",
"tests/unit/combo/reset-window-strategy-9330.test.ts",
Expand Down
83 changes: 83 additions & 0 deletions tests/unit/agy-quota-exhaustion-threshold.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
import test, { after, beforeEach } 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 previousDataDir = process.env.DATA_DIR;
const dataDir = fs.mkdtempSync(path.join(os.tmpdir(), "agy-quota-threshold-"));
process.env.DATA_DIR = dataDir;

const core = await import("../../src/lib/db/core.ts");
const cache = await import("../../src/domain/quotaCache.ts");
const { evaluateQuotaLimitPolicy } = await import("../../src/sse/services/auth.ts");
const { toProviderConnection } = await import("../../src/lib/db/providers/lazyConnectionView.ts");

beforeEach(() => cache.__clearForTests());
after(() => {
cache.__clearForTests();
core.resetDbInstance();
if (previousDataDir === undefined) delete process.env.DATA_DIR;
else process.env.DATA_DIR = previousDataDir;
fs.rmSync(dataDir, { recursive: true, force: true });
});

function seed(provider: string, remaining: number, fractionReported = true) {
const resetAt = new Date(Date.now() + 86_400_000).toISOString();
cache.setQuotaCache("threshold-account", provider, {
"gemini-3.8-flash-high": { remainingPercentage: remaining, resetAt, fractionReported },
gemini_weekly: { remainingPercentage: remaining, resetAt, fractionReported },
"claude-opus-4-6-thinking": { remainingPercentage: 0, resetAt },
claude_gpt_weekly: { remainingPercentage: 0, resetAt },
});
}

for (const provider of ["agy", "antigravity"]) {
for (const remaining of [0.01, 0.94, 1, 1.01]) {
test(`${provider}: positive ${remaining}% is not automatic exhaustion`, () => {
seed(provider, remaining);
assert.equal(
cache.isQuotaExhaustedForRequest("threshold-account", provider, "gemini-3.8-flash-high"),
false
);
assert.equal(
cache.isQuotaExhaustedForRequest("threshold-account", provider, "claude-opus-4-6-thinking"),
true
);
});
}

test(`${provider}: reported zero remains exhausted`, () => {
seed(provider, 0);
assert.equal(
cache.isQuotaExhaustedForRequest("threshold-account", provider, "gemini-3.8-flash-high"),
true
);
});

test(`${provider}: unreported zero remains unknown`, () => {
seed(provider, 0, false);
assert.equal(
cache.isQuotaExhaustedForRequest("threshold-account", provider, "gemini-3.8-flash-high"),
false
);
});

test(`${provider}: explicit 99% usage policy still blocks low remaining quota`, () => {
seed(provider, 0.94);
const decision = evaluateQuotaLimitPolicy(
provider,
toProviderConnection({
id: "threshold-account",
provider,
isActive: true,
providerSpecificData: {
limitPolicy: { enabled: true, thresholdPercent: 99, windows: ["gemini_weekly"] },
},
}),
"gemini-3.8-flash-high"
);
assert.equal(decision.blocked, true);
assert.equal(decision.reasons.length, 1);
});
}
6 changes: 3 additions & 3 deletions tests/unit/antigravity-quota-skipping.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -137,7 +137,7 @@ test("isQuotaExhaustedForRequest scopes gemini exhaustion to the requested model
);
});

test("isQuotaExhaustedForRequest treats near-zero remaining as exhausted at default threshold", () => {
test("isQuotaExhaustedForRequest keeps reported positive remaining available", () => {
const connectionId = "conn-near-zero-test";
quotaCache.setQuotaCache(connectionId, "antigravity", {
"gemini-3.7-flash-medium": { remainingPercentage: 0.00000167, resetAt: null },
Expand All @@ -149,8 +149,8 @@ test("isQuotaExhaustedForRequest treats near-zero remaining as exhausted at defa
"antigravity",
"antigravity/gemini-3.7-flash-medium"
),
true,
"effectively-zero remaining should count as exhausted"
false,
"positive quota is not exhaustion; explicit usage cutoffs are evaluated separately"
);
});

Expand Down
Loading
Loading