import fs from "node:fs" import path from "node:path" import { randomUUID } from "node:crypto" import { DatabaseSync } from "node:sqlite" export type TokenUsage = { apiKey: string accountId?: string provider: string model: string inputTokens: number outputTokens: number totalTokens: number cachedTokens?: number cacheWriteTokens?: number /** Observed increase in Codex's absolute plan-window percentage. */ planPrimaryDelta?: number planSecondaryDelta?: number recordedAt?: number } export type UsageTotals = { inputTokens: number outputTokens: number totalTokens: number requests: number } export type UsageHistoryPoint = { at: number totalTokens: number cacheHitTokens: number cacheMissTokens: number requests: number } export type ModelObservedUsage = { requests: number totalTokens: number planPercentDelta: number percentPer1MTokens: number | null } const DAY_MS = 24 * 60 * 60 * 1000 export type UpstreamErrorRecord = { id?: string timestamp?: number provider: string model: string accountId?: string route: string status: number errorType: string message: string requestId?: string correlationId?: string } export type ListedUpstreamErrors = { errors: Array total: number limit: number offset: number hasMore: boolean } type UsageDb = { db: DatabaseSync; path: string } const databases = new Map() function usageDatabasePath(stateDirectory: string): string { return path.join(stateDirectory, "usage.sqlite") } export type ConversationRouteKey = { apiKeyId: string provider: string protocol: string conversationHash: string } export type ConversationSnapshot = { route: ConversationRouteKey accountId: string model: string serviceTier?: "priority" | "fast" callerHistory: string upstreamHistory: string fixedContextDigest: string lastContextTokens: number | null generation: number lastActivityAt: number idleDueAt: number | null lastCompactedAt: number | null leaseOwner: string | null leaseExpiresAt: number | null lastCompactionError: string | null } export type ConversationLease = { acquired: boolean generation: number leaseOwner: string | null } function snapshotParams(route: ConversationRouteKey): string[] { return routeParams(route) } export function getDatabase(stateDirectory: string): UsageDb { const dbPath = usageDatabasePath(stateDirectory) const existing = databases.get(dbPath) if (existing) return existing fs.mkdirSync(stateDirectory, { recursive: true, mode: 0o700 }) try { fs.chmodSync(stateDirectory, 0o700) } catch { /* best effort */ } const db = new DatabaseSync(dbPath) try { fs.chmodSync(dbPath, 0o600) } catch { /* best effort */ } db.exec("PRAGMA busy_timeout = 5000; PRAGMA journal_mode = WAL; PRAGMA synchronous = FULL;") db.exec(`CREATE TABLE IF NOT EXISTS token_usage ( id INTEGER PRIMARY KEY, api_key TEXT NOT NULL, provider TEXT NOT NULL, model TEXT NOT NULL, input_tokens INTEGER NOT NULL DEFAULT 0, output_tokens INTEGER NOT NULL DEFAULT 0, total_tokens INTEGER NOT NULL DEFAULT 0, cached_tokens INTEGER NOT NULL DEFAULT 0, cache_write_tokens INTEGER NOT NULL DEFAULT 0, plan_primary_delta REAL, plan_secondary_delta REAL, recorded_at INTEGER NOT NULL, account_id TEXT ); CREATE INDEX IF NOT EXISTS token_usage_key_time ON token_usage(api_key, recorded_at); CREATE INDEX IF NOT EXISTS token_usage_provider_time ON token_usage(provider, recorded_at); CREATE TABLE IF NOT EXISTS conversation_routes ( api_key_id TEXT NOT NULL, provider TEXT NOT NULL, protocol TEXT NOT NULL, conversation_hash TEXT NOT NULL, account_id TEXT NOT NULL, last_seen_at INTEGER NOT NULL, PRIMARY KEY (api_key_id, provider, protocol, conversation_hash) ); CREATE INDEX IF NOT EXISTS conversation_routes_last_seen ON conversation_routes(last_seen_at); CREATE INDEX IF NOT EXISTS conversation_routes_account ON conversation_routes(account_id); CREATE TABLE IF NOT EXISTS conversation_snapshots ( api_key_id TEXT NOT NULL, provider TEXT NOT NULL, protocol TEXT NOT NULL, conversation_hash TEXT NOT NULL, account_id TEXT NOT NULL, model TEXT NOT NULL, service_tier TEXT, caller_history TEXT NOT NULL, upstream_history TEXT NOT NULL, fixed_context_digest TEXT NOT NULL, last_context_tokens INTEGER, generation INTEGER NOT NULL DEFAULT 0, last_activity_at INTEGER NOT NULL, idle_due_at INTEGER, last_compacted_at INTEGER, lease_owner TEXT, lease_expires_at INTEGER, last_compaction_error TEXT, PRIMARY KEY (api_key_id, provider, protocol, conversation_hash) ); CREATE INDEX IF NOT EXISTS conversation_snapshots_activity ON conversation_snapshots(last_activity_at); CREATE INDEX IF NOT EXISTS conversation_snapshots_idle ON conversation_snapshots(idle_due_at); CREATE TABLE IF NOT EXISTS routing_cursors ( api_key_id TEXT NOT NULL, provider TEXT NOT NULL, next_index INTEGER NOT NULL, PRIMARY KEY (api_key_id, provider) ); CREATE TABLE IF NOT EXISTS upstream_errors ( id TEXT PRIMARY KEY, recorded_at INTEGER NOT NULL, provider TEXT NOT NULL, model TEXT NOT NULL, account_id TEXT, route TEXT NOT NULL, status INTEGER NOT NULL, error_type TEXT NOT NULL, message TEXT NOT NULL, request_id TEXT, correlation_id TEXT ); CREATE INDEX IF NOT EXISTS upstream_errors_time ON upstream_errors(recorded_at DESC, id DESC);`) for (const table of ["request_metrics", "compaction_events"]) { db.exec(`CREATE TABLE IF NOT EXISTS ${table} ( id TEXT PRIMARY KEY, recorded_at INTEGER NOT NULL, conversation_hash TEXT, data TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS ${table}_time ON ${table}(recorded_at DESC, id DESC); CREATE INDEX IF NOT EXISTS ${table}_conversation ON ${table}(conversation_hash, recorded_at DESC);`) } // Older bridge versions may have created the usage database before // conversation snapshots existed. Add columns defensively rather than // replacing the private state file and losing routing/usage history. const snapshotColumns = new Set( (db.prepare("PRAGMA table_info(conversation_snapshots)").all() as Array<{ name: string }>).map( (column) => column.name, ), ) const columns: Array<[string, string]> = [ ["service_tier", "TEXT"], ["last_context_tokens", "INTEGER"], ["generation", "INTEGER NOT NULL DEFAULT 0"], ["last_activity_at", "INTEGER NOT NULL DEFAULT 0"], ["idle_due_at", "INTEGER"], ["last_compacted_at", "INTEGER"], ["lease_owner", "TEXT"], ["lease_expires_at", "INTEGER"], ["last_compaction_error", "TEXT"], ] for (const [name, definition] of columns) { if (!snapshotColumns.has(name)) db.exec(`ALTER TABLE conversation_snapshots ADD COLUMN ${name} ${definition}`) } const usageColumns = new Set( (db.prepare("PRAGMA table_info(token_usage)").all() as Array<{ name: string }>).map((column) => column.name), ) if (!usageColumns.has("account_id")) db.exec("ALTER TABLE token_usage ADD COLUMN account_id TEXT") if (!usageColumns.has("plan_primary_delta")) db.exec("ALTER TABLE token_usage ADD COLUMN plan_primary_delta REAL") if (!usageColumns.has("plan_secondary_delta")) db.exec("ALTER TABLE token_usage ADD COLUMN plan_secondary_delta REAL") db.exec("CREATE INDEX IF NOT EXISTS token_usage_account_time ON token_usage(account_id, recorded_at)") for (const sidecar of [dbPath, `${dbPath}-wal`, `${dbPath}-shm`]) { try { fs.chmodSync(sidecar, 0o600) } catch { /* sidecar may not exist yet */ } } const value = { db, path: dbPath } databases.set(dbPath, value) return value } const UPSTREAM_ERROR_RETENTION_MS = 7 * DAY_MS const UPSTREAM_ERROR_MAX_ROWS = 1000 const UPSTREAM_FIELD_LIMIT = 240 const UPSTREAM_MESSAGE_LIMIT = 1000 function safeUpstreamText(value: unknown, limit: number): string { const text = typeof value === "string" ? value : value instanceof Error ? value.message : String(value ?? "") return text .replace(/[\u0000-\u001f\u007f]/g, " ") .replace(/\s+/g, " ") .trim() .slice(0, limit) } function safeUpstreamMessage(value: unknown): string { let text = safeUpstreamText(value, UPSTREAM_MESSAGE_LIMIT) text = text.replace(/bearer\s+[a-z0-9._~+/=-]+/gi, "Bearer [redacted]") text = text.replace( /((?:api[_ -]?key|access[_ -]?token|refresh[_ -]?token|authorization)["'=: ]+)[^,;\s"'}]+/gi, "$1[redacted]", ) return text || "upstream request failed" } function pruneUpstreamErrors(db: DatabaseSync, now: number): void { db.prepare("DELETE FROM upstream_errors WHERE recorded_at < ?").run(now - UPSTREAM_ERROR_RETENTION_MS) db.exec(`DELETE FROM upstream_errors WHERE id IN ( SELECT id FROM upstream_errors ORDER BY recorded_at DESC, id DESC LIMIT -1 OFFSET ${UPSTREAM_ERROR_MAX_ROWS} )`) } export function recordUpstreamError(stateDirectory: string, error: UpstreamErrorRecord): void { const { db } = getDatabase(stateDirectory) const recordedAt = Number.isFinite(error.timestamp) ? Math.trunc(error.timestamp!) : Date.now() db.prepare( `INSERT INTO upstream_errors (id, recorded_at, provider, model, account_id, route, status, error_type, message, request_id, correlation_id) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, ).run( (error.id && safeUpstreamText(error.id, UPSTREAM_FIELD_LIMIT)) || randomUUID(), recordedAt, safeUpstreamText(error.provider, UPSTREAM_FIELD_LIMIT) || "unknown", safeUpstreamText(error.model, UPSTREAM_FIELD_LIMIT) || "unknown", error.accountId ? safeUpstreamText(error.accountId, UPSTREAM_FIELD_LIMIT) : null, safeUpstreamText(error.route, UPSTREAM_FIELD_LIMIT) || "unknown", Number.isFinite(error.status) ? Math.trunc(error.status) : 502, safeUpstreamText(error.errorType, UPSTREAM_FIELD_LIMIT) || "upstream_error", safeUpstreamMessage(error.message), error.requestId ? safeUpstreamText(error.requestId, UPSTREAM_FIELD_LIMIT) : null, error.correlationId ? safeUpstreamText(error.correlationId, UPSTREAM_FIELD_LIMIT) : null, ) pruneUpstreamErrors(db, recordedAt) } export function listUpstreamErrors( stateDirectory: string, limit = 100, offset = 0, now = Date.now(), ): ListedUpstreamErrors { const { db } = getDatabase(stateDirectory) pruneUpstreamErrors(db, now) const boundedLimit = Math.min(500, Math.max(1, Math.trunc(limit) || 100)) const boundedOffset = Math.max(0, Math.trunc(offset) || 0) const total = Number((db.prepare("SELECT COUNT(*) AS count FROM upstream_errors").get() as { count: number }).count) const rows = db .prepare( `SELECT id, recorded_at, provider, model, account_id, route, status, error_type, message, request_id, correlation_id FROM upstream_errors ORDER BY recorded_at DESC, id DESC LIMIT ? OFFSET ?`, ) .all(boundedLimit, boundedOffset) as Array> const errors = rows.map((row) => { const result: UpstreamErrorRecord & { id: string; timestamp: number } = { id: String(row.id), timestamp: Number(row.recorded_at), provider: String(row.provider), model: String(row.model), route: String(row.route), status: Number(row.status), errorType: String(row.error_type), message: String(row.message), } if (row.account_id !== null) result.accountId = String(row.account_id) if (row.request_id !== null) result.requestId = String(row.request_id) if (row.correlation_id !== null) result.correlationId = String(row.correlation_id) return result }) return { total, errors, limit: boundedLimit, offset: boundedOffset, hasMore: boundedOffset + errors.length < total } } export function recordTokenUsage(stateDirectory: string, usage: TokenUsage): void { if (!usage.apiKey || !Number.isFinite(usage.totalTokens)) return const inputTokens = Math.max(0, Math.trunc(usage.inputTokens)) const outputTokens = Math.max(0, Math.trunc(usage.outputTokens)) const totalTokens = Math.max(0, Math.trunc(usage.totalTokens)) const cacheHitTokens = Math.max(0, Math.min(inputTokens, Math.trunc(usage.cachedTokens ?? 0))) const cacheWriteTokens = Math.max(0, Math.trunc(usage.cacheWriteTokens ?? 0)) const { db } = getDatabase(stateDirectory) db.prepare( `INSERT INTO token_usage (api_key, account_id, provider, model, input_tokens, output_tokens, total_tokens, cached_tokens, cache_write_tokens, plan_primary_delta, plan_secondary_delta, recorded_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, ).run( usage.apiKey, usage.accountId ?? null, usage.provider, usage.model, inputTokens, outputTokens, totalTokens, cacheHitTokens, cacheWriteTokens, Number.isFinite(usage.planPrimaryDelta) ? Math.max(0, usage.planPrimaryDelta!) : null, Number.isFinite(usage.planSecondaryDelta) ? Math.max(0, usage.planSecondaryDelta!) : null, usage.recordedAt ?? Date.now(), ) } /** Aggregate only requests for which two comparable Codex snapshots existed. */ export function getModelObservedUsage(stateDirectory: string): Map { const { db } = getDatabase(stateDirectory) const rows = db .prepare( `SELECT model, COUNT(plan_primary_delta) AS requests, COALESCE(SUM(CASE WHEN plan_primary_delta IS NOT NULL THEN total_tokens ELSE 0 END), 0) AS total_tokens, COALESCE(SUM(plan_primary_delta), 0) AS plan_percent_delta FROM token_usage WHERE provider = 'codex' AND plan_primary_delta IS NOT NULL GROUP BY model`, ) .all() as Array<{ model: string requests: number total_tokens: number plan_percent_delta: number }> return new Map( rows.map((row) => { const totalTokens = Number(row.total_tokens) const planPercentDelta = Number(row.plan_percent_delta) return [ row.model, { requests: Number(row.requests), totalTokens, planPercentDelta, percentPer1MTokens: totalTokens > 0 ? (planPercentDelta * 1_000_000) / totalTokens : null, }, ] }), ) } export function getTokenUsage(stateDirectory: string, apiKey: string, since?: number): UsageTotals { const { db } = getDatabase(stateDirectory) const row = db .prepare( `SELECT COALESCE(SUM(input_tokens), 0) AS input_tokens, COALESCE(SUM(output_tokens), 0) AS output_tokens, COALESCE(SUM(total_tokens), 0) AS total_tokens, COUNT(*) AS requests FROM token_usage WHERE api_key = ? AND (? IS NULL OR recorded_at >= ?)`, ) .get(apiKey, since ?? null, since ?? null) as { input_tokens: number output_tokens: number total_tokens: number requests: number } return { inputTokens: Number(row.input_tokens), outputTokens: Number(row.output_tokens), totalTokens: Number(row.total_tokens), requests: Number(row.requests), } } export function getAccountTokenUsageHistory( stateDirectory: string, accountId: string, now = Date.now(), ): UsageHistoryPoint[] { const currentDay = Math.floor((Number.isFinite(now) ? now : Date.now()) / DAY_MS) * DAY_MS const start = currentDay - 29 * DAY_MS const { db } = getDatabase(stateDirectory) const rows = db .prepare( `SELECT CAST((recorded_at - ?) / ? AS INTEGER) AS bucket, COALESCE(SUM(total_tokens), 0) AS total_tokens, COALESCE(SUM(cached_tokens), 0) AS cache_hit_tokens, COALESCE(SUM(MAX(input_tokens - cached_tokens, 0)), 0) AS cache_miss_tokens, COUNT(*) AS requests FROM token_usage WHERE account_id = ? AND recorded_at >= ? AND recorded_at < ? GROUP BY bucket`, ) .all(start, DAY_MS, accountId, start, currentDay + DAY_MS) as Array<{ bucket: number total_tokens: number cache_hit_tokens: number cache_miss_tokens: number requests: number }> const byBucket = new Map(rows.map((row) => [Number(row.bucket), row])) return Array.from({ length: 30 }, (_, index) => { const row = byBucket.get(index) return { at: start + index * DAY_MS, totalTokens: row ? Number(row.total_tokens) : 0, cacheHitTokens: row ? Number(row.cache_hit_tokens) : 0, cacheMissTokens: row ? Number(row.cache_miss_tokens) : 0, requests: row ? Number(row.requests) : 0, } }) } export function getAllTokenUsage(stateDirectory: string): Map { const { db } = getDatabase(stateDirectory) const rows = db .prepare( `SELECT api_key, COALESCE(SUM(input_tokens), 0) AS input_tokens, COALESCE(SUM(output_tokens), 0) AS output_tokens, COALESCE(SUM(total_tokens), 0) AS total_tokens, COUNT(*) AS requests FROM token_usage GROUP BY api_key`, ) .all() as Array<{ api_key: string input_tokens: number output_tokens: number total_tokens: number requests: number }> return new Map( rows.map((row) => [ row.api_key, { inputTokens: Number(row.input_tokens), outputTokens: Number(row.output_tokens), totalTokens: Number(row.total_tokens), requests: Number(row.requests), }, ]), ) } export function closeTokenUsageDatabases(): void { for (const { db } of databases.values()) { db.close() } databases.clear() } function routeParams(route: ConversationRouteKey): string[] { return [route.apiKeyId, route.provider, route.protocol, route.conversationHash] } function snapshotRow(row: any): ConversationSnapshot { return { route: { apiKeyId: String(row.api_key_id), provider: String(row.provider), protocol: String(row.protocol), conversationHash: String(row.conversation_hash), }, accountId: String(row.account_id), model: String(row.model), ...(row.service_tier === "priority" || row.service_tier === "fast" ? { serviceTier: row.service_tier } : {}), callerHistory: String(row.caller_history), upstreamHistory: String(row.upstream_history), fixedContextDigest: String(row.fixed_context_digest), lastContextTokens: row.last_context_tokens === null || row.last_context_tokens === undefined ? null : Number(row.last_context_tokens), generation: Number(row.generation), lastActivityAt: Number(row.last_activity_at), idleDueAt: row.idle_due_at === null || row.idle_due_at === undefined ? null : Number(row.idle_due_at), lastCompactedAt: row.last_compacted_at === null || row.last_compacted_at === undefined ? null : Number(row.last_compacted_at), leaseOwner: row.lease_owner === null || row.lease_owner === undefined ? null : String(row.lease_owner), leaseExpiresAt: row.lease_expires_at === null || row.lease_expires_at === undefined ? null : Number(row.lease_expires_at), lastCompactionError: row.last_compaction_error === null || row.last_compaction_error === undefined ? null : String(row.last_compaction_error), } } export function lookupConversationSnapshot( stateDirectory: string, route: ConversationRouteKey, ): ConversationSnapshot | null { const { db } = getDatabase(stateDirectory) const row = db .prepare( `SELECT * FROM conversation_snapshots WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ) .get(...snapshotParams(route)) return row ? snapshotRow(row) : null } export function listDueConversationSnapshots(stateDirectory: string, now = Date.now()): ConversationSnapshot[] { const { db } = getDatabase(stateDirectory) const rows = db .prepare( `SELECT * FROM conversation_snapshots WHERE idle_due_at IS NOT NULL AND idle_due_at <= ? AND (lease_expires_at IS NULL OR lease_expires_at <= ?) ORDER BY idle_due_at ASC`, ) .all(now, now) as any[] return rows.map(snapshotRow) } export function saveConversationSnapshot( stateDirectory: string, snapshot: ConversationSnapshot, retentionDays: number, ): void { const { db } = getDatabase(stateDirectory) const now = Date.now() const cutoff = now - Math.max(1, Math.trunc(retentionDays)) * 86_400_000 db.exec("BEGIN IMMEDIATE") try { db.prepare( `INSERT INTO conversation_snapshots (api_key_id, provider, protocol, conversation_hash, account_id, model, service_tier, caller_history, upstream_history, fixed_context_digest, last_context_tokens, generation, last_activity_at, idle_due_at, last_compacted_at, lease_owner, lease_expires_at, last_compaction_error) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(api_key_id, provider, protocol, conversation_hash) DO UPDATE SET account_id = excluded.account_id, model = excluded.model, service_tier = excluded.service_tier, caller_history = excluded.caller_history, upstream_history = excluded.upstream_history, fixed_context_digest = excluded.fixed_context_digest, last_context_tokens = excluded.last_context_tokens, generation = excluded.generation, last_activity_at = excluded.last_activity_at, idle_due_at = excluded.idle_due_at, last_compacted_at = excluded.last_compacted_at, lease_owner = excluded.lease_owner, lease_expires_at = excluded.lease_expires_at, last_compaction_error = excluded.last_compaction_error`, ).run( ...snapshotParams(snapshot.route), snapshot.accountId, snapshot.model, snapshot.serviceTier ?? null, snapshot.callerHistory, snapshot.upstreamHistory, snapshot.fixedContextDigest, snapshot.lastContextTokens, snapshot.generation, snapshot.lastActivityAt, snapshot.idleDueAt, snapshot.lastCompactedAt, snapshot.leaseOwner, snapshot.leaseExpiresAt, snapshot.lastCompactionError, ) db.prepare( `DELETE FROM conversation_snapshots WHERE api_key_id = ? AND last_activity_at < ?`, ).run(snapshot.route.apiKeyId, cutoff) db.prepare( `DELETE FROM conversation_routes WHERE api_key_id = ? AND last_seen_at < ? AND NOT EXISTS ( SELECT 1 FROM conversation_snapshots s WHERE s.api_key_id = conversation_routes.api_key_id AND s.provider = conversation_routes.provider AND s.protocol = conversation_routes.protocol AND s.conversation_hash = conversation_routes.conversation_hash )`, ).run(snapshot.route.apiKeyId, cutoff) db.exec("COMMIT") } catch (error) { try { db.exec("ROLLBACK") } catch { /* preserve original failure */ } throw error } } export function acquireConversationLease( stateDirectory: string, route: ConversationRouteKey, owner: string, ttlMs = 30_000, ): ConversationLease { const { db } = getDatabase(stateDirectory) const now = Date.now() db.exec("BEGIN IMMEDIATE") try { const row = db .prepare( `SELECT generation, lease_owner, lease_expires_at FROM conversation_snapshots WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ) .get(...snapshotParams(route)) as any if (!row) { db.exec("COMMIT") return { acquired: false, generation: 0, leaseOwner: null } } const active = row.lease_owner && Number(row.lease_expires_at ?? 0) > now && row.lease_owner !== owner if (active) { db.exec("COMMIT") return { acquired: false, generation: Number(row.generation), leaseOwner: String(row.lease_owner) } } db.prepare( `UPDATE conversation_snapshots SET lease_owner = ?, lease_expires_at = ? WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ).run(owner, now + ttlMs, ...snapshotParams(route)) db.exec("COMMIT") return { acquired: true, generation: Number(row.generation), leaseOwner: owner } } catch (error) { try { db.exec("ROLLBACK") } catch { /* preserve original failure */ } throw error } } export function releaseConversationLease( stateDirectory: string, route: ConversationRouteKey, owner: string, error?: string, ): void { const { db } = getDatabase(stateDirectory) db.prepare( `UPDATE conversation_snapshots SET lease_owner = NULL, lease_expires_at = NULL, idle_due_at = CASE WHEN ? IS NULL THEN idle_due_at ELSE NULL END, last_compaction_error = ? WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ? AND lease_owner = ?`, ).run(error ?? null, error ?? null, ...snapshotParams(route), owner) } export function commitConversationTransition( stateDirectory: string, route: ConversationRouteKey, expectedAccountId: string, expectedGeneration: number, snapshot: ConversationSnapshot, destinationAccountId: string, leaseOwner?: string, ): boolean { const { db } = getDatabase(stateDirectory) const now = Date.now() db.exec("BEGIN IMMEDIATE") try { const routeRow = db .prepare( `SELECT account_id FROM conversation_routes WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ) .get(...routeParams(route)) as { account_id: string } | undefined const snapshotRowValue = db .prepare( `SELECT generation, lease_owner FROM conversation_snapshots WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ) .get(...routeParams(route)) as { generation: number; lease_owner: string | null } | undefined const ownsLease = !leaseOwner || snapshotRowValue?.lease_owner === leaseOwner if ( !routeRow || routeRow.account_id !== expectedAccountId || !snapshotRowValue || Number(snapshotRowValue.generation) !== expectedGeneration || !ownsLease ) { db.exec("ROLLBACK") return false } const nextGeneration = expectedGeneration + 1 db.prepare( `UPDATE conversation_routes SET account_id = ?, last_seen_at = ? WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ).run(destinationAccountId, now, ...routeParams(route)) db.prepare( `UPDATE conversation_snapshots SET account_id = ?, model = ?, service_tier = ?, caller_history = ?, upstream_history = ?, fixed_context_digest = ?, last_context_tokens = ?, generation = ?, last_activity_at = ?, idle_due_at = ?, last_compacted_at = ?, lease_owner = NULL, lease_expires_at = NULL, last_compaction_error = NULL WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ? AND generation = ? AND (? IS NULL OR lease_owner = ?)`, ).run( destinationAccountId, snapshot.model, snapshot.serviceTier ?? null, snapshot.callerHistory, snapshot.upstreamHistory, snapshot.fixedContextDigest, snapshot.lastContextTokens, nextGeneration, snapshot.lastActivityAt, snapshot.idleDueAt, snapshot.lastCompactedAt, ...routeParams(route), expectedGeneration, leaseOwner ?? null, leaseOwner ?? null, ) db.exec("COMMIT") return true } catch (error) { try { db.exec("ROLLBACK") } catch { /* preserve original failure */ } throw error } } export function commitConversationSnapshot( stateDirectory: string, route: ConversationRouteKey, expectedAccountId: string, expectedGeneration: number, snapshot: ConversationSnapshot, leaseOwner?: string, ): boolean { const { db } = getDatabase(stateDirectory) const now = Date.now() db.exec("BEGIN IMMEDIATE") try { const row = db .prepare( `SELECT account_id, generation, lease_owner FROM conversation_snapshots WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ) .get(...routeParams(route)) as { account_id: string; generation: number; lease_owner: string | null } | undefined if ( !row || row.account_id !== expectedAccountId || Number(row.generation) !== expectedGeneration || (leaseOwner !== undefined && row.lease_owner !== leaseOwner) ) { db.exec("ROLLBACK") return false } db.prepare( `UPDATE conversation_snapshots SET model = ?, service_tier = ?, caller_history = ?, upstream_history = ?, fixed_context_digest = ?, last_context_tokens = ?, generation = ?, last_activity_at = ?, idle_due_at = ?, last_compacted_at = ?, lease_owner = NULL, lease_expires_at = NULL, last_compaction_error = NULL WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ? AND generation = ? AND (? IS NULL OR lease_owner = ?)`, ).run( snapshot.model, snapshot.serviceTier ?? null, snapshot.callerHistory, snapshot.upstreamHistory, snapshot.fixedContextDigest, snapshot.lastContextTokens, expectedGeneration + 1, snapshot.lastActivityAt, snapshot.idleDueAt, snapshot.lastCompactedAt, ...routeParams(route), expectedGeneration, leaseOwner ?? null, leaseOwner ?? null, ) db.prepare( `UPDATE conversation_routes SET last_seen_at = ? WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ).run(now, ...routeParams(route)) db.exec("COMMIT") return true } catch (error) { try { db.exec("ROLLBACK") } catch { /* preserve original failure */ } throw error } } export function lookupConversationRoute( stateDirectory: string, route: ConversationRouteKey, touch = true, ): { accountId: string; lastSeenAt: number } | null { const { db } = getDatabase(stateDirectory) const row = db .prepare( `SELECT account_id, last_seen_at FROM conversation_routes WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ) .get(...routeParams(route)) as { account_id: string; last_seen_at: number } | undefined if (!row) return null const now = Date.now() if (touch) { db.prepare( `UPDATE conversation_routes SET last_seen_at = ? WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ).run(now, ...routeParams(route)) } return { accountId: row.account_id, lastSeenAt: touch ? now : row.last_seen_at } } export function claimConversationRoute(stateDirectory: string, route: ConversationRouteKey, accountId: string): string { const { db } = getDatabase(stateDirectory) const now = Date.now() db.exec("BEGIN IMMEDIATE") try { db.prepare( `INSERT INTO conversation_routes (api_key_id, provider, protocol, conversation_hash, account_id, last_seen_at) VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT DO NOTHING`, ).run(...routeParams(route), accountId, now) const row = db .prepare( `SELECT account_id FROM conversation_routes WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ) .get(...routeParams(route)) as { account_id: string } db.prepare( `UPDATE conversation_routes SET last_seen_at = ? WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ).run(now, ...routeParams(route)) db.exec("COMMIT") return row.account_id } catch (error) { try { db.exec("ROLLBACK") } catch { /* preserve original failure */ } throw error } } export function claimConversationRouteWithCursor( stateDirectory: string, route: ConversationRouteKey, poolAccountIds: string[], eligibleAccountIds: Set, ): { accountId: string; inserted: boolean } { if (poolAccountIds.length === 0) throw new Error("cannot route an empty account pool") const { db } = getDatabase(stateDirectory) const now = Date.now() db.exec("BEGIN IMMEDIATE") try { const existing = db .prepare( `SELECT account_id FROM conversation_routes WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ) .get(...routeParams(route)) as { account_id: string } | undefined if (existing) { db.prepare( `UPDATE conversation_routes SET last_seen_at = ? WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ).run(now, ...routeParams(route)) db.exec("COMMIT") return { accountId: existing.account_id, inserted: false } } const row = db .prepare("SELECT next_index FROM routing_cursors WHERE api_key_id = ? AND provider = ?") .get(route.apiKeyId, route.provider) as { next_index: number } | undefined const start = ((Number(row?.next_index ?? 0) % poolAccountIds.length) + poolAccountIds.length) % poolAccountIds.length const offset = poolAccountIds.findIndex((accountId, index) => eligibleAccountIds.has(accountId) && index >= start) const wrapped = poolAccountIds.findIndex((accountId) => eligibleAccountIds.has(accountId)) const selectedIndex = offset >= 0 ? offset : wrapped if (selectedIndex < 0) throw new Error("cannot route without an eligible account") const accountId = poolAccountIds[selectedIndex]! db.prepare( `INSERT INTO conversation_routes (api_key_id, provider, protocol, conversation_hash, account_id, last_seen_at) VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT DO NOTHING`, ).run(...routeParams(route), accountId, now) db.prepare( `INSERT INTO routing_cursors(api_key_id, provider, next_index) VALUES (?, ?, ?) ON CONFLICT(api_key_id, provider) DO UPDATE SET next_index = excluded.next_index`, ).run(route.apiKeyId, route.provider, (selectedIndex + 1) % poolAccountIds.length) db.exec("COMMIT") return { accountId, inserted: true } } catch (error) { try { db.exec("ROLLBACK") } catch { /* preserve original failure */ } throw error } } export function advanceEligibleRoutingCursor( stateDirectory: string, apiKeyId: string, provider: string, poolAccountIds: string[], eligibleAccountIds: Set, ): string { if (poolAccountIds.length === 0) throw new Error("cannot route an empty account pool") const { db } = getDatabase(stateDirectory) db.exec("BEGIN IMMEDIATE") try { const row = db .prepare("SELECT next_index FROM routing_cursors WHERE api_key_id = ? AND provider = ?") .get(apiKeyId, provider) as { next_index: number } | undefined const start = ((Number(row?.next_index ?? 0) % poolAccountIds.length) + poolAccountIds.length) % poolAccountIds.length const offset = poolAccountIds.findIndex((accountId, index) => eligibleAccountIds.has(accountId) && index >= start) const wrapped = poolAccountIds.findIndex((accountId) => eligibleAccountIds.has(accountId)) const selectedIndex = offset >= 0 ? offset : wrapped if (selectedIndex < 0) throw new Error("cannot route without an eligible account") const accountId = poolAccountIds[selectedIndex]! db.prepare( `INSERT INTO routing_cursors(api_key_id, provider, next_index) VALUES (?, ?, ?) ON CONFLICT(api_key_id, provider) DO UPDATE SET next_index = excluded.next_index`, ).run(apiKeyId, provider, (selectedIndex + 1) % poolAccountIds.length) db.exec("COMMIT") return accountId } catch (error) { try { db.exec("ROLLBACK") } catch { /* preserve original failure */ } throw error } } export function invalidateConversationRoute( stateDirectory: string, route: ConversationRouteKey, accountId?: string, ): void { const { db } = getDatabase(stateDirectory) if (accountId) { db.prepare( `DELETE FROM conversation_routes WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ? AND account_id = ?`, ).run(...routeParams(route), accountId) } else { db.prepare( `DELETE FROM conversation_routes WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ?`, ).run(...routeParams(route)) } } export function moveConversationRoute( stateDirectory: string, route: ConversationRouteKey, from: string, to: string, ): void { const { db } = getDatabase(stateDirectory) db.prepare( `UPDATE conversation_routes SET account_id = ?, last_seen_at = ? WHERE api_key_id = ? AND provider = ? AND protocol = ? AND conversation_hash = ? AND account_id = ?`, ).run(to, Date.now(), ...routeParams(route), from) } export function advanceRoutingCursor( stateDirectory: string, apiKeyId: string, provider: string, poolLength: number, ): number { if (poolLength <= 0) return 0 const { db } = getDatabase(stateDirectory) db.exec("BEGIN IMMEDIATE") try { const row = db .prepare("SELECT next_index FROM routing_cursors WHERE api_key_id = ? AND provider = ?") .get(apiKeyId, provider) as { next_index: number } | undefined const index = ((Number(row?.next_index ?? 0) % poolLength) + poolLength) % poolLength db.prepare( `INSERT INTO routing_cursors(api_key_id, provider, next_index) VALUES (?, ?, ?) ON CONFLICT(api_key_id, provider) DO UPDATE SET next_index = excluded.next_index`, ).run(apiKeyId, provider, (index + 1) % poolLength) db.exec("COMMIT") return index } catch (error) { try { db.exec("ROLLBACK") } catch { /* preserve original failure */ } throw error } } export function removeRoutingStateForKey(stateDirectory: string, apiKeyId: string): void { const { db } = getDatabase(stateDirectory) db.prepare("DELETE FROM conversation_snapshots WHERE api_key_id = ?").run(apiKeyId) db.prepare("DELETE FROM conversation_routes WHERE api_key_id = ?").run(apiKeyId) db.prepare("DELETE FROM routing_cursors WHERE api_key_id = ?").run(apiKeyId) } export function removeRoutingStateForAccount(_stateDirectory: string, _accountId: string): void { // Preserve prior-account provenance and private history until migration or TTL expiry. } export function peekRoutingCursor( stateDirectory: string, apiKeyId: string, provider: string, poolLength: number, ): number { if (poolLength <= 0) return 0 const { db } = getDatabase(stateDirectory) const row = db .prepare("SELECT next_index FROM routing_cursors WHERE api_key_id = ? AND provider = ?") .get(apiKeyId, provider) as { next_index: number } | undefined return ((Number(row?.next_index ?? 0) % poolLength) + poolLength) % poolLength } export function removeConversationRoutesForKeyProvider( _stateDirectory: string, _apiKeyId: string, _provider: string, _accountIds: Set, ): void { // A removed pool entry is an account transition, not permission to discard history. } export type DiagnosticTable = "request_metrics" | "compaction_events" export function writeDiagnostic(stateDirectory: string, table: DiagnosticTable, data: Record): void { try { getDatabase(stateDirectory) .db.prepare( `INSERT INTO ${table} (id, recorded_at, conversation_hash, data) VALUES (?, ?, ?, ?) ON CONFLICT(id) DO UPDATE SET data = excluded.data`, ) .run(data.id, data.timestamp, data.conversation_hash ?? null, JSON.stringify(data)) } catch { console.error(`could not persist ${table}`) } } export function listDiagnostics( stateDirectory: string, table: DiagnosticTable, limit = 100, offset = 0, conversationHash?: string, ) { limit = Number.isFinite(limit) ? Math.min(1000, Math.max(1, Math.trunc(limit))) : 100 offset = Number.isFinite(offset) ? Math.max(0, Math.trunc(offset)) : 0 const { db } = getDatabase(stateDirectory) const where = conversationHash ? "WHERE conversation_hash = ?" : "" const args = conversationHash ? [conversationHash] : [] const total = Number( (db.prepare(`SELECT COUNT(*) AS total FROM ${table} ${where}`).get(...args) as { total: number }).total, ) const rows = db .prepare(`SELECT data FROM ${table} ${where} ORDER BY recorded_at DESC, id DESC LIMIT ? OFFSET ?`) .all(...args, limit, offset) as { data: string }[] return { records: rows.map((row) => JSON.parse(row.data)), total, limit, offset, hasMore: offset + rows.length < total, } } export function latestCompactionEvent(stateDirectory: string, route: ConversationRouteKey): string | null { const row = getDatabase(stateDirectory) .db.prepare( `SELECT id FROM compaction_events WHERE conversation_hash = ? AND json_extract(data, '$.route.apiKeyId') = ? AND json_extract(data, '$.route.protocol') = ? AND json_extract(data, '$.route.provider') = ? ORDER BY recorded_at DESC, id DESC LIMIT 1`, ) .get(route.conversationHash, route.apiKeyId, route.protocol, route.provider) as { id: string } | undefined return row?.id ?? null }