diff --git a/open-sse/config/audioRegistry.ts b/open-sse/config/audioRegistry.ts index dc84b12305..338b6fb485 100644 --- a/open-sse/config/audioRegistry.ts +++ b/open-sse/config/audioRegistry.ts @@ -22,11 +22,7 @@ export interface AudioProvider { models: AudioModel[]; } -let _AUDIO_TRANSCRIPTION_PROVIDERS: Record | null = null; - -function getOrCreateTranscriptionProviders(): Record { - if (!_AUDIO_TRANSCRIPTION_PROVIDERS) { - _AUDIO_TRANSCRIPTION_PROVIDERS = { +export const AUDIO_TRANSCRIPTION_PROVIDERS: Record = { openai: { id: "openai", baseUrl: "https://api.openai.com/v1/audio/transcriptions", @@ -149,56 +145,9 @@ function getOrCreateTranscriptionProviders(): Record { { id: "elevenlabs/audio-isolation", name: "ElevenLabs Audio Isolation" }, ], }, - }; -} - return _AUDIO_TRANSCRIPTION_PROVIDERS; -} +}; -export const AUDIO_TRANSCRIPTION_PROVIDERS: Record = new Proxy({} as Record, { - get(target, key: string) { - if (key in target) { - return target[key]; - } - return getOrCreateTranscriptionProviders()[key]; - }, - set(target, key: string, value) { - target[key] = value; - getOrCreateTranscriptionProviders()[key] = value; - return true; - }, - deleteProperty(target, key: string) { - delete target[key]; - delete getOrCreateTranscriptionProviders()[key]; - return true; - }, - ownKeys(target) { - const targetKeys = Reflect.ownKeys(target); - const registryKeys = Reflect.ownKeys(getOrCreateTranscriptionProviders()); - return Array.from(new Set([...targetKeys, ...registryKeys])); - }, - has(target, key) { - return key in target || key in getOrCreateTranscriptionProviders(); - }, - getOwnPropertyDescriptor(target, key) { - if (key in target) { - return Reflect.getOwnPropertyDescriptor(target, key); - } - if (key in getOrCreateTranscriptionProviders()) { - return { configurable: true, enumerable: true, value: getOrCreateTranscriptionProviders()[key as string] }; - } - return undefined; - }, -}); - -export function getTranscriptionProviders(): Record { - return AUDIO_TRANSCRIPTION_PROVIDERS; -} - -let _AUDIO_SPEECH_PROVIDERS: Record | null = null; - -function getOrCreateSpeechProviders(): Record { - if (!_AUDIO_SPEECH_PROVIDERS) { - _AUDIO_SPEECH_PROVIDERS = { +export const AUDIO_SPEECH_PROVIDERS: Record = { openai: { id: "openai", baseUrl: "https://api.openai.com/v1/audio/speech", @@ -418,57 +367,22 @@ function getOrCreateSpeechProviders(): Record { { id: "mimo-v2.5-tts-voiceclone", name: "MiMo V2.5 Voice Clone" }, ], }, - }; -} - return _AUDIO_SPEECH_PROVIDERS; -} - -export const AUDIO_SPEECH_PROVIDERS: Record = new Proxy({} as Record, { - get(target, key: string) { - if (key in target) { - return target[key]; - } - return getOrCreateSpeechProviders()[key]; - }, - set(target, key: string, value) { - target[key] = value; - getOrCreateSpeechProviders()[key] = value; - return true; - }, - deleteProperty(target, key: string) { - delete target[key]; - delete getOrCreateSpeechProviders()[key]; - return true; - }, - ownKeys(target) { - const targetKeys = Reflect.ownKeys(target); - const registryKeys = Reflect.ownKeys(getOrCreateSpeechProviders()); - return Array.from(new Set([...targetKeys, ...registryKeys])); - }, - has(target, key) { - return key in target || key in getOrCreateSpeechProviders(); - }, - getOwnPropertyDescriptor(target, key) { - if (key in target) { - return Reflect.getOwnPropertyDescriptor(target, key); - } - if (key in getOrCreateSpeechProviders()) { - return { configurable: true, enumerable: true, value: getOrCreateSpeechProviders()[key as string] }; - } - return undefined; - }, -}); - -export function getSpeechProviders(): Record { - return AUDIO_SPEECH_PROVIDERS; -} +}; +/** + * Get transcription provider config by ID + */ export function getTranscriptionProvider(providerId: string): AudioProvider | null { return AUDIO_TRANSCRIPTION_PROVIDERS[providerId] || null; } + +/** + * Get speech provider config by ID + */ export function getSpeechProvider(providerId: string): AudioProvider | null { return AUDIO_SPEECH_PROVIDERS[providerId] || null; } + export interface ProviderNodeRow { prefix: string; name: string; diff --git a/open-sse/config/embeddingRegistry.ts b/open-sse/config/embeddingRegistry.ts index 1053bb6f02..69994c323a 100644 --- a/open-sse/config/embeddingRegistry.ts +++ b/open-sse/config/embeddingRegistry.ts @@ -45,11 +45,7 @@ export function buildDynamicEmbeddingProvider(node: EmbeddingProviderNodeRow): E }; } -let _EMBEDDING_PROVIDERS: Record | null = null; - -function getOrCreateEmbeddingProviders(): Record { - if (!_EMBEDDING_PROVIDERS) { - _EMBEDDING_PROVIDERS = { +export const EMBEDDING_PROVIDERS: Record = { cohere: { id: "cohere", baseUrl: "https://api.cohere.com/v2/embed", @@ -237,50 +233,7 @@ function getOrCreateEmbeddingProviders(): Record { { id: "jina-colbert-v2", name: "Jina ColBERT v2", dimensions: 128 }, ], }, - }; - } - return _EMBEDDING_PROVIDERS; -} - -export const EMBEDDING_PROVIDERS: Record = new Proxy({} as Record, { - get(target, key: string) { - if (key in target) { - return target[key]; - } - return getOrCreateEmbeddingProviders()[key]; - }, - set(target, key: string, value) { - target[key] = value; - getOrCreateEmbeddingProviders()[key] = value; - return true; - }, - deleteProperty(target, key: string) { - delete target[key]; - delete getOrCreateEmbeddingProviders()[key]; - return true; - }, - ownKeys(target) { - const targetKeys = Reflect.ownKeys(target); - const registryKeys = Reflect.ownKeys(getOrCreateEmbeddingProviders()); - return Array.from(new Set([...targetKeys, ...registryKeys])); - }, - has(target, key) { - return key in target || key in getOrCreateEmbeddingProviders(); - }, - getOwnPropertyDescriptor(target, key) { - if (key in target) { - return Reflect.getOwnPropertyDescriptor(target, key); - } - if (key in getOrCreateEmbeddingProviders()) { - return { configurable: true, enumerable: true, value: getOrCreateEmbeddingProviders()[key as string] }; - } - return undefined; - }, -}); - -export function getEmbeddingProviders(): Record { - return EMBEDDING_PROVIDERS; -} +}; const EMBEDDING_PROVIDER_ALIASES: Record = { jina: "jina-ai", diff --git a/open-sse/config/imageRegistry.ts b/open-sse/config/imageRegistry.ts index d37aa2c509..a32ce55638 100644 --- a/open-sse/config/imageRegistry.ts +++ b/open-sse/config/imageRegistry.ts @@ -112,11 +112,7 @@ function findImageModelConfig(providerId, modelId) { return provider.models.find((model) => model.id === modelId) || null; } -let _IMAGE_PROVIDERS: Record | null = null; - -function getOrCreateImageProviders(): Record { - if (!_IMAGE_PROVIDERS) { - _IMAGE_PROVIDERS = { +export const IMAGE_PROVIDERS: Record = { openai: { id: "openai", baseUrl: "https://api.openai.com/v1/images/generations", @@ -544,53 +540,14 @@ function getOrCreateImageProviders(): Record { supportedSizes: ["1024x1024", "1024x1280", "1280x1024"], }, }; -} -return _IMAGE_PROVIDERS; -} - -export function getImageProviders(): Record { - return IMAGE_PROVIDERS; -} - -export const IMAGE_PROVIDERS = new Proxy({} as Record, { - get(target, key: string) { - if (key in target) { - return target[key]; - } - return getOrCreateImageProviders()[key]; - }, - set(target, key: string, value) { - target[key] = value; - getOrCreateImageProviders()[key] = value; - return true; - }, - deleteProperty(target, key: string) { - delete target[key]; - delete getOrCreateImageProviders()[key]; - return true; - }, - ownKeys(target) { - const targetKeys = Reflect.ownKeys(target); - const registryKeys = Reflect.ownKeys(getOrCreateImageProviders()); - return Array.from(new Set([...targetKeys, ...registryKeys])); - }, - has(target, key) { - return key in target || key in getOrCreateImageProviders(); - }, - getOwnPropertyDescriptor(target, key) { - if (key in target) { - return Reflect.getOwnPropertyDescriptor(target, key); - } - if (key in getOrCreateImageProviders()) { - return { configurable: true, enumerable: true, value: getOrCreateImageProviders()[key as string] }; - } - return undefined; - }, -}); +/** + * Get image provider config by ID + */ export function getImageProvider(providerId) { return IMAGE_PROVIDERS[providerId] || null; } + /** * Parse image model string (format: "provider/model") * Returns { provider, model } diff --git a/open-sse/config/moderationRegistry.ts b/open-sse/config/moderationRegistry.ts index 7ff3ff8db5..4a3f621998 100644 --- a/open-sse/config/moderationRegistry.ts +++ b/open-sse/config/moderationRegistry.ts @@ -5,11 +5,7 @@ * Follows OpenAI's moderation API format. */ -let _MODERATION_PROVIDERS: Record | null = null; - -function getOrCreateModerationProviders(): Record { - if (!_MODERATION_PROVIDERS) { - _MODERATION_PROVIDERS = { +export const MODERATION_PROVIDERS = { openai: { id: "openai", baseUrl: "https://api.openai.com/v1/moderations", @@ -20,55 +16,18 @@ function getOrCreateModerationProviders(): Record { { id: "text-moderation-latest", name: "Text Moderation Latest" }, ], }, - }; - } - return _MODERATION_PROVIDERS; -} - -export const MODERATION_PROVIDERS = new Proxy({} as Record, { - get(target, key: string) { - if (key in target) { - return target[key]; - } - return getOrCreateModerationProviders()[key]; - }, - set(target, key: string, value) { - target[key] = value; - getOrCreateModerationProviders()[key] = value; - return true; - }, - deleteProperty(target, key: string) { - delete target[key]; - delete getOrCreateModerationProviders()[key]; - return true; - }, - ownKeys(target) { - const targetKeys = Reflect.ownKeys(target); - const registryKeys = Reflect.ownKeys(getOrCreateModerationProviders()); - return Array.from(new Set([...targetKeys, ...registryKeys])); - }, - has(target, key) { - return key in target || key in getOrCreateModerationProviders(); - }, - getOwnPropertyDescriptor(target, key) { - if (key in target) { - return Reflect.getOwnPropertyDescriptor(target, key); - } - if (key in getOrCreateModerationProviders()) { - return { configurable: true, enumerable: true, value: getOrCreateModerationProviders()[key as string] }; - } - return undefined; - }, -}); - -export function getModerationProviders(): Record { - return MODERATION_PROVIDERS; -} +}; +/** + * Get moderation provider config by ID + */ export function getModerationProvider(providerId) { return MODERATION_PROVIDERS[providerId] || null; } +/** + * Parse moderation model string + */ export function parseModerationModel(modelStr) { if (!modelStr) return { provider: null, model: null }; @@ -87,6 +46,9 @@ export function parseModerationModel(modelStr) { return { provider: null, model: modelStr }; } +/** + * Get all moderation models as a flat list + */ export function getAllModerationModels() { const models = []; for (const [providerId, config] of Object.entries(MODERATION_PROVIDERS)) { diff --git a/open-sse/config/musicRegistry.ts b/open-sse/config/musicRegistry.ts index 492eb444bb..7ca68d5891 100644 --- a/open-sse/config/musicRegistry.ts +++ b/open-sse/config/musicRegistry.ts @@ -23,11 +23,7 @@ interface MusicProvider { models: MusicModel[]; } -let _MUSIC_PROVIDERS: Record | null = null; - -function getOrCreateMusicProviders(): Record { - if (!_MUSIC_PROVIDERS) { - _MUSIC_PROVIDERS = { +export const MUSIC_PROVIDERS: Record = { kie: { id: "kie", baseUrl: "https://api.kie.ai", @@ -86,59 +82,25 @@ function getOrCreateMusicProviders(): Record { { id: "musicgen-medium", name: "MusicGen Medium" }, ], }, - }; -} - return _MUSIC_PROVIDERS; -} - -export const MUSIC_PROVIDERS: Record = new Proxy({} as Record, { - get(target, key: string) { - if (key in target) { - return target[key]; - } - return getOrCreateMusicProviders()[key]; - }, - set(target, key: string, value) { - target[key] = value; - getOrCreateMusicProviders()[key] = value; - return true; - }, - deleteProperty(target, key: string) { - delete target[key]; - delete getOrCreateMusicProviders()[key]; - return true; - }, - ownKeys(target) { - const targetKeys = Reflect.ownKeys(target); - const registryKeys = Reflect.ownKeys(getOrCreateMusicProviders()); - return Array.from(new Set([...targetKeys, ...registryKeys])); - }, - has(target, key) { - return key in target || key in getOrCreateMusicProviders(); - }, - getOwnPropertyDescriptor(target, key) { - if (key in target) { - return Reflect.getOwnPropertyDescriptor(target, key); - } - if (key in getOrCreateMusicProviders()) { - return { configurable: true, enumerable: true, value: getOrCreateMusicProviders()[key as string] }; - } - return undefined; - }, -}); - -export function getMusicProviders(): Record { - return MUSIC_PROVIDERS; -} +}; +/** + * Get music provider config by ID + */ export function getMusicProvider(providerId: string): MusicProvider | null { return MUSIC_PROVIDERS[providerId] || null; } +/** + * Parse music model string (format: "provider/model" or just "model") + */ export function parseMusicModel(modelStr: string | null) { return parseModelFromRegistry(modelStr, MUSIC_PROVIDERS); } +/** + * Get all music models as a flat list + */ export function getAllMusicModels() { return getAllModelsFromRegistry(MUSIC_PROVIDERS); -} \ No newline at end of file +} diff --git a/open-sse/config/rerankRegistry.ts b/open-sse/config/rerankRegistry.ts index 431afcf6ca..bf14c0cc23 100644 --- a/open-sse/config/rerankRegistry.ts +++ b/open-sse/config/rerankRegistry.ts @@ -8,11 +8,7 @@ * keyed by provider ID (e.g. "cohere", "together"). */ -let _RERANK_PROVIDERS: Record | null = null; - -function getOrCreateRerankProviders(): Record { - if (!_RERANK_PROVIDERS) { - _RERANK_PROVIDERS = { +export const RERANK_PROVIDERS = { cohere: { id: "cohere", baseUrl: "https://api.cohere.com/v2/rerank", @@ -75,50 +71,7 @@ function getOrCreateRerankProviders(): Record { { id: "jina-reranker-m0", name: "Jina Reranker m0" }, ], }, - }; -} - return _RERANK_PROVIDERS; -} - -export const RERANK_PROVIDERS = new Proxy({} as Record, { - get(target, key: string) { - if (key in target) { - return target[key]; - } - return getOrCreateRerankProviders()[key]; - }, - set(target, key: string, value) { - target[key] = value; - getOrCreateRerankProviders()[key] = value; - return true; - }, - deleteProperty(target, key: string) { - delete target[key]; - delete getOrCreateRerankProviders()[key]; - return true; - }, - ownKeys(target) { - const targetKeys = Reflect.ownKeys(target); - const registryKeys = Reflect.ownKeys(getOrCreateRerankProviders()); - return Array.from(new Set([...targetKeys, ...registryKeys])); - }, - has(target, key) { - return key in target || key in getOrCreateRerankProviders(); - }, - getOwnPropertyDescriptor(target, key) { - if (key in target) { - return Reflect.getOwnPropertyDescriptor(target, key); - } - if (key in getOrCreateRerankProviders()) { - return { configurable: true, enumerable: true, value: getOrCreateRerankProviders()[key as string] }; - } - return undefined; - }, -}); - -export function getRerankProviders(): Record { - return RERANK_PROVIDERS; -} +}; const RERANK_PROVIDER_ALIASES = { jina: "jina-ai", diff --git a/open-sse/config/searchRegistry.ts b/open-sse/config/searchRegistry.ts index b33239a559..724d349525 100644 --- a/open-sse/config/searchRegistry.ts +++ b/open-sse/config/searchRegistry.ts @@ -26,11 +26,7 @@ export interface SearchProviderConfig { cacheTTLMs: number; } -let _SEARCH_PROVIDERS: Record | null = null; - -function getOrCreateSearchProviders(): Record { - if (!_SEARCH_PROVIDERS) { - _SEARCH_PROVIDERS = { +export const SEARCH_PROVIDERS: Record = { "serper-search": { id: "serper-search", name: "Serper Search", @@ -222,52 +218,14 @@ function getOrCreateSearchProviders(): Record { timeoutMs: 10_000, cacheTTLMs: 5 * 60 * 1000, }, - }; -} - return _SEARCH_PROVIDERS; -} +}; -export const SEARCH_PROVIDERS: Record = new Proxy({} as Record, { - get(target, key: string) { - if (key in target) { - return target[key]; - } - return getOrCreateSearchProviders()[key]; - }, - set(target, key: string, value) { - target[key] = value; - getOrCreateSearchProviders()[key] = value; - return true; - }, - deleteProperty(target, key: string) { - delete target[key]; - delete getOrCreateSearchProviders()[key]; - return true; - }, - ownKeys(target) { - const targetKeys = Reflect.ownKeys(target); - const registryKeys = Reflect.ownKeys(getOrCreateSearchProviders()); - return Array.from(new Set([...targetKeys, ...registryKeys])); - }, - has(target, key) { - return key in target || key in getOrCreateSearchProviders(); - }, - getOwnPropertyDescriptor(target, key) { - if (key in target) { - return Reflect.getOwnPropertyDescriptor(target, key); - } - if (key in getOrCreateSearchProviders()) { - return { configurable: true, enumerable: true, value: getOrCreateSearchProviders()[key as string] }; - } - return undefined; - }, -}); - -export function getSearchProviders(): Record { - return SEARCH_PROVIDERS; -} - -export const SEARCH_CREDENTIAL_FALLBACKS: Record = { "perplexity-search": "perplexity", +/** + * Credential fallback mapping — search providers that can reuse credentials + * from a related provider (e.g., perplexity-search uses the same API key as perplexity chat). + */ +export const SEARCH_CREDENTIAL_FALLBACKS: Record = { + "perplexity-search": "perplexity", "ollama-search": "ollama-cloud", "zai-search": "zai", }; diff --git a/open-sse/config/videoRegistry.ts b/open-sse/config/videoRegistry.ts index 4a0da6d253..b876bb0b6e 100644 --- a/open-sse/config/videoRegistry.ts +++ b/open-sse/config/videoRegistry.ts @@ -24,11 +24,7 @@ interface VideoProvider { models: VideoModel[]; } -let _VIDEO_PROVIDERS: Record | null = null; - -function getOrCreateVideoProviders(): Record { - if (!_VIDEO_PROVIDERS) { - _VIDEO_PROVIDERS = { +export const VIDEO_PROVIDERS: Record = { kie: { id: "kie", baseUrl: "https://api.kie.ai", @@ -152,59 +148,25 @@ function getOrCreateVideoProviders(): Record { format: "runwayml", models: RUNWAYML_SUPPORTED_VIDEO_MODELS, }, - }; -} - return _VIDEO_PROVIDERS; -} - -export const VIDEO_PROVIDERS: Record = new Proxy({} as Record, { - get(target, key: string) { - if (key in target) { - return target[key]; - } - return getOrCreateVideoProviders()[key]; - }, - set(target, key: string, value) { - target[key] = value; - getOrCreateVideoProviders()[key] = value; - return true; - }, - deleteProperty(target, key: string) { - delete target[key]; - delete getOrCreateVideoProviders()[key]; - return true; - }, - ownKeys(target) { - const targetKeys = Reflect.ownKeys(target); - const registryKeys = Reflect.ownKeys(getOrCreateVideoProviders()); - return Array.from(new Set([...targetKeys, ...registryKeys])); - }, - has(target, key) { - return key in target || key in getOrCreateVideoProviders(); - }, - getOwnPropertyDescriptor(target, key) { - if (key in target) { - return Reflect.getOwnPropertyDescriptor(target, key); - } - if (key in getOrCreateVideoProviders()) { - return { configurable: true, enumerable: true, value: getOrCreateVideoProviders()[key as string] }; - } - return undefined; - }, -}); - -export function getVideoProviders(): Record { - return VIDEO_PROVIDERS; -} +}; +/** + * Get video provider config by ID + */ export function getVideoProvider(providerId: string): VideoProvider | null { return VIDEO_PROVIDERS[providerId] || null; } +/** + * Parse video model string (format: "provider/model" or just "model") + */ export function parseVideoModel(modelStr: string | null) { return parseModelFromRegistry(modelStr, VIDEO_PROVIDERS); } +/** + * Get all video models as a flat list + */ export function getAllVideoModels() { return getAllModelsFromRegistry(VIDEO_PROVIDERS); -} \ No newline at end of file +} diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index 40b92b48f8..7fa8fdf045 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -69,7 +69,7 @@ import { connectionHasExtraKeys, type KeyHealth, } from "../services/apiKeyRotator.ts"; -import { isDetailedLoggingEnabled } from "@/lib/db/detailedLogs"; + import { getCallLogPipelineCaptureStreamChunks, getChatLogTextLimit, @@ -278,13 +278,7 @@ function cloneBoundedChatLogPayload(value: unknown, depth = 0): unknown { return result; } -function isSmallEnoughForSemanticCache(value: unknown): boolean { - try { - return JSON.stringify(value).length <= 256 * 1024; - } catch { - return false; - } -} +import { estimateSizeFast, isSmallEnoughForSemanticCache } from "../utils/estimateSize.ts"; function extractMemoryTextFromResponse( response: Record | null | undefined @@ -1519,9 +1513,19 @@ export async function handleChatCore({ comboName: comboName || undefined, }); }); + const traceEnabled = + process.env.OMNIRROUTE_TRACE === "true" || process.env.DEBUG === "true"; const trace = (label: string, extra?: Record) => { + if (!traceEnabled) return; const elapsed = Date.now() - startTime; - const suffix = extra ? ` ${JSON.stringify(extra)}` : ""; + let suffix = ""; + if (extra) { + try { + suffix = ` ${JSON.stringify(extra)}`; + } catch { + suffix = " [unserializable]"; + } + } log?.info?.("STAGE_TRACE", `${traceId} ${label} t=${elapsed}ms${suffix}`); }; let tokensCompressed: number | null = null; @@ -1913,7 +1917,13 @@ export async function handleChatCore({ ); } const noLogEnabled = apiKeyInfo?.noLog === true; - const detailedLoggingEnabled = !noLogEnabled && (await isDetailedLoggingEnabled()); + // Consolidate settings reads — fetch once, reuse throughout the request + const settings = cachedSettings ?? (await getCachedSettings()); + const detailedLoggingEnabled = + !noLogEnabled && + (settings.call_log_pipeline_enabled === true || + settings.call_log_pipeline_enabled === "1" || + settings.call_log_pipeline_enabled === "true"); const capturePipelineStreamChunks = detailedLoggingEnabled && getCallLogPipelineCaptureStreamChunks(); const skillRequestId = generateRequestId(); @@ -2436,7 +2446,7 @@ export async function handleChatCore({ let cavemanOutputModeIntensity: string | null = null; let preCompressionBody: typeof body | null = null; if (body && Array.isArray(allMessages) && allMessages.length > 0) { - let estimatedTokens = estimateTokens(JSON.stringify(allMessages)); + let estimatedTokens = estimateTokens(allMessages); let promptCompressionEnabled = false; let compressionSettings: CompressionConfig | null = null; @@ -2667,7 +2677,7 @@ export async function handleChatCore({ body = outputMode.body as typeof body; cavemanOutputModeApplied = true; cavemanOutputModeIntensity = config.cavemanOutputMode.intensity; - estimatedTokens = estimateTokens(JSON.stringify(body?.messages ?? body?.input ?? [])); + estimatedTokens = estimateTokens(body?.messages ?? body?.input ?? []); log?.debug?.("COMPRESSION", "Caveman output mode instruction applied"); } else if (outputMode.skippedReason && outputMode.skippedReason !== "disabled") { log?.debug?.("COMPRESSION", `Caveman output mode skipped: ${outputMode.skippedReason}`); @@ -2877,7 +2887,7 @@ export async function handleChatCore({ const COMPRESSION_THRESHOLD = 0.7; let reservedTokens = 0; if (Array.isArray(body.tools)) { - reservedTokens = estimateTokens(JSON.stringify(body.tools)); + reservedTokens = estimateTokens(body.tools); } const threshold = Math.max( 1, @@ -4119,7 +4129,7 @@ export async function handleChatCore({ (translatedBody.conversationState?.history?.length ?? 0) + (translatedBody.conversationState?.currentMessage ? 1 : 0) || 0; - log?.debug?.("REQUEST", `${provider.toUpperCase()} | ${model} | ${msgCount} msgs`); + log?.debug?.("REQUEST", `${provider?.toUpperCase()} | ${model} | ${msgCount} msgs`); // ── Tier 2: Authoritative per-model/provider token-limit check (provider now resolved) ── if (apiKeyInfo?.id) { @@ -4360,7 +4370,7 @@ export async function handleChatCore({ }; if (newCredentials?.accessToken || newCredentials?.copilotToken) { - log?.info?.("TOKEN", `${provider.toUpperCase()} | refreshed`); + log?.info?.("TOKEN", `${provider?.toUpperCase()} | refreshed`); // Fall back to post-mutex mutation only for executors that don't route // through getAccessToken (and therefore never fire onPersist). For @@ -4414,11 +4424,11 @@ export async function handleChatCore({ // than the original 401 alone. Surface at error level with sanitization. log?.error?.( "TOKEN", - `${provider.toUpperCase()} | retry after refresh failed: ${sanitizeErrorMessage(retryErr)}` + `${provider?.toUpperCase()} | retry after refresh failed: ${sanitizeErrorMessage(retryErr)}` ); } } else { - log?.warn?.("TOKEN", `${provider.toUpperCase()} | refresh failed`); + log?.warn?.("TOKEN", `${provider?.toUpperCase()} | refresh failed`); if (isUnrecoverableRefreshError(newCredentials) && onCredentialsRefreshed) { // Front 3 (reuse-race tolerance): before deactivating, re-read the DB. // If a sibling/concurrent refresh already rotated this connection's @@ -5090,8 +5100,10 @@ export async function handleChatCore({ // Save structured call log with full payloads const cacheUsageLogMeta = buildCacheUsageLogMeta(usage); if (usage && typeof usage === "object") { - const msg = `[${new Date().toLocaleTimeString("en-US", { hour12: false, hour: "2-digit", minute: "2-digit" })}] 📊 [USAGE] ${provider.toUpperCase()} | ${formatUsageLog(usage)}${connectionId ? ` | account=${connectionId.slice(0, 8)}...` : ""}`; - console.log(`${COLORS.green}${msg}${COLORS.reset}`); + if (traceEnabled) { + const msg = `[${new Date().toLocaleTimeString("en-US", { hour12: false, hour: "2-digit", minute: "2-digit" })}] 📊 [USAGE] ${provider?.toUpperCase()} | ${formatUsageLog(usage)}${connectionId ? ` | account=${connectionId.slice(0, 8)}...` : ""}`; + console.log(`${COLORS.green}${msg}${COLORS.reset}`); + } saveRequestUsage({ provider: provider || "unknown", diff --git a/open-sse/services/apiKeyRotator.ts b/open-sse/services/apiKeyRotator.ts index 1814fb833d..58a1337f0d 100644 --- a/open-sse/services/apiKeyRotator.ts +++ b/open-sse/services/apiKeyRotator.ts @@ -21,6 +21,9 @@ const _keyIndexes = new Map(); // Tracks which connections have extra API keys (for A3 guard in chatCore.ts) // Used to prevent disabling an entire connection when only one key fails. const _connectionExtraKeys = new Map(); +// Eviction limits to prevent unbounded memory growth under heavy load +const MAX_KEY_HEALTH_ENTRIES = 500; +const MAX_CONNECTION_EXTRA_KEYS = 500; /** * Record whether a connection has extra API keys. @@ -28,6 +31,10 @@ const _connectionExtraKeys = new Map(); */ export function trackConnectionExtraKeys(connectionId: string, extraKeys: string[]): void { const validExtras = extraKeys.filter((k) => typeof k === "string" && k.trim().length > 0); + if (!_connectionExtraKeys.has(connectionId) && _connectionExtraKeys.size >= MAX_CONNECTION_EXTRA_KEYS) { + const oldest = _connectionExtraKeys.keys().next().value; + if (oldest !== undefined) _connectionExtraKeys.delete(oldest); + } _connectionExtraKeys.set(connectionId, validExtras.length > 0); } @@ -64,6 +71,10 @@ const FAILURE_THRESHOLD = 2; // Mark as invalid after 2 consecutive failures function getOrCreateHealth(connectionId: string, keyId: string): KeyHealth { const scopedKey = `${connectionId}:${keyId}`; if (!_keyHealth.has(scopedKey)) { + if (_keyHealth.size >= MAX_KEY_HEALTH_ENTRIES) { + const oldest = _keyHealth.keys().next().value; + if (oldest !== undefined) _keyHealth.delete(oldest); + } _keyHealth.set(scopedKey, { status: "active", failures: 0, @@ -265,7 +276,12 @@ export function syncHealthFromDB(connectionId: string, health?: Record= MAX_KEY_HEALTH_ENTRIES) { + const oldest = _keyHealth.keys().next().value; + if (oldest !== undefined) _keyHealth.delete(oldest); + } + _keyHealth.set(scopedKey, keyHealth); } } @@ -336,6 +352,11 @@ export function removeConnectionHealth(connectionId: string): void { export function removeConnectionIndex(connectionId: string): void { _keyIndexes.delete(connectionId); _connectionExtraKeys.delete(connectionId); + for (const key of _keyHealth.keys()) { + if (key.startsWith(`${connectionId}:`)) { + _keyHealth.delete(key); + } + } } export type { KeyHealth }; diff --git a/open-sse/services/codexQuotaFetcher.ts b/open-sse/services/codexQuotaFetcher.ts index f800f3dedb..ef40c48286 100644 --- a/open-sse/services/codexQuotaFetcher.ts +++ b/open-sse/services/codexQuotaFetcher.ts @@ -76,6 +76,7 @@ interface CodexConnectionMeta { const MAX_CONNECTIONS = 100; const connectionRegistry = new Map(); +const MAX_QUOTA_CACHE_ENTRIES = 200; /** * Register Codex connection metadata for quota fetching. @@ -85,7 +86,7 @@ const connectionRegistry = new Map(); * @param meta - Access token and optional workspace ID */ export function registerCodexConnection(connectionId: string, meta: CodexConnectionMeta): void { - if (connectionRegistry.size >= MAX_CONNECTIONS) { + if (!connectionRegistry.has(connectionId) && connectionRegistry.size >= MAX_CONNECTIONS) { const oldestKey = connectionRegistry.keys().next().value; if (oldestKey !== undefined) { quotaCache.delete(oldestKey); @@ -123,6 +124,13 @@ function getCodexConnectionMeta( if (accessToken) { const meta = { accessToken, ...(workspaceId ? { workspaceId } : {}) }; + if (!connectionRegistry.has(connectionId) && connectionRegistry.size >= MAX_CONNECTIONS) { + const oldestKey = connectionRegistry.keys().next().value; + if (oldestKey !== undefined) { + quotaCache.delete(oldestKey); + connectionRegistry.delete(oldestKey); + } + } connectionRegistry.set(connectionId, meta); return meta; } @@ -204,6 +212,10 @@ export async function fetchCodexQuota( if (!quota) return null; // Store in cache + if (!quotaCache.has(connectionId) && quotaCache.size >= MAX_QUOTA_CACHE_ENTRIES) { + const oldestCacheKey = quotaCache.keys().next().value; + if (oldestCacheKey !== undefined) quotaCache.delete(oldestCacheKey); + } quotaCache.set(connectionId, { quota, fetchedAt: Date.now() }); return quota; } catch { diff --git a/open-sse/services/combo.ts b/open-sse/services/combo.ts index 3fcb7c8fee..aafd147238 100644 --- a/open-sse/services/combo.ts +++ b/open-sse/services/combo.ts @@ -495,6 +495,10 @@ export async function validateResponseQuality( // In-memory atomic counter per combo for round-robin distribution // Resets on server restart (by design — no stale state) +// Eviction limits to prevent unbounded memory growth +const MAX_RR_COUNTERS = 500; +const MAX_RESET_AWARE_CACHE = 200; + const rrCounters = new Map(); const resetAwareConnectionCache = new Map< @@ -1627,6 +1631,10 @@ async function getQuotaAwareConnectionsForTarget( const activeConnections = Array.isArray(connections) ? (connections as Array>) : []; + 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(), @@ -1758,6 +1766,10 @@ async function fetchResetAwareQuotaWithCache({ const refreshPromise = fetcher(connectionId, connection) .then((quota) => { if (quota) { + if (!resetAwareQuotaCache.has(cacheKey) && resetAwareQuotaCache.size >= MAX_RESET_AWARE_CACHE) { + const oldest = resetAwareQuotaCache.keys().next().value; + if (oldest !== undefined) resetAwareQuotaCache.delete(oldest); + } resetAwareQuotaCache.set(cacheKey, { quota, fetchedAt: Date.now(), @@ -1771,6 +1783,10 @@ async function fetchResetAwareQuotaWithCache({ .catch((error) => { const previous = resetAwareQuotaCache.get(cacheKey); if (previous) { + if (!resetAwareQuotaCache.has(cacheKey) && resetAwareQuotaCache.size >= MAX_RESET_AWARE_CACHE) { + const oldest = resetAwareQuotaCache.keys().next().value; + if (oldest !== undefined) resetAwareQuotaCache.delete(oldest); + } resetAwareQuotaCache.set(cacheKey, { ...previous, refreshPromise: null }); } log.warn?.("COMBO", "Reset-aware quota fetch failed.", { @@ -1783,6 +1799,10 @@ async function fetchResetAwareQuotaWithCache({ return null; }); + if (!resetAwareQuotaCache.has(cacheKey) && resetAwareQuotaCache.size >= MAX_RESET_AWARE_CACHE) { + const oldest = resetAwareQuotaCache.keys().next().value; + if (oldest !== undefined) resetAwareQuotaCache.delete(oldest); + } resetAwareQuotaCache.set(cacheKey, { quota: existing?.quota ?? cached?.quota ?? null, fetchedAt: existing?.fetchedAt ?? cached?.fetchedAt ?? 0, @@ -1906,6 +1926,10 @@ async function orderTargetsByResetAwareQuota( if (tiedTargets.length > 1) { const key = `reset-aware:${comboName}`; const counter = rrCounters.get(key) || 0; + if (!rrCounters.has(key) && rrCounters.size >= MAX_RR_COUNTERS) { + const oldest = rrCounters.keys().next().value; + if (oldest !== undefined) rrCounters.delete(oldest); + } rrCounters.set(key, counter + 1); const startIndex = counter % tiedTargets.length; orderedTiedTargets = [...tiedTargets.slice(startIndex), ...tiedTargets.slice(0, startIndex)]; @@ -2065,6 +2089,10 @@ async function orderTargetsByResetWindow( const key = `reset-window:${comboName}`; const counter = rrCounters.get(key) || 0; + if (!rrCounters.has(key) && rrCounters.size >= MAX_RR_COUNTERS) { + const oldest = rrCounters.keys().next().value; + if (oldest !== undefined) rrCounters.delete(oldest); + } rrCounters.set(key, counter + 1); const startIndex = counter % tiedTargets.length; const orderedTiedTargets = [ @@ -3289,7 +3317,7 @@ export async function handleComboChat({ if (isModelAvailable) { const available = await isModelAvailable(modelStr, targetForAttempt); if (!available) { - log.info("COMBO", `Skipping ${modelStr} — no credentials available or model excluded`); + log.debug?.("COMBO", `Skipping ${modelStr} — no credentials available or model excluded`); if (i > 0) fallbackCount++; return null; } @@ -3802,7 +3830,7 @@ export async function handleComboChat({ ? Math.min(cooldownMs, fallbackDelayMs) : 0; if ([502, 503, 504].includes(result.status) && fallbackWaitMs > 0) { - log.info("COMBO", `Waiting ${fallbackWaitMs}ms before fallback to next model`); + log.debug?.("COMBO", `Waiting ${fallbackWaitMs}ms before fallback to next model`); await new Promise((resolve) => { const timer = setTimeout(resolve, fallbackWaitMs); signal?.addEventListener( @@ -3990,6 +4018,10 @@ async function handleRoundRobinCombo({ // Get and increment atomic counter const counter = rrCounters.get(combo.name) || 0; + if (!rrCounters.has(combo.name) && rrCounters.size >= MAX_RR_COUNTERS) { + const oldest = rrCounters.keys().next().value; + if (oldest !== undefined) rrCounters.delete(oldest); + } rrCounters.set(combo.name, counter + 1); const startIndex = counter % modelCount; @@ -4026,7 +4058,7 @@ async function handleRoundRobinCombo({ if (isModelAvailable) { const available = await isModelAvailable(modelStr, targetForAttempt); if (!available) { - log.info("COMBO-RR", `Skipping ${modelStr} — no credentials available or model excluded`); + log.debug?.("COMBO-RR", `Skipping ${modelStr} — no credentials available or model excluded`); if (offset > 0) fallbackCount++; continue; } @@ -4267,7 +4299,7 @@ async function handleRoundRobinCombo({ isAllAccountsRateLimited); if (providerExhausted) { exhaustedProviders.add(provider); - log.info("COMBO-RR", `Provider ${provider} quota exhausted — marking for skip (#1731)`); + log.debug?.("COMBO-RR", `Provider ${provider} quota exhausted — marking for skip (#1731)`); } else if ( result.status === 429 && !isTokenLimitBreach && @@ -4324,7 +4356,7 @@ async function handleRoundRobinCombo({ ? Math.min(cooldownMs, fallbackDelayMs) : 0; if ([502, 503, 504].includes(result.status) && fallbackWaitMs > 0) { - log.info("COMBO-RR", `Waiting ${fallbackWaitMs}ms before fallback to next model`); + log.debug?.("COMBO-RR", `Waiting ${fallbackWaitMs}ms before fallback to next model`); await new Promise((resolve) => { const timer = setTimeout(resolve, fallbackWaitMs); signal?.addEventListener( diff --git a/open-sse/services/rateLimitManager.ts b/open-sse/services/rateLimitManager.ts index a46a016574..c5a4a181ae 100644 --- a/open-sse/services/rateLimitManager.ts +++ b/open-sse/services/rateLimitManager.ts @@ -192,6 +192,20 @@ function reconcileEnabledConnections( function watchdogTick() { const now = Date.now(); + // Clean up idle limiters that haven't been used recently + for (const [key, limiter] of Array.from(limiters)) { + const lastUsed = limiterLastUsed.get(key) ?? 0; + if (now - lastUsed > INACTIVE_LIMITER_MS) { + const counts = limiter.counts(); + if (counts.QUEUED === 0 && counts.RUNNING === 0 && counts.EXECUTING === 0) { + limiters.delete(key); + lastDispatchAt.delete(key); + limiterLastUsed.delete(key); + logRateLimit(`🧹 [RATE-LIMIT] Evicting idle limiter: ${key} (inactive for ${Math.round((now - lastUsed) / 1000)}s)`); + trackAsyncOperation(limiter.disconnect()); + } + } + } for (const [key, limiter] of Array.from(limiters)) { const counts = limiter.counts(); if (counts.QUEUED === 0) continue; @@ -211,6 +225,7 @@ function watchdogTick() { ); limiters.delete(key); lastDispatchAt.delete(key); + limiterLastUsed.delete(key); // Do NOT call limiter.stop() — it permanently rejects future .schedule() calls with // "This limiter has been stopped". In-flight requests still holding a reference to // the old instance cannot be redirected to a new one, causing spurious 502 bursts. @@ -257,6 +272,7 @@ function shutdownLimiters(): void { } limiters.clear(); lastDispatchAt.clear(); + limiterLastUsed.clear(); } // Only register shutdown handlers when there are active limiters to shut down. @@ -349,6 +365,7 @@ export function disableRateLimitProtection(connectionId) { if (key.includes(connectionId)) { limiters.delete(key); lastDispatchAt.delete(key); + limiterLastUsed.delete(key); trackAsyncOperation(limiter.disconnect()); } } @@ -403,8 +420,10 @@ function getLimiter(provider, connectionId, model = null) { limiters.set(key, limiter); lastDispatchAt.set(key, Date.now()); + limiterLastUsed.set(key, Date.now()); } + limiterLastUsed.set(key, Date.now()); return limiters.get(key); } @@ -626,6 +645,7 @@ export function updateFromHeaders(provider, connectionId, headers, status, model // the abandoned Bottleneck; under sustained quota pressure that is a real leak. limiters.delete(limiterKey); lastDispatchAt.delete(limiterKey); + limiterLastUsed.delete(limiterKey); trackAsyncOperation(limiter.disconnect()); return; } @@ -799,6 +819,7 @@ export async function __resetRateLimitManagerForTests() { enabledConnections.clear(); initialized = false; lastDispatchAt.clear(); + limiterLastUsed.clear(); shutdownHandlersRegistered = false; for (const key of Object.keys(learnedLimits)) { diff --git a/open-sse/utils/estimateSize.ts b/open-sse/utils/estimateSize.ts new file mode 100644 index 0000000000..ac320f9aad --- /dev/null +++ b/open-sse/utils/estimateSize.ts @@ -0,0 +1,35 @@ +/** + * Fast object-tree size estimator — walks without JSON.stringify. + * Safe for circular references (uses WeakSet). + * Early-exits at 256KB to avoid wasting CPU on huge payloads. + */ +export function estimateSizeFast(value: unknown): number { + let bytes = 0; + const stack: unknown[] = [value]; + const seen = new WeakSet(); + while (stack.length > 0) { + const v = stack.pop(); + if (v === null || v === undefined) continue; + if (typeof v === "string") { + bytes += v.length; + if (bytes > 262144) return bytes; + } else if (typeof v === "number") bytes += 8; + else if (typeof v === "boolean") bytes += 4; + else if (typeof v === "object") { + if (seen.has(v as object)) continue; + seen.add(v as object); + if (Array.isArray(v)) { + for (let i = 0; i < v.length; i++) stack.push(v[i]); + } else { + for (const key in v) { + if (Object.prototype.hasOwnProperty.call(v, key)) stack.push((v as Record)[key]); + } + } + } + } + return bytes; +} + +export function isSmallEnoughForSemanticCache(value: unknown): boolean { + return estimateSizeFast(value) <= 256 * 1024; +} diff --git a/tests/unit/estimateSizeFast.test.ts b/tests/unit/estimateSizeFast.test.ts new file mode 100644 index 0000000000..d84fc75086 --- /dev/null +++ b/tests/unit/estimateSizeFast.test.ts @@ -0,0 +1,111 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +const { estimateSizeFast, isSmallEnoughForSemanticCache } = await import( + "../../open-sse/utils/estimateSize.ts" +); + +test("estimateSizeFast returns 0 for null/undefined", () => { + assert.equal(estimateSizeFast(null), 0); + assert.equal(estimateSizeFast(undefined), 0); +}); + +test("estimateSizeFast counts string lengths", () => { + assert.equal(estimateSizeFast("hello"), 5); + assert.equal(estimateSizeFast(""), 0); +}); + +test("estimateSizeFast counts numbers as 8 bytes", () => { + assert.equal(estimateSizeFast(42), 8); + assert.equal(estimateSizeFast(0), 8); + assert.equal(estimateSizeFast(3.14), 8); +}); + +test("estimateSizeFast counts booleans as 4 bytes", () => { + assert.equal(estimateSizeFast(true), 4); + assert.equal(estimateSizeFast(false), 4); +}); + +test("estimateSizeFast walks arrays recursively", () => { + const arr = ["abc", "de", 42]; + assert.equal(estimateSizeFast(arr), 3 + 2 + 8); // 13 +}); + +test("estimateSizeFast walks objects recursively", () => { + const obj = { a: "hello", b: 42 }; + assert.equal(estimateSizeFast(obj), 5 + 8); // 13 +}); + +test("estimateSizeFast walks nested structures", () => { + const nested = { messages: [{ role: "user", content: "hi" }] }; + // role=4, content=2 + assert.equal(estimateSizeFast(nested), 4 + 2); // 6 +}); + +test("estimateSizeFast handles circular references without infinite loop", () => { + const circular: Record = { a: "test" }; + circular.self = circular; // Create circular ref + // Should not hang — WeakSet skips already-visited objects + const result = estimateSizeFast(circular); + assert.equal(result, 4); // Only "test" (4) counted; circular ref skipped +}); + +test("estimateSizeFast handles deeply nested circular refs", () => { + const a: Record = { val: "x" }; + const b: Record = { ref: a }; + a.back = b; + const result = estimateSizeFast({ root: a }); + assert.equal(result, 1); // "x" = 1 +}); + +test("estimateSizeFast early-exits at 262144 bytes (256KB)", () => { + // Create a string > 256KB + const bigStr = "x".repeat(300_000); + const result = estimateSizeFast(bigStr); + assert.ok(result >= 262144, `Should early-exit, got ${result}`); +}); + +test("estimateSizeFast handles mixed object/array nesting", () => { + const data = { + choices: [ + { + delta: { content: "Hello world" }, + index: 0, + }, + ], + }; + // content=11, index=8 (number), delta keys: content+delta=7, choices=8 + const result = estimateSizeFast(data); + assert.ok(result > 0); + assert.ok(result < 100); +}); + +test("estimateSizeFast does not count keys, only values", () => { + // Object with long keys but short values + const obj = { aLongKeyName: "x", anotherLongKeyName: "y" }; + assert.equal(estimateSizeFast(obj), 2); // "x" + "y" +}); + +test("isSmallEnoughForSemanticCache returns true for small payloads", () => { + assert.ok(isSmallEnoughForSemanticCache({ msg: "hi" })); +}); + +test("isSmallEnoughForSemanticCache returns false for huge payloads", () => { + const huge = { data: "x".repeat(300_000) }; + assert.ok(!isSmallEnoughForSemanticCache(huge)); +}); + +test("isSmallEnoughForSemanticCache handles circular refs gracefully", () => { + const circular: Record = {}; + circular.self = circular; + // Should not hang; estimateSizeFast has WeakSet protection + const result = isSmallEnoughForSemanticCache(circular); + assert.equal(result, true); // 0 bytes < 256KB +}); + +test("estimateSizeFast handles Map-like objects (no infinite loop on iterables)", () => { + const map = new Map([["key", "value"]]); + // Maps are objects but have no enumerable own properties via for-in + const result = estimateSizeFast(map); + assert.ok(typeof result === "number"); +}); diff --git a/tests/unit/eviction-guards-apiKeyRotator.test.ts b/tests/unit/eviction-guards-apiKeyRotator.test.ts new file mode 100644 index 0000000000..59e086db80 --- /dev/null +++ b/tests/unit/eviction-guards-apiKeyRotator.test.ts @@ -0,0 +1,81 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +const rotator = await import("../../open-sse/services/apiKeyRotator.ts"); +const { + trackConnectionExtraKeys, + connectionHasExtraKeys, + getAllKeyHealth, + syncHealthFromDB, + resetKeyStatus, + removeConnectionIndex, +} = rotator; + +test("trackConnectionExtraKeys: inserting into a full map does not evict existing key being updated", () => { + // Fill the map to capacity by inserting many unique connection IDs + for (let i = 0; i < 510; i++) { + trackConnectionExtraKeys(`conn-${i}`, [`key-${i}`]); + } + + // Now update an existing key — this should NOT evict it + trackConnectionExtraKeys("conn-0", ["key-0", "key-new"]); + assert.ok(connectionHasExtraKeys("conn-0", ["key-0"]), "Existing key should not be evicted on update"); +}); + +test("trackConnectionExtraKeys: evicts oldest when inserting NEW key at capacity", () => { + // The map was filled above. Insert a brand new key — oldest should be evicted + const before = connectionHasExtraKeys("conn-1", ["key-1"]); + trackConnectionExtraKeys("conn-NEW-UNIQUE-XYZ", ["new-key"]); + // conn-1 may or may not be evicted depending on insertion order, but the map should not grow unbounded + assert.ok(typeof before === "boolean"); +}); + +test("syncHealthFromDB: does not evict when updating existing scopedKey", () => { + // Seed with some health entries + for (let i = 0; i < 5; i++) { + resetKeyStatus("test-conn", `key-${i}`); + } + + // Sync health for existing entries — should not evict them + const health = { + "key-0": { status: "active" as const, failures: 0, lastFailure: 0 }, + "key-1": { status: "active" as const, failures: 0, lastFailure: 0 }, + }; + syncHealthFromDB("test-conn", health); + + const all = getAllKeyHealth(); + assert.ok(all["test-conn:key-0"], "key-0 should still exist after sync"); + assert.ok(all["test-conn:key-1"], "key-1 should still exist after sync"); +}); + +test("syncHealthFromDB: evicts oldest when inserting NEW scopedKey at capacity", () => { + // Fill the map by syncing many unique entries + for (let i = 0; i < 505; i++) { + syncHealthFromDB(`bulk-conn-${i}`, { + [`bulk-key-${i}`]: { status: "active" as const, failures: 0, lastFailure: 0 }, + }); + } + // Insert one more brand new entry — should trigger eviction of oldest + syncHealthFromDB("bulk-conn-NEW", { + "bulk-key-NEW": { status: "active" as const, failures: 0, lastFailure: 0 }, + }); + const all = getAllKeyHealth(); + assert.ok(all["bulk-conn-NEW:bulk-key-NEW"], "New entry should exist"); +}); + +test("removeConnectionIndex cleans all 3 maps (keyIndexes, connectionExtraKeys, keyHealth)", () => { + // Seed data + trackConnectionExtraKeys("cleanup-conn", ["k1", "k2"]); + resetKeyStatus("cleanup-conn", "k1"); + + // Verify data exists via the in-memory cache (no extraKeys arg) + assert.ok(connectionHasExtraKeys("cleanup-conn")); + + // Remove via removeConnectionIndex (cleans all 3 maps) + removeConnectionIndex("cleanup-conn"); + + // Verify cleaned — in-memory cache should be empty + assert.ok(!connectionHasExtraKeys("cleanup-conn")); + const all = getAllKeyHealth(); + assert.ok(!all["cleanup-conn:k1"], "Health entry should be removed"); +}); diff --git a/tests/unit/eviction-guards-codexQuotaFetcher.test.ts b/tests/unit/eviction-guards-codexQuotaFetcher.test.ts new file mode 100644 index 0000000000..9f06aab365 --- /dev/null +++ b/tests/unit/eviction-guards-codexQuotaFetcher.test.ts @@ -0,0 +1,49 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +const codex = await import("../../open-sse/services/codexQuotaFetcher.ts"); +const { registerCodexConnection, unregisterCodexConnection, getCodexConnectionMeta } = codex; + +// getCodexConnectionMeta is not exported — use registerCodexConnection + internal verify +// Let's check what is exported +const exportedKeys = Object.keys(codex).filter((k) => typeof codex[k] === "function"); +assert.ok(exportedKeys.length > 0, "codexQuotaFetcher should export functions"); + +test("registerCodexConnection: inserting into a full map does not evict the key being updated", () => { + // Fill registry to capacity + for (let i = 0; i < 210; i++) { + registerCodexConnection(`conn-${i}`, { accessToken: `tok-${i}` }); + } + + // Update an existing connection — should NOT trigger eviction + registerCodexConnection("conn-0", { accessToken: "tok-0-updated" }); + + // If conn-0 was evicted, unregistering it would be a no-op. + // The key point: this should not throw or corrupt state. + unregisterCodexConnection("conn-0"); +}); + +test("registerCodexConnection: evicts oldest when inserting NEW connection at capacity", () => { + // Re-fill to capacity + for (let i = 0; i < 210; i++) { + registerCodexConnection(`fill-${i}`, { accessToken: `tok-${i}` }); + } + + // Insert a brand new one — should evict the oldest + registerCodexConnection("fill-NEW-UNIQUE", { accessToken: "tok-new" }); + + // The new entry should be registered (no throw) + unregisterCodexConnection("fill-NEW-UNIQUE"); +}); + +test("registerCodexConnection does not throw on normal usage", () => { + registerCodexConnection("test-conn", { accessToken: "test-token" }); + unregisterCodexConnection("test-conn"); +}); + +test("unregisterCodexConnection is idempotent", () => { + registerCodexConnection("idempotent-test", { accessToken: "tok" }); + unregisterCodexConnection("idempotent-test"); + // Should not throw on double unregister + unregisterCodexConnection("idempotent-test"); +}); diff --git a/tests/unit/rateLimitManager-idle-eviction.test.ts b/tests/unit/rateLimitManager-idle-eviction.test.ts new file mode 100644 index 0000000000..e316002c3a --- /dev/null +++ b/tests/unit/rateLimitManager-idle-eviction.test.ts @@ -0,0 +1,86 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +const rlm = await import("../../open-sse/services/rateLimitManager.ts"); +const { + enableRateLimitProtection, + disableRateLimitProtection, + isRateLimitEnabled, + withRateLimit, + updateFromHeaders, + getAllRateLimitStatus, + startRateLimitWatchdog, + stopRateLimitWatchdog, + __resetRateLimitManagerForTests, + __getLimiterStateForTests, +} = rlm; + +// Clean slate before each test +test.beforeEach(async () => { + await __resetRateLimitManagerForTests(); +}); + +test("enableRateLimitProtection creates limiter", async () => { + enableRateLimitProtection("test-conn-1"); + assert.ok(isRateLimitEnabled("test-conn-1")); +}); + +test("disableRateLimitProtection cleans up limiters and limiterLastUsed", async () => { + // Create a limiter by using withRateLimit + enableRateLimitProtection("test-conn-2"); + await withRateLimit("openai", "test-conn-2", "gpt-4", async () => "ok"); + + // Verify limiter exists + const before = getAllRateLimitStatus(); + const hasKey = Object.keys(before).some((k) => k.includes("test-conn-2")); + + // Disable — should clean up all 3 Maps (limiters, lastDispatchAt, limiterLastUsed) + disableRateLimitProtection("test-conn-2"); + assert.ok(!isRateLimitEnabled("test-conn-2")); +}); + +test("limiterLastUsed is populated on each withRateLimit call", async () => { + enableRateLimitProtection("test-conn-3"); + const result = await withRateLimit("anthropic", "test-conn-3", "claude-3", async () => "response"); + assert.equal(result, "response"); + + // Second call should also work (limiterLastUsed prevents eviction) + const result2 = await withRateLimit("anthropic", "test-conn-3", "claude-3", async () => "response2"); + assert.equal(result2, "response2"); +}); + +test("updateFromHeaders with 429 triggers limiter disconnect and cleanup", async () => { + enableRateLimitProtection("test-conn-4"); + await withRateLimit("openai", "test-conn-4", "gpt-4", async () => "ok"); + + // Simulate a 429 response — this should disconnect the old limiter + const headers = new Headers({ "retry-after": "60" }); + updateFromHeaders("openai", "test-conn-4", headers, 429, "gpt-4"); + + // Should not throw — old limiter was properly cleaned up + assert.ok(true); +}); + +test("shutdown clears all maps including limiterLastUsed", async () => { + enableRateLimitProtection("test-conn-5"); + await withRateLimit("openai", "test-conn-5", "gpt-4", async () => "ok"); + + // __resetRateLimitManagerForTests calls shutdown internally + await __resetRateLimitManagerForTests(); + + // All limiters should be gone + const after = getAllRateLimitStatus(); + assert.equal(Object.keys(after).length, 0); +}); + +test("multiple providers/connections create separate limiters", async () => { + enableRateLimitProtection("conn-a"); + enableRateLimitProtection("conn-b"); + await withRateLimit("openai", "conn-a", "gpt-4", async () => "a"); + await withRateLimit("anthropic", "conn-b", "claude-3", async () => "b"); + + const status = getAllRateLimitStatus(); + const keys = Object.keys(status); + // Should have at least 2 separate limiters + assert.ok(keys.length >= 2, `Expected >=2 limiters, got ${keys.length}`); +}); diff --git a/tests/unit/registry-direct-exports.test.ts b/tests/unit/registry-direct-exports.test.ts new file mode 100644 index 0000000000..3c784966b1 --- /dev/null +++ b/tests/unit/registry-direct-exports.test.ts @@ -0,0 +1,63 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +// Verify all 8 registries export plain objects (no Proxy, no lazy getter) +const registries = [ + { name: "audio", mod: await import("../../open-sse/config/audioRegistry.ts"), keys: ["AUDIO_TRANSCRIPTION_PROVIDERS", "AUDIO_SPEECH_PROVIDERS"] }, + { name: "embedding", mod: await import("../../open-sse/config/embeddingRegistry.ts"), keys: ["EMBEDDING_PROVIDERS"] }, + { name: "image", mod: await import("../../open-sse/config/imageRegistry.ts"), keys: ["IMAGE_PROVIDERS"] }, + { name: "moderation", mod: await import("../../open-sse/config/moderationRegistry.ts"), keys: ["MODERATION_PROVIDERS"] }, + { name: "music", mod: await import("../../open-sse/config/musicRegistry.ts"), keys: ["MUSIC_PROVIDERS"] }, + { name: "rerank", mod: await import("../../open-sse/config/rerankRegistry.ts"), keys: ["RERANK_PROVIDERS"] }, + { name: "search", mod: await import("../../open-sse/config/searchRegistry.ts"), keys: ["SEARCH_PROVIDERS"] }, + { name: "video", mod: await import("../../open-sse/config/videoRegistry.ts"), keys: ["VIDEO_PROVIDERS"] }, +]; + +for (const { name, mod, keys } of registries) { + for (const key of keys) { + test(`${name} registry: ${key} is a plain object, not a Proxy`, () => { + const registry = mod[key]; + assert.ok(registry, `${key} should be exported`); + assert.equal(typeof registry, "object"); + + // Proxy traps break Object.keys() — plain objects return keys immediately + const firstKey = Object.keys(registry)[0]; + assert.ok(firstKey, `${key} should have at least one entry`); + + // Direct property access should work without trap overhead + const entry = registry[firstKey]; + assert.ok(entry, `First entry should be accessible`); + }); + + test(`${name} registry: ${key} entries are mutable (no Proxy freeze)`, () => { + const registry = mod[key]; + const firstKey = Object.keys(registry)[0]; + const original = registry[firstKey]; + + // Should be able to mutate without Proxy restrictions + registry[firstKey] = { ...original, _test: true }; + assert.ok(registry[firstKey]._test === true); + + // Restore + registry[firstKey] = original; + }); + } +} + +// Verify registries don't use lazy initialization patterns +test("registries do not contain getOrCreate* functions", async () => { + for (const { name, mod } of registries) { + const fns = Object.keys(mod).filter((k) => typeof mod[k] === "function"); + const lazyFns = fns.filter((fn) => fn.startsWith("getOrCreate")); + assert.equal(lazyFns.length, 0, `${name} has lazy getter: ${lazyFns.join(", ")}`); + } +}); + +test("audioRegistry exports per-type provider lookup functions", async () => { + const { getTranscriptionProvider, getSpeechProvider } = await import( + "../../open-sse/config/audioRegistry.ts" + ); + // Should return null for unknown providers (not throw) + assert.equal(getTranscriptionProvider("nonexistent-provider"), null); + assert.equal(getSpeechProvider("nonexistent-provider"), null); +});