From c2e6b618e6376ce463a2bebe0cf13b04e3edae75 Mon Sep 17 00:00:00 2001 From: diegosouzapw Date: Thu, 18 Jun 2026 10:19:42 -0300 Subject: [PATCH] =?UTF-8?q?refactor(sse):=20split=20chatCore.ts=20pure=20h?= =?UTF-8?q?elpers=20into=20chatCore/=20modules=20(=E2=88=92561=20LOC)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Extracts the pure, side-effect-free helpers out of the chatCore.ts god-file into focused open-sse/handlers/chatCore/ modules (headers, logTruncation, memoryExtraction, nonStreamingSse, passthroughToolNames, upstreamTimeouts), each with its own unit test. Behavior-preserving — chatCore.ts re-imports them. Re-synced cleanly onto release/v3.8.29: the prior branch tip carried a stale-merge regression that re-bloated open-sse/config/freeModelCatalog.data.ts (462 -> 4530 lines) and reverted the free-tier refresh + provider/docs changes. This commit keeps ONLY the chatCore split (zero contamination). --- config/quality/file-size-baseline.json | 2 +- open-sse/handlers/chatCore.ts | 622 +----------------- open-sse/handlers/chatCore/headers.ts | 16 + open-sse/handlers/chatCore/logTruncation.ts | 84 +++ .../handlers/chatCore/memoryExtraction.ts | 131 ++++ open-sse/handlers/chatCore/nonStreamingSse.ts | 141 ++++ .../handlers/chatCore/passthroughToolNames.ts | 65 ++ .../handlers/chatCore/upstreamTimeouts.ts | 158 +++++ open-sse/services/tokenLimitCounter.ts | 2 +- tests/unit/chatcore-headers.test.ts | 20 + tests/unit/chatcore-log-truncation.test.ts | 46 ++ tests/unit/chatcore-memory-extraction.test.ts | 42 ++ tests/unit/chatcore-non-streaming-sse.test.ts | 72 ++ .../chatcore-passthrough-tool-names.test.ts | 41 ++ tests/unit/chatcore-upstream-timeouts.test.ts | 51 ++ 15 files changed, 898 insertions(+), 595 deletions(-) create mode 100644 open-sse/handlers/chatCore/headers.ts create mode 100644 open-sse/handlers/chatCore/logTruncation.ts create mode 100644 open-sse/handlers/chatCore/memoryExtraction.ts create mode 100644 open-sse/handlers/chatCore/nonStreamingSse.ts create mode 100644 open-sse/handlers/chatCore/passthroughToolNames.ts create mode 100644 open-sse/handlers/chatCore/upstreamTimeouts.ts create mode 100644 tests/unit/chatcore-headers.test.ts create mode 100644 tests/unit/chatcore-log-truncation.test.ts create mode 100644 tests/unit/chatcore-memory-extraction.test.ts create mode 100644 tests/unit/chatcore-non-streaming-sse.test.ts create mode 100644 tests/unit/chatcore-passthrough-tool-names.test.ts create mode 100644 tests/unit/chatcore-upstream-timeouts.test.ts diff --git a/config/quality/file-size-baseline.json b/config/quality/file-size-baseline.json index 4215fefff4..357cd0e0f4 100644 --- a/config/quality/file-size-baseline.json +++ b/config/quality/file-size-baseline.json @@ -52,7 +52,7 @@ "open-sse/executors/muse-spark-web.ts": 1284, "open-sse/executors/perplexity-web.ts": 1013, "open-sse/handlers/audioSpeech.ts": 965, - "open-sse/handlers/chatCore.ts": 6009, + "open-sse/handlers/chatCore.ts": 5445, "open-sse/handlers/imageGeneration.ts": 3777, "open-sse/handlers/responseSanitizer.ts": 1103, "open-sse/handlers/search.ts": 1546, diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index 416c43641a..45f3b13625 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -2,6 +2,13 @@ import { injectMemoryAndSkills } from "./chatCore/memorySkillsInjection.ts"; import { checkIdempotencyCache } from "./chatCore/idempotency.ts"; import { checkSemanticCache } from "./chatCore/semanticCache.ts"; import { sanitizeChatRequestBody } from "./chatCore/sanitization.ts"; +import { cloneBoundedChatLogPayload, truncateForLog } from "./chatCore/logTruncation.ts"; +import { getHeaderValueCaseInsensitive } from "./chatCore/headers.ts"; +import { + extractMemoryTextFromResponse, + extractMemoryTextFromRequestBody, + resolveMemoryOwnerId, +} from "./chatCore/memoryExtraction.ts"; import { CORS_HEADERS } from "../utils/cors.ts"; import { HEAP_PRESSURE_THRESHOLD_MB } from "../utils/heapPressure.ts"; import { normalizeHeaders } from "../utils/headers.ts"; @@ -59,7 +66,6 @@ import { import { COOLDOWN_MS, HTTP_STATUS, - FETCH_TIMEOUT_MS, FETCH_BODY_TIMEOUT_MS, MAX_TOOLS_LIMIT, PROVIDER_MAX_TOKENS, @@ -89,10 +95,6 @@ import { import { getCallLogPipelineCaptureStreamChunks, getCallLogPipelineMaxSizeBytes, - getChatLogTextLimit, - getChatLogArrayTailItems, - getChatLogMaxDepth, - getChatLogMaxObjectKeys, } from "@/lib/logEnv"; import { logAuditEvent } from "@/lib/compliance"; import { emit } from "@/lib/events/eventBus"; @@ -106,16 +108,29 @@ import { saveCallLog, } from "@/lib/usageDb"; import { finalizePendingScope, updatePendingScope } from "@/lib/usage/pendingRequestScope"; -import { - formatUsageLog, - getLoggedInputTokens, - getLoggedOutputTokens, - getReasoningTokens, -} from "@/lib/usage/tokenAccounting"; +import { formatUsageLog } from "@/lib/usage/tokenAccounting"; import { recordCost } from "@/domain/costRules"; import { calculateCost } from "@/lib/usage/costCalculator"; import { buildOmniRouteResponseMetaHeaders } from "@/domain/omnirouteResponseMeta"; -import { CLAUDE_OAUTH_TOOL_PREFIX } from "../translator/request/openai-to-claude.ts"; +import { + buildClaudePassthroughToolNameMap, + restoreClaudePassthroughToolNames, + mergeResponseToolNameMap, +} from "./chatCore/passthroughToolNames.ts"; +import { + parseNonStreamingSSEPayload, + normalizeNonStreamingEventPayload, + shouldTreatBufferedEventResponseAsExpected, + appendNonStreamingSseTerminalSignal, + type NonStreamingSseTerminalState, +} from "./chatCore/nonStreamingSse.ts"; +import { + createBodyTimeoutError, + readStreamChunkWithTimeout, + computeBillableTokens, + normalizeExecutorResult, + executeWithUpstreamStartTimeout, +} from "./chatCore/upstreamTimeouts.ts"; import { getModelNormalizeToolCallId, getModelPreserveOpenAIDeveloperRole, @@ -164,12 +179,7 @@ import { import { invalidateCodexQuotaCache } from "../services/codexQuotaFetcher.ts"; import { translateNonStreamingResponse } from "./responseTranslator.ts"; import { extractUsageFromResponse } from "./usageExtractor.ts"; -import { - extractSSEErrorMessage, - parseSSEToClaudeResponse, - parseSSEToOpenAIResponse, - parseSSEToResponsesOutput, -} from "./sseParser.ts"; +import { extractSSEErrorMessage } from "./sseParser.ts"; import { sanitizeOpenAIResponse, sanitizeResponsesApiResponse } from "./responseSanitizer.ts"; import { withRateLimit, @@ -240,8 +250,6 @@ import { } from "../services/modelscopePolicy.ts"; import { incrementRequestCount } from "../services/geminiRateLimitTracker.ts"; -const MEMORY_EXTRACTION_TEXT_LIMIT = 64 * 1024; - // ── Global memory pressure guard ──────────────────────────────────────── // Prevents OOM by rejecting new requests when V8 heap exceeds threshold. // Self-healing: no counters to leak, no cleanup needed. The threshold @@ -249,204 +257,7 @@ const MEMORY_EXTRACTION_TEXT_LIMIT = 64 * 1024; // it tracks --max-old-space-size across 1GB/2GB/large VPS instead of a fixed // 200MB that sat below the app's own ~260MB baseline and rejected every request. -function capMemoryExtractionText(value: string): string { - if (value.length <= MEMORY_EXTRACTION_TEXT_LIMIT) return value; - return value.slice(-MEMORY_EXTRACTION_TEXT_LIMIT); -} - -function truncateChatLogText(value: string): string { - const limit = getChatLogTextLimit(); - if (value.length <= limit) return value; - const head = value.slice(0, Math.floor(limit / 2)); - const tail = value.slice(-Math.ceil(limit / 2)); - return `${head}\n[...truncated ${value.length - limit} chars...]\n${tail}`; -} - -function cloneBoundedChatLogPayload(value: unknown, depth = 0): unknown { - if (value === null || value === undefined) return value; - if (typeof value === "string") return truncateChatLogText(value); - if (typeof value !== "object") return value; - if (depth >= getChatLogMaxDepth()) return "[MaxDepth]"; - - const maxTailItems = getChatLogArrayTailItems(); - - if (Array.isArray(value)) { - const retained = value.length > maxTailItems ? value.slice(-maxTailItems) : value; - const cloned = retained.map((item) => cloneBoundedChatLogPayload(item, depth + 1)); - if (value.length > maxTailItems) { - return [ - { - _omniroute_truncated_array: true, - originalLength: value.length, - retainedTailItems: maxTailItems, - }, - ...cloned, - ]; - } - return cloned; - } - - const result: Record = {}; - const entries = Object.entries(value as Record); - const maxKeys = getChatLogMaxObjectKeys(); - for (const [key, item] of maxKeys > 0 ? entries.slice(0, maxKeys) : entries) { - result[key] = cloneBoundedChatLogPayload(item, depth + 1); - } - if (maxKeys > 0 && entries.length > maxKeys) { - result._omniroute_truncated_keys = entries.length - maxKeys; - } - return result; -} - -import { estimateSizeFast, isSmallEnoughForSemanticCache } from "../utils/estimateSize.ts"; - -const MAX_LOG_BODY_CHARS = 8 * 1024; // 8KB cap for logged request/response bodies -/** - * Truncate a large object for logging. If its JSON representation exceeds - * MAX_LOG_BODY_CHARS, return a lightweight summary instead of the full clone. - * This prevents persistAttemptLogs from holding multi-MB references to - * translatedBody across 17 call sites per request. - */ -function truncateForLog(value: unknown): Record | null | undefined { - if (value === null || value === undefined) return value as null | undefined; - if (typeof value !== "object") return value as unknown as Record; - const estimatedSize = estimateSizeFast(value); - if (estimatedSize <= MAX_LOG_BODY_CHARS) return value as Record; - // Object is too large — return a summary instead of a deep clone - const obj = value as Record; - const summary: Record = { - _truncated: true, - _originalBytes: estimatedSize, - }; - if (typeof obj.model === "string") summary.model = obj.model; - if (typeof obj.provider === "string") summary.provider = obj.provider; - if (Array.isArray(obj.messages)) summary.messageCount = obj.messages.length; - if (Array.isArray(obj.contents)) summary.contentCount = obj.contents.length; - if (typeof obj.stream === "boolean") summary.stream = obj.stream; - return summary; -} - -function extractMemoryTextFromResponse( - response: Record | null | undefined -): string { - if (!response || typeof response !== "object") return ""; - - const openAIText = response?.choices?.[0]?.message?.content; - if (typeof openAIText === "string") { - return capMemoryExtractionText(openAIText.trim()); - } - - if (Array.isArray(response?.content)) { - const contentText = response.content - .filter( - (part: Record) => part?.type === "text" && typeof part?.text === "string" - ) - .map((part: Record) => String(part.text).trim()) - .filter(Boolean) - .join("\n"); - if (contentText) return capMemoryExtractionText(contentText); - } - - if (typeof response?.output_text === "string") { - return capMemoryExtractionText(response.output_text.trim()); - } - - return ""; -} - -function extractMemoryTextFromRequestBody( - body: Record | null | undefined -): string { - if (!body || typeof body !== "object") return ""; - - const messages = Array.isArray(body.messages) ? body.messages : null; - if (messages && messages.length > 0) { - for (let i = messages.length - 1; i >= 0; i -= 1) { - const msg = messages[i] as Record; - if (msg?.role !== "user") continue; - - if (typeof msg.content === "string" && msg.content.trim().length > 0) { - return capMemoryExtractionText(msg.content.trim()); - } - - if (Array.isArray(msg.content)) { - const text = msg.content - .map((part: Record) => { - if (typeof part?.text === "string") return part.text.trim(); - if (part?.type === "input_text" && typeof part?.text === "string") - return part.text.trim(); - return ""; - }) - .filter(Boolean) - .join("\n") - .trim(); - if (text) return capMemoryExtractionText(text); - } - } - } - - const input = Array.isArray(body.input) ? body.input : null; - if (input && input.length > 0) { - for (let i = input.length - 1; i >= 0; i -= 1) { - const item = input[i] as Record; - const role = typeof item?.role === "string" ? item.role.trim().toLowerCase() : ""; - const itemType = typeof item?.type === "string" ? item.type.trim().toLowerCase() : ""; - if (role && role !== "user") continue; - if (itemType && itemType !== "message") continue; - - if (typeof item?.content === "string" && item.content.trim()) { - return capMemoryExtractionText(item.content.trim()); - } - if (Array.isArray(item?.content)) { - const text = item.content - .map((part: Record) => { - if (typeof part?.text === "string") return part.text.trim(); - if (part?.type === "input_text" && typeof part?.text === "string") - return part.text.trim(); - return ""; - }) - .filter(Boolean) - .join("\n") - .trim(); - if (text) return capMemoryExtractionText(text); - } - } - - const tailChunks: string[] = []; - let tailLength = 0; - for (let i = input.length - 1; i >= 0 && tailLength < MEMORY_EXTRACTION_TEXT_LIMIT; i -= 1) { - const item = input[i] as Record; - const text = (() => { - const role = typeof item?.role === "string" ? item.role.trim().toLowerCase() : ""; - const itemType = typeof item?.type === "string" ? item.type.trim().toLowerCase() : ""; - if (role && role !== "user") return ""; - if (itemType && itemType !== "message") return ""; - - if (typeof item?.content === "string") return item.content.trim(); - if (Array.isArray(item?.content)) { - return item.content - .map((part: Record) => { - if (typeof part?.text === "string") return part.text.trim(); - if (part?.type === "input_text" && typeof part?.text === "string") - return part.text.trim(); - return ""; - }) - .filter(Boolean) - .join("\n") - .trim(); - } - return ""; - })(); - if (!text) continue; - tailChunks.unshift(text); - tailLength += text.length + 1; - } - const chunks = tailChunks.join("\n").trim(); - if (chunks) return capMemoryExtractionText(chunks); - } - - return ""; -} +import { isSmallEnoughForSemanticCache } from "../utils/estimateSize.ts"; async function forwardDashboardEventToLiveWs(event: string, payload: unknown): Promise { const port = process.env.LIVE_WS_PORT || "20129"; @@ -489,14 +300,6 @@ async function maybeSyncClaudeExtraUsageState({ } } -function resolveMemoryOwnerId(apiKeyInfo: Record | null): string | null { - const rawId = apiKeyInfo?.id; - if (typeof rawId === "string" && rawId.trim().length > 0) { - return rawId; - } - return null; -} - export function shouldUseNativeCodexPassthrough({ provider, sourceFormat, @@ -577,70 +380,6 @@ export function isClaudeCodeSemanticPassthroughRequest({ return typeof sessionId === "string" && sessionId.trim().length > 0; } -function buildClaudePassthroughToolNameMap(body: Record | null | undefined) { - if (!body || !Array.isArray(body.tools)) return null; - - const toolNameMap = new Map(); - for (const tool of body.tools) { - const toolRecord = tool as Record; - const toolData = - toolRecord?.type === "function" && - toolRecord.function && - typeof toolRecord.function === "object" - ? (toolRecord.function as Record) - : toolRecord; - const originalName = typeof toolData?.name === "string" ? toolData.name.trim() : ""; - if (!originalName) continue; - toolNameMap.set(`${CLAUDE_OAUTH_TOOL_PREFIX}${originalName}`, originalName); - } - - return toolNameMap.size > 0 ? toolNameMap : null; -} - -function restoreClaudePassthroughToolNames( - responseBody: Record, - toolNameMap: Map | null -) { - if (!toolNameMap || !Array.isArray(responseBody?.content)) return responseBody; - - let changed = false; - const content = responseBody.content.map((block: Record) => { - if (block?.type !== "tool_use" || typeof block?.name !== "string") return block; - const restoredName = toolNameMap.get(block.name) ?? block.name; - if (restoredName === block.name) return block; - changed = true; - return { - ...block, - name: restoredName, - }; - }); - - if (!changed) return responseBody; - return { - ...responseBody, - content, - }; -} - -function mergeResponseToolNameMap( - baseToolNameMap: Map | null, - transformedBody: Record | null | undefined -) { - const executorToolNameMap = - transformedBody && transformedBody._toolNameMap instanceof Map - ? (transformedBody._toolNameMap as Map) - : null; - - if (!executorToolNameMap?.size) return baseToolNameMap; - if (!baseToolNameMap?.size) return executorToolNameMap; - - const merged = new Map(baseToolNameMap); - for (const [toolName, originalName] of executorToolNameMap.entries()) { - merged.set(toolName, originalName); - } - return merged; -} - const STREAMING_RESPONSE_HEADER_DENYLIST = new Set([ "content-type", "content-encoding", @@ -717,292 +456,6 @@ function getSkillsModelIdForFormat(format: string): string { } } -function parseNonStreamingSSEPayload( - rawBody: string, - preferredFormat: string, - fallbackModel: string -): { body: Record; format: string } | null { - const formatsToTry: string[] = []; - const seen = new Set(); - const queueFormat = (format: string) => { - if (!format || seen.has(format)) return; - seen.add(format); - formatsToTry.push(format); - }; - - queueFormat(preferredFormat); - queueFormat(FORMATS.OPENAI_RESPONSES); - queueFormat(FORMATS.CLAUDE); - queueFormat(FORMATS.OPENAI); - - for (const format of formatsToTry) { - const parsed = - format === FORMATS.OPENAI_RESPONSES - ? parseSSEToResponsesOutput(rawBody, fallbackModel) - : format === FORMATS.CLAUDE - ? parseSSEToClaudeResponse(rawBody, fallbackModel) - : parseSSEToOpenAIResponse(rawBody, fallbackModel); - if (parsed && typeof parsed === "object") { - return { - body: parsed as Record, - format, - }; - } - } - - return null; -} - -function convertNDJSONToSSE(rawBody: string): string { - const chunks = String(rawBody || "") - .split(/\r?\n/) - .map((line) => line.trim()) - .filter((line) => line.length > 0); - - if (chunks.length === 0) return rawBody; - - return `${chunks.map((chunk) => `data: ${chunk}\n`).join("\n")}\n`; -} - -function normalizeNonStreamingEventPayload(rawBody: string, contentType: string): string { - if (contentType.includes("application/x-ndjson")) { - return convertNDJSONToSSE(rawBody); - } - return rawBody; -} - -function isTruthyStreamBody(body: unknown): boolean { - return !!body && typeof body === "object" && (body as { stream?: unknown }).stream === true; -} - -function isEventStreamAccepted(headers: Record | Headers | null | undefined) { - return (getHeaderValueCaseInsensitive(headers, "accept") || "") - .toLowerCase() - .includes("text/event-stream"); -} - -function shouldTreatBufferedEventResponseAsExpected( - upstreamStream: boolean, - providerHeaders: Record | Headers | null | undefined, - finalBody: unknown -): boolean { - return upstreamStream || isEventStreamAccepted(providerHeaders) || isTruthyStreamBody(finalBody); -} - -const NON_STREAMING_SSE_TERMINAL_TYPES = new Set([ - "message_stop", - "response.completed", - "response.done", - "response.cancelled", - "response.canceled", - "response.failed", - "response.incomplete", -]); - -type NonStreamingSseTerminalState = { - currentEvent: string; - pendingLine: string; -}; - -function processNonStreamingSseTerminalLine( - state: NonStreamingSseTerminalState, - rawLine: string -): boolean { - const trimmed = rawLine.trim(); - if (!trimmed || trimmed.startsWith(":")) { - if (!trimmed) state.currentEvent = ""; - return false; - } - - if (trimmed.startsWith("event:")) { - state.currentEvent = trimmed.slice(6).trim(); - return false; - } - - if (!trimmed.startsWith("data:")) return false; - const data = trimmed.slice(5).trim(); - if (data === "[DONE]") return true; - if (!data) return false; - - try { - const parsed = JSON.parse(data); - const eventType = - parsed && typeof parsed === "object" && typeof parsed.type === "string" - ? parsed.type - : state.currentEvent; - return NON_STREAMING_SSE_TERMINAL_TYPES.has(eventType); - } catch { - // Keep reading malformed data so the parser can report a useful upstream error. - return false; - } -} - -function appendNonStreamingSseTerminalSignal( - state: NonStreamingSseTerminalState, - chunk: string -): boolean { - const lines = `${state.pendingLine}${chunk}`.split(/\r?\n/); - state.pendingLine = lines.pop() ?? ""; - - for (const rawLine of lines) { - if (processNonStreamingSseTerminalLine(state, rawLine)) return true; - } - - return false; -} - -function createBodyTimeoutError(timeoutMs: number): Error { - const err = new Error(`Response body read timeout after ${timeoutMs}ms`); - err.name = "BodyTimeoutError"; - return err; -} - -function readStreamChunkWithTimeout( - reader: ReadableStreamDefaultReader, - timeoutMs: number -): Promise<{ done: boolean; value?: Uint8Array }> { - if (timeoutMs <= 0) return reader.read(); - - return new Promise((resolve, reject) => { - const timeout = setTimeout(() => reject(createBodyTimeoutError(timeoutMs)), timeoutMs); - reader.read().then( - (value) => { - clearTimeout(timeout); - resolve(value); - }, - (error) => { - clearTimeout(timeout); - reject(error); - } - ); - }); -} - -function createUpstreamStartTimeoutError( - timeoutMs: number, - provider: string, - model: string -): Error { - const err = new Error( - `Upstream request did not return response headers after ${timeoutMs}ms (${provider}/${model})` - ); - err.name = "TimeoutError"; - return err; -} - -function createAbortError(signal: AbortSignal): Error { - const reason = signal.reason; - if (reason instanceof Error) return reason; - const err = new Error(typeof reason === "string" ? reason : "The operation was aborted"); - err.name = "AbortError"; - return err; -} - -/** Billable token total — mirrors the columns persisted by saveRequestUsage so the - * live token-limit counter stays consistent with usage_history seed-on-miss. */ -function computeBillableTokens(usage: unknown): number { - // Cache read/creation tokens are a BREAKDOWN already contained inside - // getLoggedInputTokens (prompt_tokens / input_tokens). Adding them here would - // double-count. Canonical billable total = input + output + reasoning, matching - // the columns persisted by saveRequestUsage and seedWindowUsageFromHistory. - return getLoggedInputTokens(usage) + getLoggedOutputTokens(usage) + getReasoningTokens(usage); -} - -function getExecutorTimeoutMs(executor: unknown): number { - const getTimeoutMs = (executor as { getTimeoutMs?: () => unknown } | null)?.getTimeoutMs; - if (typeof getTimeoutMs !== "function") return FETCH_TIMEOUT_MS; - - try { - const timeoutMs = getTimeoutMs.call(executor); - if (typeof timeoutMs !== "number" || !Number.isFinite(timeoutMs)) return FETCH_TIMEOUT_MS; - return Math.max(0, Math.floor(timeoutMs)); - } catch { - return FETCH_TIMEOUT_MS; - } -} - -function normalizeExecutorResult( - result: - | Response - | { - response: Response; - url?: string; - headers?: Record; - transformedBody?: unknown; - } -): { response: Response; url: string; headers: Record; transformedBody: unknown } { - if (result instanceof Response) { - return { response: result, url: "", headers: {}, transformedBody: null }; - } - return { - response: result.response, - url: result.url || "", - headers: result.headers || {}, - transformedBody: result.transformedBody ?? null, - }; -} - -async function executeWithUpstreamStartTimeout({ - executor, - provider, - model, - signal, - log, - execute, -}: { - executor: unknown; - provider: string; - model: string; - signal: AbortSignal; - log?: { warn?: (tag: string, message: string) => void } | null; - execute: (signal: AbortSignal) => Promise; -}): Promise { - const timeoutMs = getExecutorTimeoutMs(executor); - if (timeoutMs <= 0) return execute(signal); - if (signal.aborted) throw createAbortError(signal); - - const timeoutController = new AbortController(); - const combinedController = new AbortController(); - const timeoutError = createUpstreamStartTimeoutError(timeoutMs, provider, model); - - let timeoutId: ReturnType | null = null; - let abortListener: (() => void) | null = null; - let timeoutAbortListener: (() => void) | null = null; - - const abortCombined = (source: AbortSignal) => { - if (combinedController.signal.aborted) return; - const reason = source.reason instanceof Error ? source.reason : createAbortError(source); - combinedController.abort(reason); - }; - - abortListener = () => abortCombined(signal); - timeoutAbortListener = () => abortCombined(timeoutController.signal); - signal.addEventListener("abort", abortListener, { once: true }); - timeoutController.signal.addEventListener("abort", timeoutAbortListener, { once: true }); - - const timeoutPromise = new Promise((_, reject) => { - timeoutId = setTimeout(() => { - log?.warn?.("TIMEOUT", timeoutError.message); - timeoutController.abort(timeoutError); - reject(timeoutError); - }, timeoutMs); - }); - - const abortPromise = new Promise((_, reject) => { - signal.addEventListener("abort", () => reject(createAbortError(signal)), { once: true }); - }); - - try { - return await Promise.race([execute(combinedController.signal), timeoutPromise, abortPromise]); - } finally { - if (timeoutId) clearTimeout(timeoutId); - if (abortListener) signal.removeEventListener("abort", abortListener); - if (timeoutAbortListener) { - timeoutController.signal.removeEventListener("abort", timeoutAbortListener); - } - } -} - /** * Strip hop-by-hop headers that describe the upstream wire encoding. * @@ -1082,23 +535,6 @@ async function readNonStreamingResponseBody( return rawBody; } -function getHeaderValueCaseInsensitive( - headers: Record | Headers | null | undefined, - targetName: string -) { - if (!headers || typeof headers !== "object") return null; - if (headers instanceof Headers) { - return headers.get(targetName); - } - const lowered = targetName.toLowerCase(); - for (const [key, value] of Object.entries(headers)) { - if (key.toLowerCase() === lowered && typeof value === "string" && value.trim()) { - return value.trim(); - } - } - return null; -} - function toFiniteNumberOrNull(value: unknown): number | null { if (typeof value === "number" && Number.isFinite(value)) { return value; diff --git a/open-sse/handlers/chatCore/headers.ts b/open-sse/handlers/chatCore/headers.ts new file mode 100644 index 0000000000..9e79a67ec4 --- /dev/null +++ b/open-sse/handlers/chatCore/headers.ts @@ -0,0 +1,16 @@ +export function getHeaderValueCaseInsensitive( + headers: Record | Headers | null | undefined, + targetName: string +) { + if (!headers || typeof headers !== "object") return null; + if (headers instanceof Headers) { + return headers.get(targetName); + } + const lowered = targetName.toLowerCase(); + for (const [key, value] of Object.entries(headers)) { + if (key.toLowerCase() === lowered && typeof value === "string" && value.trim()) { + return value.trim(); + } + } + return null; +} diff --git a/open-sse/handlers/chatCore/logTruncation.ts b/open-sse/handlers/chatCore/logTruncation.ts new file mode 100644 index 0000000000..5f4171948b --- /dev/null +++ b/open-sse/handlers/chatCore/logTruncation.ts @@ -0,0 +1,84 @@ +import { + getChatLogTextLimit, + getChatLogMaxDepth, + getChatLogArrayTailItems, + getChatLogMaxObjectKeys, +} from "@/lib/logEnv"; +import { estimateSizeFast } from "../../utils/estimateSize.ts"; + +export const MEMORY_EXTRACTION_TEXT_LIMIT = 64 * 1024; +const MAX_LOG_BODY_CHARS = 8 * 1024; // 8KB cap for logged request/response bodies + +export function capMemoryExtractionText(value: string): string { + if (value.length <= MEMORY_EXTRACTION_TEXT_LIMIT) return value; + return value.slice(-MEMORY_EXTRACTION_TEXT_LIMIT); +} + +export function truncateChatLogText(value: string): string { + const limit = getChatLogTextLimit(); + if (value.length <= limit) return value; + const head = value.slice(0, Math.floor(limit / 2)); + const tail = value.slice(-Math.ceil(limit / 2)); + return `${head}\n[...truncated ${value.length - limit} chars...]\n${tail}`; +} + +export function cloneBoundedChatLogPayload(value: unknown, depth = 0): unknown { + if (value === null || value === undefined) return value; + if (typeof value === "string") return truncateChatLogText(value); + if (typeof value !== "object") return value; + if (depth >= getChatLogMaxDepth()) return "[MaxDepth]"; + + const maxTailItems = getChatLogArrayTailItems(); + + if (Array.isArray(value)) { + const retained = value.length > maxTailItems ? value.slice(-maxTailItems) : value; + const cloned = retained.map((item) => cloneBoundedChatLogPayload(item, depth + 1)); + if (value.length > maxTailItems) { + return [ + { + _omniroute_truncated_array: true, + originalLength: value.length, + retainedTailItems: maxTailItems, + }, + ...cloned, + ]; + } + return cloned; + } + + const result: Record = {}; + const entries = Object.entries(value as Record); + const maxKeys = getChatLogMaxObjectKeys(); + for (const [key, item] of maxKeys > 0 ? entries.slice(0, maxKeys) : entries) { + result[key] = cloneBoundedChatLogPayload(item, depth + 1); + } + if (maxKeys > 0 && entries.length > maxKeys) { + result._omniroute_truncated_keys = entries.length - maxKeys; + } + return result; +} + +/** + * Truncate a large object for logging. If its JSON representation exceeds + * MAX_LOG_BODY_CHARS, return a lightweight summary instead of the full clone. + * This prevents persistAttemptLogs from holding multi-MB references to + * translatedBody across 17 call sites per request. + */ +export function truncateForLog(value: unknown): Record | null | undefined { + if (value === null || value === undefined) return value as null | undefined; + if (typeof value !== "object") return value as unknown as Record; + const estimatedSize = estimateSizeFast(value); + if (estimatedSize <= MAX_LOG_BODY_CHARS) return value as Record; + // Object is too large — return a summary instead of a deep clone + const obj = value as Record; + const summary: Record = { + _truncated: true, + _originalBytes: estimatedSize, + }; + if (typeof obj.model === "string") summary.model = obj.model; + if (typeof obj.provider === "string") summary.provider = obj.provider; + if (Array.isArray(obj.messages)) summary.messageCount = obj.messages.length; + if (Array.isArray(obj.contents)) summary.contentCount = obj.contents.length; + if (typeof obj.stream === "boolean") summary.stream = obj.stream; + return summary; +} diff --git a/open-sse/handlers/chatCore/memoryExtraction.ts b/open-sse/handlers/chatCore/memoryExtraction.ts new file mode 100644 index 0000000000..7ca6f66da6 --- /dev/null +++ b/open-sse/handlers/chatCore/memoryExtraction.ts @@ -0,0 +1,131 @@ +import { capMemoryExtractionText, MEMORY_EXTRACTION_TEXT_LIMIT } from "./logTruncation.ts"; + +export function extractMemoryTextFromResponse( + response: Record | null | undefined +): string { + if (!response || typeof response !== "object") return ""; + + const openAIText = response?.choices?.[0]?.message?.content; + if (typeof openAIText === "string") { + return capMemoryExtractionText(openAIText.trim()); + } + + if (Array.isArray(response?.content)) { + const contentText = response.content + .filter( + (part: Record) => part?.type === "text" && typeof part?.text === "string" + ) + .map((part: Record) => String(part.text).trim()) + .filter(Boolean) + .join("\n"); + if (contentText) return capMemoryExtractionText(contentText); + } + + if (typeof response?.output_text === "string") { + return capMemoryExtractionText(response.output_text.trim()); + } + + return ""; +} + +export function extractMemoryTextFromRequestBody( + body: Record | null | undefined +): string { + if (!body || typeof body !== "object") return ""; + + const messages = Array.isArray(body.messages) ? body.messages : null; + if (messages && messages.length > 0) { + for (let i = messages.length - 1; i >= 0; i -= 1) { + const msg = messages[i] as Record; + if (msg?.role !== "user") continue; + + if (typeof msg.content === "string" && msg.content.trim().length > 0) { + return capMemoryExtractionText(msg.content.trim()); + } + + if (Array.isArray(msg.content)) { + const text = msg.content + .map((part: Record) => { + if (typeof part?.text === "string") return part.text.trim(); + if (part?.type === "input_text" && typeof part?.text === "string") + return part.text.trim(); + return ""; + }) + .filter(Boolean) + .join("\n") + .trim(); + if (text) return capMemoryExtractionText(text); + } + } + } + + const input = Array.isArray(body.input) ? body.input : null; + if (input && input.length > 0) { + for (let i = input.length - 1; i >= 0; i -= 1) { + const item = input[i] as Record; + const role = typeof item?.role === "string" ? item.role.trim().toLowerCase() : ""; + const itemType = typeof item?.type === "string" ? item.type.trim().toLowerCase() : ""; + if (role && role !== "user") continue; + if (itemType && itemType !== "message") continue; + + if (typeof item?.content === "string" && item.content.trim()) { + return capMemoryExtractionText(item.content.trim()); + } + if (Array.isArray(item?.content)) { + const text = item.content + .map((part: Record) => { + if (typeof part?.text === "string") return part.text.trim(); + if (part?.type === "input_text" && typeof part?.text === "string") + return part.text.trim(); + return ""; + }) + .filter(Boolean) + .join("\n") + .trim(); + if (text) return capMemoryExtractionText(text); + } + } + + const tailChunks: string[] = []; + let tailLength = 0; + for (let i = input.length - 1; i >= 0 && tailLength < MEMORY_EXTRACTION_TEXT_LIMIT; i -= 1) { + const item = input[i] as Record; + const text = (() => { + const role = typeof item?.role === "string" ? item.role.trim().toLowerCase() : ""; + const itemType = typeof item?.type === "string" ? item.type.trim().toLowerCase() : ""; + if (role && role !== "user") return ""; + if (itemType && itemType !== "message") return ""; + + if (typeof item?.content === "string") return item.content.trim(); + if (Array.isArray(item?.content)) { + return item.content + .map((part: Record) => { + if (typeof part?.text === "string") return part.text.trim(); + if (part?.type === "input_text" && typeof part?.text === "string") + return part.text.trim(); + return ""; + }) + .filter(Boolean) + .join("\n") + .trim(); + } + return ""; + })(); + if (!text) continue; + tailChunks.unshift(text); + tailLength += text.length + 1; + } + const chunks = tailChunks.join("\n").trim(); + if (chunks) return capMemoryExtractionText(chunks); + } + + return ""; +} + +export function resolveMemoryOwnerId(apiKeyInfo: Record | null): string | null { + const rawId = apiKeyInfo?.id; + if (typeof rawId === "string" && rawId.trim().length > 0) { + return rawId; + } + return null; +} diff --git a/open-sse/handlers/chatCore/nonStreamingSse.ts b/open-sse/handlers/chatCore/nonStreamingSse.ts new file mode 100644 index 0000000000..eb513cf7f2 --- /dev/null +++ b/open-sse/handlers/chatCore/nonStreamingSse.ts @@ -0,0 +1,141 @@ +import { FORMATS } from "../../translator/formats.ts"; +import { + parseSSEToResponsesOutput, + parseSSEToClaudeResponse, + parseSSEToOpenAIResponse, +} from "../sseParser.ts"; +import { getHeaderValueCaseInsensitive } from "./headers.ts"; + +export function parseNonStreamingSSEPayload( + rawBody: string, + preferredFormat: string, + fallbackModel: string +): { body: Record; format: string } | null { + const formatsToTry: string[] = []; + const seen = new Set(); + const queueFormat = (format: string) => { + if (!format || seen.has(format)) return; + seen.add(format); + formatsToTry.push(format); + }; + + queueFormat(preferredFormat); + queueFormat(FORMATS.OPENAI_RESPONSES); + queueFormat(FORMATS.CLAUDE); + queueFormat(FORMATS.OPENAI); + + for (const format of formatsToTry) { + const parsed = + format === FORMATS.OPENAI_RESPONSES + ? parseSSEToResponsesOutput(rawBody, fallbackModel) + : format === FORMATS.CLAUDE + ? parseSSEToClaudeResponse(rawBody, fallbackModel) + : parseSSEToOpenAIResponse(rawBody, fallbackModel); + if (parsed && typeof parsed === "object") { + return { + body: parsed as Record, + format, + }; + } + } + + return null; +} + +export function convertNDJSONToSSE(rawBody: string): string { + const chunks = String(rawBody || "") + .split(/\r?\n/) + .map((line) => line.trim()) + .filter((line) => line.length > 0); + + if (chunks.length === 0) return rawBody; + + return `${chunks.map((chunk) => `data: ${chunk}\n`).join("\n")}\n`; +} + +export function normalizeNonStreamingEventPayload(rawBody: string, contentType: string): string { + if (contentType.includes("application/x-ndjson")) { + return convertNDJSONToSSE(rawBody); + } + return rawBody; +} + +export function isTruthyStreamBody(body: unknown): boolean { + return !!body && typeof body === "object" && (body as { stream?: unknown }).stream === true; +} + +export function isEventStreamAccepted(headers: Record | Headers | null | undefined) { + return (getHeaderValueCaseInsensitive(headers, "accept") || "") + .toLowerCase() + .includes("text/event-stream"); +} + +export function shouldTreatBufferedEventResponseAsExpected( + upstreamStream: boolean, + providerHeaders: Record | Headers | null | undefined, + finalBody: unknown +): boolean { + return upstreamStream || isEventStreamAccepted(providerHeaders) || isTruthyStreamBody(finalBody); +} + +const NON_STREAMING_SSE_TERMINAL_TYPES = new Set([ + "message_stop", + "response.completed", + "response.done", + "response.cancelled", + "response.canceled", + "response.failed", + "response.incomplete", +]); + +export type NonStreamingSseTerminalState = { + currentEvent: string; + pendingLine: string; +}; + +function processNonStreamingSseTerminalLine( + state: NonStreamingSseTerminalState, + rawLine: string +): boolean { + const trimmed = rawLine.trim(); + if (!trimmed || trimmed.startsWith(":")) { + if (!trimmed) state.currentEvent = ""; + return false; + } + + if (trimmed.startsWith("event:")) { + state.currentEvent = trimmed.slice(6).trim(); + return false; + } + + if (!trimmed.startsWith("data:")) return false; + const data = trimmed.slice(5).trim(); + if (data === "[DONE]") return true; + if (!data) return false; + + try { + const parsed = JSON.parse(data); + const eventType = + parsed && typeof parsed === "object" && typeof parsed.type === "string" + ? parsed.type + : state.currentEvent; + return NON_STREAMING_SSE_TERMINAL_TYPES.has(eventType); + } catch { + // Keep reading malformed data so the parser can report a useful upstream error. + return false; + } +} + +export function appendNonStreamingSseTerminalSignal( + state: NonStreamingSseTerminalState, + chunk: string +): boolean { + const lines = `${state.pendingLine}${chunk}`.split(/\r?\n/); + state.pendingLine = lines.pop() ?? ""; + + for (const rawLine of lines) { + if (processNonStreamingSseTerminalLine(state, rawLine)) return true; + } + + return false; +} diff --git a/open-sse/handlers/chatCore/passthroughToolNames.ts b/open-sse/handlers/chatCore/passthroughToolNames.ts new file mode 100644 index 0000000000..0ab6b17d4d --- /dev/null +++ b/open-sse/handlers/chatCore/passthroughToolNames.ts @@ -0,0 +1,65 @@ +import { CLAUDE_OAUTH_TOOL_PREFIX } from "../../translator/request/openai-to-claude.ts"; + +export function buildClaudePassthroughToolNameMap(body: Record | null | undefined) { + if (!body || !Array.isArray(body.tools)) return null; + + const toolNameMap = new Map(); + for (const tool of body.tools) { + const toolRecord = tool as Record; + const toolData = + toolRecord?.type === "function" && + toolRecord.function && + typeof toolRecord.function === "object" + ? (toolRecord.function as Record) + : toolRecord; + const originalName = typeof toolData?.name === "string" ? toolData.name.trim() : ""; + if (!originalName) continue; + toolNameMap.set(`${CLAUDE_OAUTH_TOOL_PREFIX}${originalName}`, originalName); + } + + return toolNameMap.size > 0 ? toolNameMap : null; +} + +export function restoreClaudePassthroughToolNames( + responseBody: Record, + toolNameMap: Map | null +) { + if (!toolNameMap || !Array.isArray(responseBody?.content)) return responseBody; + + let changed = false; + const content = responseBody.content.map((block: Record) => { + if (block?.type !== "tool_use" || typeof block?.name !== "string") return block; + const restoredName = toolNameMap.get(block.name) ?? block.name; + if (restoredName === block.name) return block; + changed = true; + return { + ...block, + name: restoredName, + }; + }); + + if (!changed) return responseBody; + return { + ...responseBody, + content, + }; +} + +export function mergeResponseToolNameMap( + baseToolNameMap: Map | null, + transformedBody: Record | null | undefined +) { + const executorToolNameMap = + transformedBody && transformedBody._toolNameMap instanceof Map + ? (transformedBody._toolNameMap as Map) + : null; + + if (!executorToolNameMap?.size) return baseToolNameMap; + if (!baseToolNameMap?.size) return executorToolNameMap; + + const merged = new Map(baseToolNameMap); + for (const [toolName, originalName] of executorToolNameMap.entries()) { + merged.set(toolName, originalName); + } + return merged; +} diff --git a/open-sse/handlers/chatCore/upstreamTimeouts.ts b/open-sse/handlers/chatCore/upstreamTimeouts.ts new file mode 100644 index 0000000000..6f3636bb1f --- /dev/null +++ b/open-sse/handlers/chatCore/upstreamTimeouts.ts @@ -0,0 +1,158 @@ +import { FETCH_TIMEOUT_MS } from "../../config/constants.ts"; +import { + getLoggedInputTokens, + getLoggedOutputTokens, + getReasoningTokens, +} from "@/lib/usage/tokenAccounting"; + +export function createBodyTimeoutError(timeoutMs: number): Error { + const err = new Error(`Response body read timeout after ${timeoutMs}ms`); + err.name = "BodyTimeoutError"; + return err; +} + +export function readStreamChunkWithTimeout( + reader: ReadableStreamDefaultReader, + timeoutMs: number +): Promise<{ done: boolean; value?: Uint8Array }> { + if (timeoutMs <= 0) return reader.read(); + + return new Promise((resolve, reject) => { + const timeout = setTimeout(() => reject(createBodyTimeoutError(timeoutMs)), timeoutMs); + reader.read().then( + (value) => { + clearTimeout(timeout); + resolve(value); + }, + (error) => { + clearTimeout(timeout); + reject(error); + } + ); + }); +} + +export function createUpstreamStartTimeoutError( + timeoutMs: number, + provider: string, + model: string +): Error { + const err = new Error( + `Upstream request did not return response headers after ${timeoutMs}ms (${provider}/${model})` + ); + err.name = "TimeoutError"; + return err; +} + +export function createAbortError(signal: AbortSignal): Error { + const reason = signal.reason; + if (reason instanceof Error) return reason; + const err = new Error(typeof reason === "string" ? reason : "The operation was aborted"); + err.name = "AbortError"; + return err; +} + +/** Billable token total — mirrors the columns persisted by saveRequestUsage so the + * live token-limit counter stays consistent with usage_history seed-on-miss. */ +export function computeBillableTokens(usage: unknown): number { + // Cache read/creation tokens are a BREAKDOWN already contained inside + // getLoggedInputTokens (prompt_tokens / input_tokens). Adding them here would + // double-count. Canonical billable total = input + output + reasoning, matching + // the columns persisted by saveRequestUsage and seedWindowUsageFromHistory. + return getLoggedInputTokens(usage) + getLoggedOutputTokens(usage) + getReasoningTokens(usage); +} + +export function getExecutorTimeoutMs(executor: unknown): number { + const getTimeoutMs = (executor as { getTimeoutMs?: () => unknown } | null)?.getTimeoutMs; + if (typeof getTimeoutMs !== "function") return FETCH_TIMEOUT_MS; + + try { + const timeoutMs = getTimeoutMs.call(executor); + if (typeof timeoutMs !== "number" || !Number.isFinite(timeoutMs)) return FETCH_TIMEOUT_MS; + return Math.max(0, Math.floor(timeoutMs)); + } catch { + return FETCH_TIMEOUT_MS; + } +} + +export function normalizeExecutorResult( + result: + | Response + | { + response: Response; + url?: string; + headers?: Record; + transformedBody?: unknown; + } +): { response: Response; url: string; headers: Record; transformedBody: unknown } { + if (result instanceof Response) { + return { response: result, url: "", headers: {}, transformedBody: null }; + } + return { + response: result.response, + url: result.url || "", + headers: result.headers || {}, + transformedBody: result.transformedBody ?? null, + }; +} + +export async function executeWithUpstreamStartTimeout({ + executor, + provider, + model, + signal, + log, + execute, +}: { + executor: unknown; + provider: string; + model: string; + signal: AbortSignal; + log?: { warn?: (tag: string, message: string) => void } | null; + execute: (signal: AbortSignal) => Promise; +}): Promise { + const timeoutMs = getExecutorTimeoutMs(executor); + if (timeoutMs <= 0) return execute(signal); + if (signal.aborted) throw createAbortError(signal); + + const timeoutController = new AbortController(); + const combinedController = new AbortController(); + const timeoutError = createUpstreamStartTimeoutError(timeoutMs, provider, model); + + let timeoutId: ReturnType | null = null; + let abortListener: (() => void) | null = null; + let timeoutAbortListener: (() => void) | null = null; + + const abortCombined = (source: AbortSignal) => { + if (combinedController.signal.aborted) return; + const reason = source.reason instanceof Error ? source.reason : createAbortError(source); + combinedController.abort(reason); + }; + + abortListener = () => abortCombined(signal); + timeoutAbortListener = () => abortCombined(timeoutController.signal); + signal.addEventListener("abort", abortListener, { once: true }); + timeoutController.signal.addEventListener("abort", timeoutAbortListener, { once: true }); + + const timeoutPromise = new Promise((_, reject) => { + timeoutId = setTimeout(() => { + log?.warn?.("TIMEOUT", timeoutError.message); + timeoutController.abort(timeoutError); + reject(timeoutError); + }, timeoutMs); + }); + + const abortPromise = new Promise((_, reject) => { + signal.addEventListener("abort", () => reject(createAbortError(signal)), { once: true }); + }); + + try { + return await Promise.race([execute(combinedController.signal), timeoutPromise, abortPromise]); + } finally { + if (timeoutId) clearTimeout(timeoutId); + if (abortListener) signal.removeEventListener("abort", abortListener); + if (timeoutAbortListener) { + timeoutController.signal.removeEventListener("abort", timeoutAbortListener); + } + } +} diff --git a/open-sse/services/tokenLimitCounter.ts b/open-sse/services/tokenLimitCounter.ts index ad6fa38e8c..cd69a9385b 100644 --- a/open-sse/services/tokenLimitCounter.ts +++ b/open-sse/services/tokenLimitCounter.ts @@ -50,7 +50,7 @@ export function seedWindowUsageFromHistory(limit: TokenLimit, now = Date.now()): // Canonical billable total = input + output + reasoning. tokens_cache_read and // tokens_cache_creation are a BREAKDOWN already inside tokens_input (see migration // 012_fix_token_input_cache_tokens.sql) — summing them again would double-count. - // This must mirror computeBillableTokens() in chatCore.ts. + // This must mirror computeBillableTokens() in chatCore/upstreamTimeouts.ts. const tokenSum = `COALESCE(SUM( COALESCE(tokens_input, 0) + COALESCE(tokens_output, 0) + COALESCE(tokens_reasoning, 0) diff --git a/tests/unit/chatcore-headers.test.ts b/tests/unit/chatcore-headers.test.ts new file mode 100644 index 0000000000..9858d749e3 --- /dev/null +++ b/tests/unit/chatcore-headers.test.ts @@ -0,0 +1,20 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { getHeaderValueCaseInsensitive } from "../../open-sse/handlers/chatCore/headers.ts"; + +test("getHeaderValueCaseInsensitive reads Headers and plain objects, case-insensitively", () => { + const h = new Headers({ "Content-Type": "text/event-stream" }); + assert.equal(getHeaderValueCaseInsensitive(h, "content-type"), "text/event-stream"); + + const obj = { Accept: "text/event-stream", "X-Foo": "bar" }; + assert.equal(getHeaderValueCaseInsensitive(obj, "accept"), "text/event-stream"); + assert.equal(getHeaderValueCaseInsensitive(obj, "x-foo"), "bar"); + + // plain-object values are trimmed + assert.equal(getHeaderValueCaseInsensitive({ Accept: " v " }, "accept"), "v"); + // missing / non-object -> null + assert.equal(getHeaderValueCaseInsensitive({ Accept: "x" }, "missing"), null); + assert.equal(getHeaderValueCaseInsensitive(null, "accept"), null); + assert.equal(getHeaderValueCaseInsensitive(undefined, "accept"), null); +}); diff --git a/tests/unit/chatcore-log-truncation.test.ts b/tests/unit/chatcore-log-truncation.test.ts new file mode 100644 index 0000000000..704b57e252 --- /dev/null +++ b/tests/unit/chatcore-log-truncation.test.ts @@ -0,0 +1,46 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { + capMemoryExtractionText, + truncateForLog, + cloneBoundedChatLogPayload, + MEMORY_EXTRACTION_TEXT_LIMIT, +} from "../../open-sse/handlers/chatCore/logTruncation.ts"; + +test("capMemoryExtractionText keeps short strings and tail-truncates long ones", () => { + assert.equal(capMemoryExtractionText("hello"), "hello"); + const long = "x".repeat(MEMORY_EXTRACTION_TEXT_LIMIT + 100); + const capped = capMemoryExtractionText(long); + assert.equal(capped.length, MEMORY_EXTRACTION_TEXT_LIMIT); + assert.ok(capped.endsWith("x")); +}); + +test("truncateForLog summarizes oversized objects and passes through small ones", () => { + const small = { model: "gpt-4o", messages: [{ role: "user", content: "hi" }] }; + assert.equal(truncateForLog(small), small); + + const huge = { + model: "gpt-4o", + provider: "openai", + stream: true, + // Use Array.from to create distinct object references so estimateSizeFast + // (which deduplicates via WeakSet) counts every message individually. + messages: Array.from({ length: 50000 }, () => ({ role: "user", content: "x".repeat(64) })), + }; + const summary = truncateForLog(huge) as Record; + assert.equal(summary._truncated, true); + assert.equal(summary.model, "gpt-4o"); + assert.equal(summary.provider, "openai"); + assert.equal(summary.messageCount, 50000); + assert.equal(summary.stream, true); +}); + +test("cloneBoundedChatLogPayload truncates long tail arrays with a marker", () => { + const cloned = cloneBoundedChatLogPayload({ items: new Array(1000).fill("a") }) as { + items: unknown[]; + }; + const marker = cloned.items[0] as Record; + assert.equal(marker._omniroute_truncated_array, true); + assert.equal(marker.originalLength, 1000); +}); diff --git a/tests/unit/chatcore-memory-extraction.test.ts b/tests/unit/chatcore-memory-extraction.test.ts new file mode 100644 index 0000000000..5eece2bf0c --- /dev/null +++ b/tests/unit/chatcore-memory-extraction.test.ts @@ -0,0 +1,42 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { + extractMemoryTextFromResponse, + extractMemoryTextFromRequestBody, + resolveMemoryOwnerId, +} from "../../open-sse/handlers/chatCore/memoryExtraction.ts"; + +test("extractMemoryTextFromResponse reads OpenAI, Claude-array and responses output_text", () => { + assert.equal( + extractMemoryTextFromResponse({ choices: [{ message: { content: " hi " } }] }), + "hi" + ); + assert.equal( + extractMemoryTextFromResponse({ content: [{ type: "text", text: " a " }, { type: "image" }] }), + "a" + ); + assert.equal(extractMemoryTextFromResponse({ output_text: " out " }), "out"); + assert.equal(extractMemoryTextFromResponse(null), ""); +}); + +test("extractMemoryTextFromRequestBody returns the last user message text", () => { + const body = { + messages: [ + { role: "user", content: "first" }, + { role: "assistant", content: "ignored" }, + { role: "user", content: "second" }, + ], + }; + assert.equal(extractMemoryTextFromRequestBody(body), "second"); + const inputBody = { + input: [{ role: "user", type: "message", content: [{ type: "input_text", text: "hey" }] }], + }; + assert.equal(extractMemoryTextFromRequestBody(inputBody), "hey"); +}); + +test("resolveMemoryOwnerId returns trimmed id or null", () => { + assert.equal(resolveMemoryOwnerId({ id: "key_123" }), "key_123"); + assert.equal(resolveMemoryOwnerId({ id: " " }), null); + assert.equal(resolveMemoryOwnerId(null), null); +}); diff --git a/tests/unit/chatcore-non-streaming-sse.test.ts b/tests/unit/chatcore-non-streaming-sse.test.ts new file mode 100644 index 0000000000..2f4965e1c8 --- /dev/null +++ b/tests/unit/chatcore-non-streaming-sse.test.ts @@ -0,0 +1,72 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { + convertNDJSONToSSE, + normalizeNonStreamingEventPayload, + isTruthyStreamBody, + isEventStreamAccepted, + shouldTreatBufferedEventResponseAsExpected, + parseNonStreamingSSEPayload, + appendNonStreamingSseTerminalSignal, + type NonStreamingSseTerminalState, +} from "../../open-sse/handlers/chatCore/nonStreamingSse.ts"; +import { FORMATS } from "../../open-sse/translator/formats.ts"; + +test("convertNDJSONToSSE wraps each non-empty line as a data: frame", () => { + const out = convertNDJSONToSSE('{"a":1}\n{"b":2}\n'); + assert.ok(out.includes('data: {"a":1}\n')); + assert.ok(out.includes('data: {"b":2}\n')); + assert.equal(convertNDJSONToSSE(""), ""); +}); + +test("normalizeNonStreamingEventPayload only converts x-ndjson content", () => { + const raw = '{"a":1}'; + assert.equal(normalizeNonStreamingEventPayload(raw, "application/json"), raw); + assert.notEqual(normalizeNonStreamingEventPayload(raw, "application/x-ndjson"), raw); +}); + +test("stream-body and event-stream acceptance predicates", () => { + assert.equal(isTruthyStreamBody({ stream: true }), true); + assert.equal(isTruthyStreamBody({ stream: false }), false); + assert.equal(isTruthyStreamBody(null), false); + + assert.equal(isEventStreamAccepted({ accept: "text/event-stream" }), true); + assert.equal(isEventStreamAccepted({ accept: "application/json" }), false); + + assert.equal( + shouldTreatBufferedEventResponseAsExpected( + false, + { accept: "application/json" }, + { stream: true } + ), + true + ); + assert.equal( + shouldTreatBufferedEventResponseAsExpected(false, { accept: "application/json" }, {}), + false + ); +}); + +test("appendNonStreamingSseTerminalSignal detects [DONE] and terminal event types", () => { + const done: NonStreamingSseTerminalState = { currentEvent: "", pendingLine: "" }; + assert.equal(appendNonStreamingSseTerminalSignal(done, "data: [DONE]\n"), true); + + const stop: NonStreamingSseTerminalState = { currentEvent: "", pendingLine: "" }; + assert.equal(appendNonStreamingSseTerminalSignal(stop, "event: message_stop\ndata: {}\n"), true); + + const delta: NonStreamingSseTerminalState = { currentEvent: "", pendingLine: "" }; + assert.equal( + appendNonStreamingSseTerminalSignal(delta, 'data: {"type":"content_block_delta"}\n'), + false + ); +}); + +test("parseNonStreamingSSEPayload parses an OpenAI-format SSE buffer", () => { + const raw = + 'data: {"id":"x","choices":[{"delta":{"content":"hi"},"finish_reason":"stop"}]}\n\ndata: [DONE]\n'; + const result = parseNonStreamingSSEPayload(raw, FORMATS.OPENAI, "gpt-4o"); + assert.ok(result !== null); + assert.equal(result?.format, FORMATS.OPENAI); + assert.equal(typeof result?.body, "object"); +}); diff --git a/tests/unit/chatcore-passthrough-tool-names.test.ts b/tests/unit/chatcore-passthrough-tool-names.test.ts new file mode 100644 index 0000000000..255240d1a6 --- /dev/null +++ b/tests/unit/chatcore-passthrough-tool-names.test.ts @@ -0,0 +1,41 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { + buildClaudePassthroughToolNameMap, + restoreClaudePassthroughToolNames, + mergeResponseToolNameMap, +} from "../../open-sse/handlers/chatCore/passthroughToolNames.ts"; +import { CLAUDE_OAUTH_TOOL_PREFIX } from "../../open-sse/translator/request/openai-to-claude.ts"; + +test("buildClaudePassthroughToolNameMap maps prefixed -> original names", () => { + const map = buildClaudePassthroughToolNameMap({ + tools: [{ type: "function", function: { name: "get_weather" } }], + }); + assert.ok(map); + assert.equal(map?.get(`${CLAUDE_OAUTH_TOOL_PREFIX}get_weather`), "get_weather"); + assert.equal(buildClaudePassthroughToolNameMap({ tools: [] }), null); + assert.equal(buildClaudePassthroughToolNameMap(null), null); +}); + +test("restoreClaudePassthroughToolNames rewrites tool_use block names", () => { + const map = new Map([[`${CLAUDE_OAUTH_TOOL_PREFIX}x`, "x"]]); + const restored = restoreClaudePassthroughToolNames( + { content: [{ type: "tool_use", name: `${CLAUDE_OAUTH_TOOL_PREFIX}x` }] }, + map + ) as { content: { name: string }[] }; + assert.equal(restored.content[0].name, "x"); + const body = { content: [{ type: "tool_use", name: "y" }] }; + assert.equal(restoreClaudePassthroughToolNames(body, null), body); +}); + +test("mergeResponseToolNameMap unions base with executor _toolNameMap", () => { + const base = new Map([["a", "1"]]); + const merged = mergeResponseToolNameMap(base, { _toolNameMap: new Map([["b", "2"]]) }) as Map< + string, + string + >; + assert.equal(merged.get("a"), "1"); + assert.equal(merged.get("b"), "2"); + assert.equal(mergeResponseToolNameMap(base, {}), base); +}); diff --git a/tests/unit/chatcore-upstream-timeouts.test.ts b/tests/unit/chatcore-upstream-timeouts.test.ts new file mode 100644 index 0000000000..7fa90b3331 --- /dev/null +++ b/tests/unit/chatcore-upstream-timeouts.test.ts @@ -0,0 +1,51 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { + createBodyTimeoutError, + createUpstreamStartTimeoutError, + createAbortError, + computeBillableTokens, + getExecutorTimeoutMs, + normalizeExecutorResult, +} from "../../open-sse/handlers/chatCore/upstreamTimeouts.ts"; + +test("error factories set name and message", () => { + const body = createBodyTimeoutError(1234); + assert.equal(body.name, "BodyTimeoutError"); + assert.match(body.message, /1234ms/); + + const start = createUpstreamStartTimeoutError(500, "openai", "gpt-4o"); + assert.equal(start.name, "TimeoutError"); + assert.match(start.message, /openai\/gpt-4o/); + + const ctrl = new AbortController(); + ctrl.abort("nope"); + const ab = createAbortError(ctrl.signal); + assert.equal(ab.name, "AbortError"); +}); + +test("computeBillableTokens sums input+output+reasoning (no cache double-count)", () => { + const total = computeBillableTokens({ + prompt_tokens: 10, + completion_tokens: 5, + reasoning_tokens: 2, + }); + assert.equal(total, 17); +}); + +test("getExecutorTimeoutMs floors valid values and falls back to default", () => { + assert.equal(getExecutorTimeoutMs({ getTimeoutMs: () => 1234.9 }), 1234); + assert.equal(getExecutorTimeoutMs({ getTimeoutMs: () => NaN }), getExecutorTimeoutMs(null)); + assert.ok(Number.isFinite(getExecutorTimeoutMs(null))); +}); + +test("normalizeExecutorResult wraps bare Response and passes through rich result", () => { + const r = new Response("x"); + const wrapped = normalizeExecutorResult(r); + assert.equal(wrapped.response, r); + assert.equal(wrapped.url, ""); + const rich = normalizeExecutorResult({ response: r, url: "u", headers: { a: "b" } }); + assert.equal(rich.url, "u"); + assert.equal(rich.headers.a, "b"); +});