import { dirname } from "node:path"; import { mkdir, readFile, rename, rm, stat, writeFile } from "node:fs/promises"; import type { LoadResult, ParsedExtraQuotaLimit, ParsedQuotaMetadata, PendingQuotaObservation, QuotaDimensionObservation, QuotaObservation, QuotaState, } from "./types.js"; const LOCK_STALE_MS = 10_000; const LOCK_WAIT_MS = 5_000; export function emptyState(): QuotaState { return { version: 1, observations: {} }; } export function normalizeState(value: unknown): QuotaState { if (!isRecord(value)) return emptyState(); const observations = normalizeObservationMap(value.observations); const pendingObservations = normalizePendingObservationMap( value.pendingObservations, ); return { version: 1, observations, ...(Object.keys(pendingObservations).length > 0 ? { pendingObservations } : {}), }; } function normalizeObservationMap( value: unknown, ): Record { if (!isRecord(value)) return {}; return Object.fromEntries( Object.entries(value).flatMap(([key, observation]) => { const normalized = normalizeObservation(observation); return normalized ? [[key, normalized]] : []; }), ); } function normalizePendingObservationMap( value: unknown, ): Record { if (!isRecord(value)) return {}; return Object.fromEntries( Object.entries(value).flatMap(([key, pending]) => { const normalized = normalizePendingObservation(pending); return normalized ? [[key, normalized]] : []; }), ); } function normalizeObservation(value: unknown): QuotaObservation | undefined { if (!isRecord(value)) return undefined; const provider = nonEmptyString(value.provider); const model = nonEmptyString(value.model); const source = observationSource(value.source); const observedAt = finiteNumber(value.observedAt); const updatedAt = finiteNumber(value.updatedAt); const dimensions = Array.isArray(value.dimensions) ? value.dimensions .map(normalizeDimension) .filter((dimension): dimension is QuotaDimensionObservation => Boolean(dimension), ) : []; if ( !provider || !model || !source || observedAt === undefined || updatedAt === undefined || dimensions.length === 0 ) return undefined; const status = finiteNumber(value.status); const adapter = nonEmptyString(value.adapter); const metadata = normalizeMetadata(value.metadata); return { provider, model, ...(adapter ? { adapter } : {}), source, ...(status !== undefined ? { status } : {}), observedAt, updatedAt, dimensions, ...(metadata ? { metadata } : {}), }; } function normalizeDimension( value: unknown, ): QuotaDimensionObservation | undefined { if (!isRecord(value)) return undefined; const name = nonEmptyString(value.name); const source = observationSource(value.source); const observedAt = finiteNumber(value.observedAt); if (!name || !source || observedAt === undefined) return undefined; const limit = finiteNumber(value.limit); const remaining = finiteNumber(value.remaining); const resetAt = finiteNumber(value.resetAt); return { name, ...(limit !== undefined ? { limit } : {}), ...(remaining !== undefined ? { remaining } : {}), ...(resetAt !== undefined ? { resetAt } : {}), observedAt, source, }; } function normalizePendingObservation( value: unknown, ): PendingQuotaObservation | undefined { if (!isRecord(value)) return undefined; const observation = normalizeObservation(value.observation); const reason = nonEmptyString(value.reason); const createdAt = finiteNumber(value.createdAt); const updatedAt = finiteNumber(value.updatedAt); if (!observation || !reason || createdAt === undefined || updatedAt === undefined) return undefined; const priorRemaining = finiteNumber(value.priorRemaining); const newRemaining = finiteNumber(value.newRemaining); const resetAt = finiteNumber(value.resetAt); return { observation, reason, createdAt, updatedAt, ...(priorRemaining !== undefined ? { priorRemaining } : {}), ...(newRemaining !== undefined ? { newRemaining } : {}), ...(resetAt !== undefined ? { resetAt } : {}), }; } function normalizeMetadata(value: unknown): ParsedQuotaMetadata | undefined { if (!isRecord(value)) return undefined; const extraLimits = Array.isArray(value.extraLimits) ? value.extraLimits .map(normalizeExtraQuotaLimit) .filter((limit): limit is ParsedExtraQuotaLimit => Boolean(limit)) : []; const metadata = { ...(typeof value.accountHeaderSent === "boolean" ? { accountHeaderSent: value.accountHeaderSent } : {}), ...(typeof value.allowed === "boolean" ? { allowed: value.allowed } : {}), ...(typeof value.limitReached === "boolean" ? { limitReached: value.limitReached } : {}), ...(typeof value.rateLimitReachedType === "string" || value.rateLimitReachedType === null ? { rateLimitReachedType: value.rateLimitReachedType } : {}), ...(typeof value.codexCliRpcFallback === "boolean" ? { codexCliRpcFallback: value.codexCliRpcFallback } : {}), ...(extraLimits.length > 0 ? { extraLimits } : {}), }; return Object.keys(metadata).length > 0 ? metadata : undefined; } function normalizeExtraQuotaLimit( value: unknown, ): ParsedExtraQuotaLimit | undefined { if (!isRecord(value)) return undefined; const name = nonEmptyString(value.name); if (!name) return undefined; const usedPercent = finiteNumber(value.usedPercent); const remaining = finiteNumber(value.remaining); const resetAt = finiteNumber(value.resetAt); return { name, ...(usedPercent !== undefined ? { usedPercent } : {}), ...(remaining !== undefined ? { remaining } : {}), ...(resetAt !== undefined ? { resetAt } : {}), }; } function isRecord(value: unknown): value is Record { return Boolean(value && typeof value === "object" && !Array.isArray(value)); } function nonEmptyString(value: unknown): string | undefined { return typeof value === "string" && value.length > 0 ? value : undefined; } function finiteNumber(value: unknown): number | undefined { return typeof value === "number" && Number.isFinite(value) ? value : undefined; } function observationSource( value: unknown, ): QuotaDimensionObservation["source"] | undefined { return value === "headers" || value === "429" || value === "fallback" || value === "subscription" ? value : undefined; } export async function readJsonFile( file: string, fallback: T, ): Promise> { try { const text = await readFile(file, "utf8"); return { path: file, value: JSON.parse(text) as T }; } catch (error) { if (isCode(error, "ENOENT")) return { path: file, value: cloneFallback(fallback) }; return { path: file, value: cloneFallback(fallback), error: errorToString(error), }; } } export async function loadState( stateFile: string, ): Promise> { const result = await readJsonFile(stateFile, emptyState()); return { path: stateFile, value: normalizeState(result.value), error: result.error ?? stateStructuralError(result.value), }; } export async function writeJsonFileAtomic( file: string, value: unknown, options: { exclusive?: boolean } = {}, ): Promise { await mkdir(dirname(file), { recursive: true }); const text = `${JSON.stringify(value, null, 2)}\n`; if (options.exclusive) { try { await writeFile(file, text, { encoding: "utf8", flag: "wx" }); return; } catch (error) { if (!isCode(error, "EEXIST")) throw error; return; } } const tempFile = `${file}.${process.pid}.${Date.now()}.${Math.random().toString(16).slice(2)}.tmp`; try { await writeFile(tempFile, text, "utf8"); await rename(tempFile, file); } catch (error) { await rm(tempFile, { force: true }).catch(() => undefined); throw error; } } export async function mergeStateFile( stateFile: string, mutator: ( state: QuotaState, ) => QuotaState | void | Promise, ): Promise { await mkdir(dirname(stateFile), { recursive: true }); return withFileLock(`${stateFile}.lock`, async () => { const current = await loadState(stateFile); const cloned = normalizeState(current.value); const next = (await mutator(cloned)) ?? cloned; const normalized = normalizeState(next); await writeJsonFileAtomic(stateFile, normalized); return normalized; }); } async function withFileLock( lockDir: string, callback: () => Promise, ): Promise { const start = Date.now(); while (true) { try { await mkdir(lockDir, { recursive: false }); break; } catch (error) { if (!isCode(error, "EEXIST")) throw error; await removeStaleLock(lockDir); if (Date.now() - start > LOCK_WAIT_MS) throw new Error(`Timed out waiting for lock: ${lockDir}`); await sleep(25 + Math.floor(Math.random() * 50)); } } try { return await callback(); } finally { await rm(lockDir, { recursive: true, force: true }).catch(() => undefined); } } async function removeStaleLock(lockDir: string): Promise { try { const info = await stat(lockDir); if (Date.now() - info.mtimeMs > LOCK_STALE_MS) await rm(lockDir, { recursive: true, force: true }); } catch { // Missing or unreadable lock will be retried by the caller. } } function sleep(ms: number): Promise { return new Promise((resolve) => setTimeout(resolve, ms)); } function stateStructuralError(value: unknown): string | undefined { if (!isRecord(value)) return "Invalid state: expected a JSON object."; if (value.observations !== undefined && !isRecord(value.observations)) return "Invalid state: observations must be an object."; if (isRecord(value.observations)) { for (const [key, observation] of Object.entries(value.observations)) { if (!isRecord(observation)) return `Invalid state: observations[${key}] must be an object.`; if (!Array.isArray(observation.dimensions)) return `Invalid state: observations[${key}].dimensions must be an array.`; if (observation.dimensions.some((dimension) => !isRecord(dimension))) return `Invalid state: observations[${key}].dimensions contains an invalid entry.`; } } if ( value.pendingObservations !== undefined && !isRecord(value.pendingObservations) ) return "Invalid state: pendingObservations must be an object."; if (isRecord(value.pendingObservations)) { for (const [key, pending] of Object.entries(value.pendingObservations)) { if (!isRecord(pending)) return `Invalid state: pendingObservations[${key}] must be an object.`; if (!isRecord(pending.observation)) return `Invalid state: pendingObservations[${key}].observation must be an object.`; if (!Array.isArray(pending.observation.dimensions)) return `Invalid state: pendingObservations[${key}].observation.dimensions must be an array.`; } } return undefined; } function cloneFallback(value: T): T { if (value === undefined) return value; return structuredClone(value) as T; } function isCode(error: unknown, code: string): boolean { return Boolean( error && typeof error === "object" && "code" in error && (error as { code?: string }).code === code, ); } function errorToString(error: unknown): string { return error instanceof Error ? error.message : String(error); }