import { getProviderConnectionById, getProviderConnections, updateProviderConnection, } from "@/lib/db/providers"; import { getSettings, resolveProxyForConnection, updateSettings } from "@/lib/db/settings"; import { getAllProviderLimitsCache, getProviderLimitsCache, setProviderLimitsCache, setProviderLimitsCacheBatch, type ProviderLimitsCacheEntry, } from "@/lib/db/providerLimits"; import { syncToCloud } from "@/lib/cloudSync"; import { setQuotaCache } from "@/domain/quotaCache"; import { buildClaudeExtraUsageConnectionUpdate, CLAUDE_EXTRA_USAGE_ERROR_SOURCE, isClaudeExtraUsageBlockEnabled, isClaudeExtraUsageQueued, } from "@/lib/providers/claudeExtraUsage"; import { isConnectionUnavailableToAuxiliaryActivity } from "@/lib/exclusiveLeaseIsolation"; import { clearRecoveredProviderState } from "@/sse/services/auth"; import { getMachineId } from "@/shared/utils/machine"; import { USAGE_SUPPORTED_PROVIDERS } from "@/shared/constants/providers"; import { mergeProviderLimitsCacheEntry, toProviderLimitsCacheEntry } from "./providerLimitsCache"; import { getExecutor } from "@omniroute/open-sse/executors/index.ts"; import { getUsageForProvider } from "@omniroute/open-sse/services/usage.ts"; import { rotationGroupFor, serializeRefresh, } from "@omniroute/open-sse/services/refreshSerializer.ts"; import { extractCodeAssistOnboardTierId, extractCodeAssistSubscriptionTier, } from "@omniroute/open-sse/services/codeAssistSubscription.ts"; import { extractAntigravityProjectIdFromPayload, getStoredAntigravityProjectId, } from "@omniroute/open-sse/services/antigravityProjectPersistence.ts"; import { runWithProxyContext } from "@omniroute/open-sse/utils/proxyFetch.ts"; import { onUsageRecorded } from "./usageEvents"; import { isRecord, isUsageQuotaKeyAllowed, normalizeUsageQuotasForProvider, sanitizeUsageQuotasForProvider, } from "./providerLimits/quotaNormalize"; import { syncInChunksWithSpacing } from "./providerLimits/chunkedSpacingSync"; type JsonRecord = Record; type SyncSource = "manual" | "scheduled"; interface ProviderConnectionLike { id: string; provider: string; authType?: string; accessToken?: string; refreshToken?: string; expiresAt?: string; tokenExpiresAt?: string; providerSpecificData?: JsonRecord; testStatus?: string; isActive?: boolean; lastError?: string | null; lastErrorAt?: string | null; lastErrorType?: string | null; lastErrorSource?: string | null; errorCode?: string | number | null; rateLimitedUntil?: string | null; backoffLevel?: number; } const PROVIDER_LIMITS_APIKEY_PROVIDERS = new Set([ "glm", "glm-cn", "zai", "glmt", "opencode-go", "ollama-cloud", "minimax", "minimax-cn", "crof", "nanogpt", "deepseek", "xiaomi-mimo", "vertex", "vertex-partner", "kimi-coding-apikey", "kiro", // Qoder connections are PAT-based (authType "apikey"); the usage fetcher // exchanges the PAT for a job token and reads openapi.qoder.sh/user/status. "qoder", "promptql", // PromptQL playground JWT → getCreditSummary USD credits "pql", // Adobe Firefly: web-cookie / JWT stored as apikey → credits/balance "adobe-firefly", "firefly", // HyperAgent session cookie → billing/usage creditBlocks "hyperagent", "ha", "firecrawl", // Command Code API key → /alpha/billing/credits + windowLimits "command-code", "conol-web", "cnl", // Alibaba Coding Plan (console API key) + Qwen personal Token Plan (console cookie) — #9603 "bailian-coding-plan", "qwen-cloud-token-plan", // AgentRouter (New-API) console System Access Token + New-Api-User id (providerSpecificData) "agentrouter", ]); const DEFAULT_PROVIDER_LIMITS_SYNC_INTERVAL_MINUTES = 70; const PROVIDER_LIMITS_AUTO_SYNC_SETTING_KEY = "provider_limits_auto_sync_last_run"; const DEFAULT_PROVIDER_LIMITS_POST_USAGE_REFRESH_DELAY_MS = 5_000; const pendingPostUsageRefreshes = new Set(); function getProviderLimitsPostUsageRefreshDelayMs(): number { const raw = Number(process.env.PROVIDER_LIMITS_POST_USAGE_REFRESH_DELAY_MS ?? ""); return Number.isFinite(raw) && raw >= 0 ? raw : DEFAULT_PROVIDER_LIMITS_POST_USAGE_REFRESH_DELAY_MS; } function scheduleProviderLimitsPostUsageRefresh(connectionId: string): void { if (!connectionId || pendingPostUsageRefreshes.has(connectionId)) return; pendingPostUsageRefreshes.add(connectionId); const timer = setTimeout(() => { pendingPostUsageRefreshes.delete(connectionId); void fetchAndPersistProviderLimits(connectionId, "scheduled", { allowRotatingRefresh: true, }).catch((error) => { const message = error instanceof Error ? error.message : String(error); console.warn( `[ProviderLimits] Post-usage refresh failed for connection ${connectionId}: ${message}` ); }); }, getProviderLimitsPostUsageRefreshDelayMs()); timer.unref?.(); } export function notifyProviderUsageRecorded( provider: string | null | undefined, connectionId: string | null | undefined ): void { if ((provider !== "antigravity" && provider !== "agy") || !connectionId) return; scheduleProviderLimitsPostUsageRefresh(connectionId); } // Subscribe at module load so usageHistory can emit usage events without importing // this module (and its executors/translator import graph). This module is loaded by // the provider-limits route and the background auto-sync scheduler at server boot. onUsageRecorded(notifyProviderUsageRecorded); function hasRetrieveUserQuotaSource( provider: string, cache: ProviderLimitsCacheEntry | undefined ): boolean { if (provider !== "antigravity" && provider !== "agy") return true; if (!cache?.quotas) return false; return Object.values(cache.quotas).some((quota) => { if (!isRecord(quota)) return false; return quota.quotaSource === "retrieveUserQuota"; }); } function sanitizeProviderLimitsCacheForConnection( connection: ProviderConnectionLike | null | undefined, entry: ProviderLimitsCacheEntry | null ): ProviderLimitsCacheEntry | null { if (!connection || !entry || !entry.quotas) return entry; if (connection.provider !== "antigravity" && connection.provider !== "agy") return entry; const sanitizedQuotas = normalizeUsageQuotasForProvider(connection.provider, entry.quotas); return sanitizedQuotas === entry.quotas ? entry : { ...entry, quotas: sanitizedQuotas }; } function shouldRefreshProviderLimitsCache( connection: ProviderConnectionLike, cache: ProviderLimitsCacheEntry | undefined ): boolean { if (!cache?.quotas) return true; if (connection.provider !== "antigravity" && connection.provider !== "agy") return false; return ( !hasRetrieveUserQuotaSource(connection.provider, cache) || Object.keys(cache.quotas).some( (quotaKey) => !isUsageQuotaKeyAllowed(connection.provider, quotaKey) ) ); } export function isSupportedUsageConnection(connection: ProviderConnectionLike | null): boolean { if ( !connection || !connection.provider || !USAGE_SUPPORTED_PROVIDERS.includes(connection.provider) ) { return false; } if (connection.authType === "oauth") return true; return ( (connection.authType === "apikey" || connection.authType === "api_key") && PROVIDER_LIMITS_APIKEY_PROVIDERS.has(connection.provider) ); } function withStatus(error: Error, status: number): Error & { status: number } { return Object.assign(error, { status }); } async function syncToCloudIfEnabled() { try { const machineId = await getMachineId(); if (!machineId) return; await syncToCloud(machineId); } catch (error) { console.error("[ProviderLimits] Error syncing refreshed credentials to cloud:", error); } } /** * Whether the quota path may refresh this provider's token. Exported for testing. * * Rotating-refresh providers (Codex/OpenAI share one Auth0 client_id, etc.) mint a * single-use refresh_token on every refresh. The BULK quota-sync path runs many * connections concurrently; refreshing sibling accounts in parallel makes Auth0 * revoke the whole token family (openai/codex#9648) and kills every account but * the last (#3019). So the bulk path never refreshes rotating providers * (`allowRotatingRefresh` falsy). The on-demand, per-connection path opts in and * is made safe by `serializeRefresh` (one token mint at a time per rotation group, * so even N concurrent per-account requests can never refresh siblings in * parallel). Non-rotating providers are always eligible. */ export function shouldAttemptRotatingRefresh( provider: string, allowRotatingRefresh: boolean | undefined ): boolean { if (rotationGroupFor(provider) === null) return true; return allowRotatingRefresh === true; } export async function refreshAndUpdateCredentials( connection: ProviderConnectionLike, opts: { allowRotatingRefresh?: boolean; force?: boolean } = {} ) { if (!shouldAttemptRotatingRefresh(connection.provider, opts.allowRotatingRefresh)) { return { connection, refreshed: false }; } const executor = await getExecutor(connection.provider); const credentials = { connectionId: connection.id, accessToken: connection.accessToken, refreshToken: connection.refreshToken, expiresAt: connection.tokenExpiresAt || connection.expiresAt || null, providerSpecificData: connection.providerSpecificData, copilotToken: connection.providerSpecificData?.copilotToken, copilotTokenExpiresAt: connection.providerSpecificData?.copilotTokenExpiresAt, }; // `force` is used ONLY on the reactive 401 recovery path (a usage fetch came // back unauthorized) — it bypasses the proactive `needsRefresh` heuristic so // imported accounts (expiresAt=null, where needsRefresh is always false) can // still re-mint. The mint stays serialized per rotation group; this never // refreshes proactively from the bulk path (#3019 guard above is unchanged). if (!opts.force && !executor.needsRefresh(credentials)) { return { connection, refreshed: false }; } // Serialize the actual token mint per rotation group so two sibling accounts // never hit Auth0 concurrently (passthrough for non-rotating providers). const refreshResult = (await serializeRefresh(connection.provider, () => executor.refreshCredentials(credentials, console) )) as | (JsonRecord & { accessToken?: string; refreshToken?: string; expiresIn?: number; expiresAt?: string; copilotToken?: string; copilotTokenExpiresAt?: string; }) | null; if (!refreshResult) { // Refresh failed but we still have an accessToken — fall back to the // existing token for ANY OAuth provider (graceful degradation) instead of // hard-failing. Previously this was qualified to `provider === "github"`, // which left every other provider stuck on a transient refresh failure even // when a usable access token was still on hand. if (connection.accessToken) { return { connection, refreshed: false }; } throw withStatus( new Error("Failed to refresh credentials. Please re-authorize the connection."), 401 ); } const updateData: JsonRecord = { updatedAt: new Date().toISOString(), }; if (refreshResult.accessToken) { updateData.accessToken = refreshResult.accessToken; } if (refreshResult.refreshToken) { updateData.refreshToken = refreshResult.refreshToken; } if (refreshResult.expiresIn) { const expiresAt = new Date(Date.now() + refreshResult.expiresIn * 1000).toISOString(); updateData.expiresAt = expiresAt; updateData.tokenExpiresAt = expiresAt; } else if (refreshResult.expiresAt) { updateData.expiresAt = refreshResult.expiresAt; updateData.tokenExpiresAt = refreshResult.expiresAt; } if (refreshResult.copilotToken || refreshResult.copilotTokenExpiresAt) { updateData.providerSpecificData = { ...(connection.providerSpecificData || {}), copilotToken: refreshResult.copilotToken, copilotTokenExpiresAt: refreshResult.copilotTokenExpiresAt, }; } await updateProviderConnection(connection.id, updateData); return { connection: { ...connection, ...updateData, providerSpecificData: (updateData.providerSpecificData as JsonRecord | undefined) || connection.providerSpecificData, }, refreshed: true, }; } function isUsageAuthError(message: unknown): boolean { if (typeof message !== "string") return false; const m = message.toLowerCase(); return ( m.includes("token expired") || m.includes("unauthorized") || m.includes("re-authenticate") || m.includes("access denied") || m.includes("invalidated") || m.includes("401") ); } function isNetworkFailureMessage(message: unknown): boolean { if (typeof message !== "string") return false; return ( message.includes("fetch failed") || message.includes("ECONNREFUSED") || message.includes("ETIMEDOUT") || message.includes("Proxy unreachable") || message.includes("UND_ERR_CONNECT_TIMEOUT") ); } function isAccountScopedProxyResolution(proxyInfo: unknown): boolean { if (!isRecord(proxyInfo)) return false; if (!proxyInfo.proxy) return false; return proxyInfo.level === "key" || proxyInfo.level === "account"; } function shouldFailClosedForProviderLimitsProxy( connection: ProviderConnectionLike, proxyInfo: unknown ): boolean { return connection.authType === "oauth" && isAccountScopedProxyResolution(proxyInfo); } /** * Decide whether the quota-sync path should flag a connection `expired` from an * auth-style usage error. Exported for unit testing. * * Rotating-refresh providers (Codex/OpenAI/Claude/etc. — see refreshSerializer's * ROTATION_LOCK_GROUP) have their access_token deliberately NOT proactively * refreshed in this quota path (#3019, to avoid the Auth0 family-revocation * cascade). So a "token expired" from the quota fetch is a recoverable * false-negative: the credential is still valid (its `expires_at` is in the * future) and the reactive, serialized 401 path refreshes the access_token on * next use. Flagging it `expired` hides a healthy account from the quota page * (observed: freshly-added Codex accounts flagged expired while a providers-page * refresh turns them green). So never mark a rotating provider expired from the * quota sync — leave its status to the reactive path / connection test. */ export function quotaPathShouldMarkExpired( provider: string, usageMessage: unknown, currentTestStatus: string | null | undefined ): boolean { if (currentTestStatus === "expired") return false; const message = typeof usageMessage === "string" ? usageMessage.toLowerCase() : ""; const isAuthError = message.includes("token expired") || message.includes("access denied") || message.includes("re-authenticate") || message.includes("unauthorized"); if (!isAuthError) return false; if (rotationGroupFor(provider) !== null) return false; return true; } const TERMINAL_STATUSES_FOR_QUOTA_RECOVERY = new Set([ "credits_exhausted", "banned", "expired", "deactivated", ]); function isTerminalStatusForQuotaRecovery(testStatus: string | null | undefined): boolean { if (!testStatus) return false; return TERMINAL_STATUSES_FOR_QUOTA_RECOVERY.has(testStatus); } export function hasUsableQuota(usage: JsonRecord): boolean { const quotas = usage?.quotas; if (!isRecord(quotas)) return false; for (const value of Object.values(quotas)) { if (!isRecord(value)) continue; if (value.unlimited === true) return true; const remaining = typeof value.remaining === "number" ? value.remaining : typeof value.remainingPercentage === "number" ? value.remainingPercentage : null; if (remaining !== null && remaining > 0) return true; } return false; } // A window "still blocks" recovery when it governs quota and is either still // exhausted with a real reset that hasn't passed yet, or exhausted with no // parseable real reset at all (unknown-reset windows stay locked, matching // the pre-existing kimi-coding partial-refresh semantics). function windowStillExhaustedAfterRealReset(value: unknown, nowMs: number): boolean { if (!isRecord(value)) return false; if (value.unlimited === true) return false; const remaining = typeof value.remaining === "number" ? value.remaining : typeof value.remainingPercentage === "number" ? value.remainingPercentage : null; if (remaining !== null && remaining > 0) return false; if (value.resetAt == null) return true; const resetMs = Date.parse(String(value.resetAt)); if (Number.isNaN(resetMs)) return true; return resetMs > nowMs; } export async function maybeClearRecoveredQuotaState( connection: ProviderConnectionLike, usage: JsonRecord ): Promise { if (!hasUsableQuota(usage)) return connection; if (isTerminalStatusForQuotaRecovery(connection.testStatus)) return connection; if (connection.lastErrorType === "quota_exhausted") { if ( connection.lastErrorSource === CLAUDE_EXTRA_USAGE_ERROR_SOURCE && isClaudeExtraUsageBlockEnabled(connection.provider, connection.providerSpecificData) && isClaudeExtraUsageQueued(usage) ) { // Claude's pay-as-you-go extra-usage block is orthogonal to the // session/weekly quota windows checked below: the upstream can report a // fully recovered quota window while extraUsage.queued is still true. // Only syncClaudeExtraUsageStateIfNeeded (buildClaudeExtraUsageConnectionUpdate) // owns clearing this specific state — the general window-recovery logic // below must not release it just because some quota window looks fresh. return connection; } const quotas = usage?.quotas; if (isRecord(quotas)) { // Honor the REAL per-window resetAt from the freshly fetched quota // instead of the synthetic cooldown persisted at failure time (e.g. // Claude's flat 1h SUBSCRIPTION_QUOTA_COOLDOWN_MS when no upstream // reset was parseable). Only stay locked if some window that governs // this connection's quota is still demonstrably exhausted. const anyStillBlocking = Object.values(quotas).some((value) => windowStillExhaustedAfterRealReset(value, Date.now()) ); if (anyStillBlocking) return connection; } else if ( connection.rateLimitedUntil && new Date(connection.rateLimitedUntil).getTime() > Date.now() ) { // No quota object at all (degraded/failed fetch shape) — fall back to // the previous synthetic-cooldown guard. return connection; } } const hasTransientState = connection.testStatus === "unavailable" || Boolean(connection.rateLimitedUntil) || Boolean(connection.lastError) || Boolean(connection.errorCode) || Boolean(connection.lastErrorType) || Boolean(connection.lastErrorSource) || (connection.backoffLevel ?? 0) > 0; if (!hasTransientState) return connection; let cleared = true; try { const result = await clearRecoveredProviderState( { connectionId: connection.id, testStatus: connection.testStatus, lastError: connection.lastError ?? null, rateLimitedUntil: connection.rateLimitedUntil ?? null, errorCode: connection.errorCode ?? null, lastErrorType: connection.lastErrorType ?? null, lastErrorSource: connection.lastErrorSource ?? null, }, { testStatus: connection.testStatus ?? null, lastErrorAt: connection.lastErrorAt ?? null, rateLimitedUntil: connection.rateLimitedUntil ?? null, } ); cleared = result.applied; } catch (dbError) { console.warn("[ProviderLimits] Failed to clear recovered quota state:", dbError); return connection; } if (!cleared) { // CAS miss — a concurrent writer (markAccountUnavailable, etc.) updated // the row between our read and the clear. Return the original snapshot; // the next read from DB will surface the fresh state. return connection; } return { ...connection, testStatus: "active", lastError: null, lastErrorAt: null, lastErrorType: null, lastErrorSource: null, errorCode: null, rateLimitedUntil: null, backoffLevel: 0, }; } async function syncExpiredStatusIfNeeded( connection: ProviderConnectionLike, usage: JsonRecord ): Promise { if (!quotaPathShouldMarkExpired(connection.provider, usage.message, connection.testStatus)) { return connection; } try { await updateProviderConnection(connection.id, { testStatus: "expired", lastErrorType: "token_expired", lastErrorAt: new Date().toISOString(), }); } catch (dbError) { console.error("[ProviderLimits] Failed to sync expired status to DB:", dbError); return connection; } return { ...connection, testStatus: "expired", lastErrorType: "token_expired", }; } async function syncClaudeExtraUsageStateIfNeeded( connection: ProviderConnectionLike, usage: JsonRecord ): Promise { const update = buildClaudeExtraUsageConnectionUpdate(connection, usage); if (!update) return connection; await updateProviderConnection(connection.id, update); return { ...connection, ...update, }; } /** Persist Antigravity tier from live loadCodeAssist on quota refresh (not only OAuth). */ async function syncAntigravitySubscriptionIfNeeded( connection: ProviderConnectionLike, usage: JsonRecord ): Promise { if (connection.provider !== "antigravity" && connection.provider !== "agy") return connection; const subscriptionInfo = usage.subscriptionInfo; if (!subscriptionInfo) return connection; const psd = (connection.providerSpecificData || {}) as JsonRecord; const nextPsd: JsonRecord = { ...psd }; let changed = false; const tierId = extractCodeAssistOnboardTierId(subscriptionInfo); if (tierId && tierId !== "legacy-tier" && psd.tier !== tierId) { nextPsd.tier = tierId; changed = true; } const subscriptionTier = extractCodeAssistSubscriptionTier(subscriptionInfo); if (subscriptionTier && psd.subscriptionTier !== subscriptionTier) { nextPsd.subscriptionTier = subscriptionTier; changed = true; } const plan = typeof usage.plan === "string" ? usage.plan.trim() : ""; if (plan && psd.plan !== plan) { nextPsd.plan = plan; changed = true; } const discoveredProjectId = extractAntigravityProjectIdFromPayload( subscriptionInfo as Record ); const storedProjectId = getStoredAntigravityProjectId(connection); let nextProjectId: string | undefined; if (discoveredProjectId && !storedProjectId) { nextPsd.projectId = discoveredProjectId; nextProjectId = discoveredProjectId; changed = true; } if (!changed) return connection; await updateProviderConnection(connection.id, { ...(nextProjectId ? { projectId: nextProjectId, errorCode: null, lastError: null } : {}), providerSpecificData: nextPsd, }); return { ...connection, ...(nextProjectId ? { projectId: nextProjectId, errorCode: null, lastError: null } : {}), providerSpecificData: nextPsd, }; } /** Persist refreshed Claude bootstrap fields into psd; writes only on diff. */ async function syncClaudeBootstrapIfNeeded( connection: ProviderConnectionLike, usage: JsonRecord ): Promise { if (connection.provider !== "claude") return connection; const bootstrap = usage?.bootstrap as Record | null | undefined; if (!bootstrap || typeof bootstrap !== "object") return connection; const psd = (connection.providerSpecificData || {}) as JsonRecord; const mapping: Array<[keyof typeof bootstrap, string]> = [ ["account_uuid", "accountUUID"], ["organization_uuid", "organizationUUID"], ["organization_name", "organizationName"], ["organization_type", "organizationType"], ["organization_rate_limit_tier", "organizationRateLimitTier"], ]; const nextPsd: JsonRecord = { ...psd }; let changed = false; for (const [bsKey, psdKey] of mapping) { const next = bootstrap[bsKey]; if (typeof next === "string" && next.length > 0 && psd[psdKey] !== next) { nextPsd[psdKey] = next; changed = true; } } if (!changed) return connection; await updateProviderConnection(connection.id, { providerSpecificData: nextPsd }); return { ...connection, providerSpecificData: nextPsd, }; } export function getProviderLimitsSyncIntervalMinutes(): number { const raw = Number.parseInt(process.env.PROVIDER_LIMITS_SYNC_INTERVAL_MINUTES ?? "", 10); return Number.isFinite(raw) && raw > 0 ? raw : DEFAULT_PROVIDER_LIMITS_SYNC_INTERVAL_MINUTES; } export function getProviderLimitsSyncIntervalMs(): number { return getProviderLimitsSyncIntervalMinutes() * 60 * 1000; } /** Default gap (ms) inserted between two consecutive OAuth quota fetches. */ const DEFAULT_PROVIDER_LIMITS_SYNC_SPACING_MS = 1500; /** * Spacing (ms) applied between consecutive provider-limits fetch batches in a * bulk sync, for BOTH the OAuth and local/API-key paths. * * OAuth providers (Codex/Claude/Kimi-coding/…) are fetched ONE AT A TIME with * this gap so a single host never bursts several simultaneous usage/refresh * requests to the same upstream — bursts read as automated traffic and * contribute to session termination / anomaly flags (and, for rotating-token * providers, to the Auth0 family-revocation race). Local/API-key connections * (e.g. Ollama) keep their fast in-chunk concurrent path, but the gap is now * also applied BETWEEN concurrency chunks so a local endpoint isn't hit by a * simultaneous refresh burst either (#6916). Tunable via * `PROVIDER_LIMITS_SYNC_SPACING_MS`; set to `"0"` to opt out on either path. */ export function getProviderLimitsSyncSpacingMs(): number { const rawEnv = process.env.PROVIDER_LIMITS_SYNC_SPACING_MS; if (rawEnv === undefined || rawEnv === "") return DEFAULT_PROVIDER_LIMITS_SYNC_SPACING_MS; const raw = Number(rawEnv); return Number.isFinite(raw) && raw >= 0 ? raw : DEFAULT_PROVIDER_LIMITS_SYNC_SPACING_MS; } export async function getLastProviderLimitsAutoSyncTime(): Promise { try { const settings = await getSettings(); const value = settings[PROVIDER_LIMITS_AUTO_SYNC_SETTING_KEY]; return typeof value === "string" && value.trim() ? value : null; } catch { return null; } } async function setLastProviderLimitsAutoSyncTime(timestamp: string): Promise { await updateSettings({ [PROVIDER_LIMITS_AUTO_SYNC_SETTING_KEY]: timestamp }); } export function getCachedProviderLimitsMap(): Record { return getAllProviderLimitsCache(); } export async function getSanitizedCachedProviderLimitsMap(): Promise< Record > { const caches = getAllProviderLimitsCache(); // Sanitization only rewrites Antigravity/agy quota keys; every other provider's cache // entry is returned untouched (see sanitizeProviderLimitsCacheForConnection). The // dashboard polls this on an auto-refresh interval, so avoid the unconditional // `SELECT * FROM provider_connections` + per-row credential decryption that the // previous implementation paid on every poll: skip the scan entirely when nothing is // cached, and otherwise fetch ONLY the Antigravity/agy connections. For any other // provider, byId.get(id) is undefined and the entry is returned verbatim — identical // output to scanning every active connection, but without decrypting unrelated keys. // (LEDGER-2 / #3821-review) const connectionIds = Object.keys(caches); if (connectionIds.length === 0) return {}; const sanitizableConnections = [ ...((await getProviderConnections({ isActive: true, provider: "antigravity", })) as unknown as ProviderConnectionLike[]), ...((await getProviderConnections({ isActive: true, provider: "agy", })) as unknown as ProviderConnectionLike[]), ]; if (sanitizableConnections.length === 0) { // No connection can change the cache → return the raw entries unchanged. return { ...caches }; } const byId = new Map(sanitizableConnections.map((conn) => [conn.id, conn])); const sanitized: Record = {}; for (const [connectionId, entry] of Object.entries(caches)) { sanitized[connectionId] = sanitizeProviderLimitsCacheForConnection(byId.get(connectionId), entry) || entry; } return sanitized; } export async function fetchLiveProviderLimits(connectionId: string): Promise<{ connection: ProviderConnectionLike; usage: JsonRecord; }> { return fetchLiveProviderLimitsWithOptions(connectionId, { forceRefresh: false }); } async function fetchLiveProviderLimitsWithOptions( connectionId: string, options: { forceRefresh?: boolean; allowRotatingRefresh?: boolean } = {} ): Promise<{ connection: ProviderConnectionLike; usage: JsonRecord; }> { if (await isConnectionUnavailableToAuxiliaryActivity(connectionId)) { throw withStatus(new Error("Usage refresh deferred while an exclusive lease is active"), 409); } let connection = (await getProviderConnectionById( connectionId )) as unknown as ProviderConnectionLike | null; if (!connection) { throw withStatus(new Error("Connection not found"), 404); } if (!isSupportedUsageConnection(connection)) { throw withStatus(new Error("Usage not available for this connection"), 400); } if (connection.authType !== "oauth") { // L3: route the API-key usage/quota fetch through the connection's proxy context, // mirroring the OAuth branch below (proxyInfo?.proxy ?? null). Without this, API-key // usage egresses on the host IP, ignoring the connection's assigned proxy. const apiKeyProxy = await resolveProxyForConnection(connectionId); const usage = sanitizeUsageQuotasForProvider( connection.provider, (await runWithProxyContext(apiKeyProxy?.proxy ?? null, () => getUsageForProvider(connection as unknown as JsonRecord, options) )) as JsonRecord ); if (isRecord(usage.quotas)) { setQuotaCache(connectionId, connection.provider, usage.quotas); } connection = await syncExpiredStatusIfNeeded(connection, usage); connection = await syncClaudeExtraUsageStateIfNeeded(connection, usage); connection = await syncClaudeBootstrapIfNeeded(connection, usage); connection = await syncAntigravitySubscriptionIfNeeded(connection, usage); connection = await maybeClearRecoveredQuotaState(connection, usage); return { connection, usage }; } const proxyInfo = await resolveProxyForConnection(connectionId); const fetchUsageWithContext = async (proxyConfig: unknown) => runWithProxyContext(proxyConfig, async () => { let conn = connection as ProviderConnectionLike; let wasRefreshed = false; const result = await refreshAndUpdateCredentials(conn, { allowRotatingRefresh: options.allowRotatingRefresh, }); conn = result.connection; wasRefreshed = result.refreshed; if (wasRefreshed) { await syncToCloudIfEnabled(); } let usageData = sanitizeUsageQuotasForProvider( conn.provider, (await getUsageForProvider(conn as unknown as JsonRecord, options)) as JsonRecord ); // Reactive 401 recovery (on-demand/force path only): an unauthorized usage // response means the access token is actually dead. Force ONE serialized // re-mint and retry once. This recovers imported accounts (expiresAt=null, // where the proactive needsRefresh heuristic never fires) without ever // refreshing proactively from the bulk path. if (options.allowRotatingRefresh && !wasRefreshed && isUsageAuthError(usageData?.message)) { const forced = await refreshAndUpdateCredentials(conn, { allowRotatingRefresh: true, force: true, }); if (forced.refreshed) { conn = forced.connection; await syncToCloudIfEnabled(); usageData = sanitizeUsageQuotasForProvider( conn.provider, (await getUsageForProvider(conn as unknown as JsonRecord, options)) as JsonRecord ); } } connection = conn; return { usage: usageData }; }); let result: { usage: JsonRecord }; const proxyConfig = proxyInfo?.proxy || null; const failClosedOnProxyFailure = shouldFailClosedForProviderLimitsProxy(connection, proxyInfo); try { result = await fetchUsageWithContext(proxyConfig); } catch (error: any) { const isThrownNetworkError = error?.message === "fetch failed" || error?.code === "PROXY_UNREACHABLE" || error?.code === "UND_ERR_CONNECT_TIMEOUT" || error?.cause?.code === "ECONNREFUSED"; if (proxyConfig && isThrownNetworkError) { if (failClosedOnProxyFailure) { console.warn( "[ProviderLimits] Account-scoped %s proxy fetch failed for %s; failing closed without direct retry:", connection.provider, connectionId, error?.message ); throw error; } console.warn( "[ProviderLimits] Proxy fetch threw for %s, retrying without proxy:", connectionId, error?.message ); result = await fetchUsageWithContext(null); } else { throw error; } } if (proxyConfig && isNetworkFailureMessage(result.usage?.message)) { if (failClosedOnProxyFailure) { const message = typeof result.usage.message === "string" ? result.usage.message : "Provider-limits proxy request failed"; console.warn( "[ProviderLimits] Account-scoped %s proxy usage failed for %s; failing closed without direct retry:", connection.provider, connectionId, message ); throw withStatus(new Error(message), 503); } console.warn( "[ProviderLimits] Proxy usage returned network error for %s, retrying without proxy:", connectionId, result.usage.message ); result = await fetchUsageWithContext(null); } if (isRecord(result.usage.quotas)) { setQuotaCache(connectionId, connection.provider, result.usage.quotas); } connection = await syncExpiredStatusIfNeeded(connection, result.usage); connection = await syncClaudeExtraUsageStateIfNeeded(connection, result.usage); connection = await syncClaudeBootstrapIfNeeded(connection, result.usage); connection = await syncAntigravitySubscriptionIfNeeded(connection, result.usage); connection = await maybeClearRecoveredQuotaState(connection, result.usage); return { connection, usage: result.usage, }; } export async function fetchAndPersistProviderLimits( connectionId: string, source: SyncSource = "manual", opts: { allowRotatingRefresh?: boolean } = {} ): Promise<{ connection: ProviderConnectionLike; usage: JsonRecord; cache: ProviderLimitsCacheEntry; }> { const { connection, usage } = await fetchLiveProviderLimitsWithOptions(connectionId, { forceRefresh: source === "manual", allowRotatingRefresh: opts.allowRotatingRefresh, }); const newCache = toProviderLimitsCacheEntry(usage, source); const previous = getProviderLimitsCache(connectionId); const cache = mergeProviderLimitsCacheEntry(connection.provider, newCache, previous); // Don't persist error-only entries (429 etc.) — would wipe prior good cache. // Serve the prior entry instead; only successful fetches update the cache. if (cache === previous && newCache.message) { const staleUsage: JsonRecord = { ...usage, quotas: previous.quotas, plan: previous.plan ?? usage.plan ?? null, bankedResetCredits: previous.bankedResetCredits, billing: previous.billing, message: null, _stale: true, _staleSince: previous.fetchedAt, _staleReason: newCache.message, }; return { connection, usage: staleUsage, cache: previous }; } const mergedUsage: JsonRecord = { ...usage, ...(cache.billing ? { billing: cache.billing } : {}), }; setProviderLimitsCache(connectionId, cache); return { connection, usage: mergedUsage, cache }; } export async function syncAllProviderLimits( options: { source?: SyncSource; concurrency?: number; } = {} ): Promise<{ total: number; succeeded: number; failed: number; caches: Record; errors: Record; }> { const { source = "manual", concurrency = 5 } = options; const connectionRows = (await getProviderConnections({ isActive: true, })) as unknown as ProviderConnectionLike[]; const connections = ( await Promise.all( connectionRows.map(async (connection) => ({ connection, blocked: await isConnectionUnavailableToAuxiliaryActivity(connection.id), })) ) ) .filter(({ connection, blocked }) => isSupportedUsageConnection(connection) && !blocked) .map(({ connection }) => connection); const cacheEntries: Array<{ connectionId: string; entry: ProviderLimitsCacheEntry }> = []; const caches: Record = {}; const errors: Record = {}; const recordResult = ( connectionId: string, result: PromiseSettledResult<{ connectionId: string; cache: ProviderLimitsCacheEntry }> ) => { if (result.status === "fulfilled") { const { cache } = result.value; const previous = getProviderLimitsCache(connectionId); if (cache === previous) { caches[connectionId] = cache; return; } cacheEntries.push({ connectionId, entry: cache }); caches[connectionId] = cache; return; } const reason = result.reason as { message?: string } | undefined; errors[connectionId] = reason?.message || "Failed to refresh provider limits"; }; const fetchOne = async (connection: ProviderConnectionLike) => { const existingCache = getProviderLimitsCache(connection.id); const forceRefresh = source === "manual" || shouldRefreshProviderLimitsCache(connection, existingCache || undefined); const { usage } = await fetchLiveProviderLimitsWithOptions(connection.id, { forceRefresh, }); const nextCache = toProviderLimitsCacheEntry(usage, source); const cache = mergeProviderLimitsCacheEntry(connection.provider, nextCache, existingCache); return { connectionId: connection.id, cache }; }; // OAuth connections are processed STRICTLY SEQUENTIALLY (chunk size 1) with a // spacing gap so a single host never bursts simultaneous usage/refresh // requests to the same upstream (anomaly/session-termination guard; see // getProviderLimitsSyncSpacingMs). Local/API-key connections keep their fast // in-chunk concurrent path, spaced BETWEEN chunks (#6916). const oauthConnections = connections.filter((c) => c.authType === "oauth"); const otherConnections = connections.filter((c) => c.authType !== "oauth"); const spacingMs = getProviderLimitsSyncSpacingMs(); const recordChunk = ( chunk: ProviderConnectionLike[], results: PromiseSettledResult<{ connectionId: string; cache: ProviderLimitsCacheEntry }>[] ) => { results.forEach((result, index) => { const connectionId = chunk[index]?.id; if (connectionId) recordResult(connectionId, result); }); }; await syncInChunksWithSpacing(otherConnections, concurrency, spacingMs, fetchOne, recordChunk); await syncInChunksWithSpacing(oauthConnections, 1, spacingMs, fetchOne, recordChunk); if (cacheEntries.length > 0) { setProviderLimitsCacheBatch(cacheEntries); } if (source === "scheduled") { await setLastProviderLimitsAutoSyncTime(new Date().toISOString()); } return { total: connections.length, succeeded: cacheEntries.length, failed: connections.length - cacheEntries.length, caches, errors, }; }