From 5c95da552d3b2199fd511f03fb853baf40146041 Mon Sep 17 00:00:00 2001 From: Diego Rodrigues de Sa e Souza <8016841+diegosouzapw@users.noreply.github.com> Date: Tue, 23 Jun 2026 07:17:40 -0300 Subject: [PATCH] refactor(chatCore): extrai persistAttemptLogs para leaf (#3501) (#4717) chatCore #3501: extract persistAttemptLogs to leaf (attemptLogging.ts). Rebased onto release tip post-#4708 (resolved imports conflict: kept tip's resolveCompressionHeader from compression Phase 3, dropped now-unused logTruncation import moved into the leaf). 288/288 chatcore tests, typecheck/cycles/file-size green. Integrated into release/v3.8.35. --- open-sse/handlers/chatCore.ts | 138 +++----------- open-sse/handlers/chatCore/attemptLogging.ts | 178 +++++++++++++++++++ tests/unit/chatcore-attempt-logging.test.ts | 99 +++++++++++ 3 files changed, 301 insertions(+), 114 deletions(-) create mode 100644 open-sse/handlers/chatCore/attemptLogging.ts create mode 100644 tests/unit/chatcore-attempt-logging.test.ts diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index dbd6727c64..01c454fc6f 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -6,7 +6,6 @@ export { extractSystemRoleMessages } from "./chatCore/claudeSystemRole.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, isNoMemoryRequested, @@ -134,12 +133,16 @@ import { getUpstreamErrorIdentifier, } from "./chatCore/streamErrorResult.ts"; import { wrapReadableStreamWithFinalize } from "./chatCore/streamFinalize.ts"; -import { buildCacheUsageLogMeta, attachLogMeta } from "./chatCore/cacheUsageMeta.ts"; +import { buildCacheUsageLogMeta } from "./chatCore/cacheUsageMeta.ts"; import { buildExecutorClientHeaders } from "./chatCore/executorClientHeaders.ts"; import { resolveExecutionCredentials as resolveExecutionCredentialsFor } from "./chatCore/executionCredentials.ts"; import { resolveExecutorWithProxy as resolveExecutorWithProxyFor } from "./chatCore/executorProxy.ts"; import type { ClaudeMessage } from "./chatCore/claudeMessageTypes.ts"; import { normalizeClaudeUpstreamMessages as normalizeClaudeUpstreamMessagesFor } from "./chatCore/claudeUpstreamMessages.ts"; +import { + persistAttemptLogs as persistAttemptLogsFor, + type PersistAttemptLogsArgs, +} from "./chatCore/attemptLogging.ts"; import { getCallLogPipelineCaptureStreamChunks, @@ -147,16 +150,10 @@ import { } from "@/lib/logEnv"; import { logAuditEvent } from "@/lib/compliance"; import { emit } from "@/lib/events/eventBus"; -import { extractProviderWarnings } from "@/lib/compliance/providerAudit"; import { adaptBodyForCompression } from "../services/compression/bodyAdapter.ts"; import { ensureEngineBreakdown } from "../services/compression/engineBreakdown.ts"; import { handleBypassRequest } from "../utils/bypassHandler.ts"; -import { - saveRequestUsage, - trackPendingRequest, - appendRequestLog, - saveCallLog, -} from "@/lib/usageDb"; +import { saveRequestUsage, trackPendingRequest, appendRequestLog } from "@/lib/usageDb"; import { finalizePendingScope, updatePendingScope } from "@/lib/usage/pendingRequestScope"; import { formatUsageLog } from "@/lib/usage/tokenAccounting"; import { recordCost } from "@/domain/costRules"; @@ -709,118 +706,31 @@ export async function handleChatCore({ clientRawRequest?.headers ?? null, "x-omniroute-session-id" )) || skillRequestId; - const persistAttemptLogs = ({ - status, - tokens, - responseBody, - error, - providerRequest, - providerResponse, - clientResponse, - claudeCacheMeta, - claudeCacheUsageMeta, - cacheSource, - }: { - status: number; - tokens?: unknown; - responseBody?: unknown; - error?: string | null; - providerRequest?: unknown; - providerResponse?: unknown; - clientResponse?: unknown; - claudeCacheMeta?: Record; - claudeCacheUsageMeta?: Record; - cacheSource?: "upstream" | "semantic"; - }) => { - const providerWarnings = extractProviderWarnings( - providerResponse, - clientResponse, - responseBody - ); - if (providerWarnings.length > 0) { - logAuditEvent({ - action: "provider.warning", - actor: "system", - target: [provider, connectionId].filter(Boolean).join(":") || provider || model, - resourceType: "provider_warning", - status: "warning", - requestId: skillRequestId, - details: { - provider, - model, - connectionId, - httpStatus: status, - warnings: providerWarnings, - }, - }); - } - - const pipelinePayloads = detailedLoggingEnabled - ? (reqLogger?.getPipelinePayloads?.() ?? {}) - : null; - - if (pipelinePayloads) { - if (providerRequest !== undefined && !pipelinePayloads.providerRequest) { - pipelinePayloads.providerRequest = providerRequest as Record; - } - if (providerResponse !== undefined && !pipelinePayloads.providerResponse) { - pipelinePayloads.providerResponse = providerResponse as Record; - } - if (clientResponse !== undefined) { - pipelinePayloads.clientResponse = clientResponse as Record; - } - if (error) { - pipelinePayloads.error = { - ...(typeof pipelinePayloads.error === "object" && pipelinePayloads.error - ? (pipelinePayloads.error as Record) - : {}), - message: error, - }; - } - } - - saveCallLog({ - id: pendingRequestId, - method: "POST", - path: clientRawRequest?.endpoint || "/v1/chat/completions", - status, - model, - requestedModel, + // persistAttemptLogs extracted to chatCore/attemptLogging.ts (#3501); bind the per-request context + // once so the 16 call sites keep passing only the per-attempt args (byte-identical). + const persistAttemptLogs = (args: PersistAttemptLogsArgs) => + persistAttemptLogsFor(args, { provider, - connectionId: connectionId || credentials?.connectionId || undefined, - duration: Date.now() - startTime, - tokens: tokens || {}, - requestBody: cloneBoundedChatLogPayload( - attachLogMeta(truncateForLog(body as Record), { - claudePromptCache: claudeCacheMeta, - }) - ), - responseBody: cloneBoundedChatLogPayload( - attachLogMeta(truncateForLog(responseBody as Record), { - claudePromptCache: claudeCacheMeta - ? { - applied: claudeCacheMeta.applied, - totalBreakpoints: claudeCacheMeta.totalBreakpoints, - anthropicBeta: claudeCacheMeta.anthropicBeta, - } - : null, - claudePromptCacheUsage: claudeCacheUsageMeta, - }) - ), - error: error || null, + connectionId, + model, + skillRequestId, + detailedLoggingEnabled, + reqLogger, + pendingRequestId, + clientRawRequest, + requestedModel, + credentials, + startTime, + body, sourceFormat, targetFormat, comboName, comboStepId, comboExecutionKey, tokensCompressed, - cacheSource: cacheSource === "semantic" ? "semantic" : "upstream", - apiKeyId: apiKeyInfo?.id || null, - apiKeyName: apiKeyInfo?.name || null, - noLog: noLogEnabled, - pipelinePayloads, - }).catch(() => {}); - }; + apiKeyInfo, + noLogEnabled, + }); // Primary path: merge client model id + alias target so config on either key applies; resolved // id wins on same header name. T5 family fallback uses only (nextModel, resolveModelAlias(next)) diff --git a/open-sse/handlers/chatCore/attemptLogging.ts b/open-sse/handlers/chatCore/attemptLogging.ts new file mode 100644 index 0000000000..7414b951da --- /dev/null +++ b/open-sse/handlers/chatCore/attemptLogging.ts @@ -0,0 +1,178 @@ +/** + * chatCore per-attempt logging persistence (Quality Gate v2 / Fase 9 — chatCore god-file + * decomposition, #3501). + * + * Extracted from handleChatCore: persists one attempt's call log. Emits a provider.warning audit + * event when the provider response carries warnings, fills the detailed pipeline payloads (when + * detailed logging is on), and writes the bounded/truncated call-log row (request/response bodies + * with the Claude prompt-cache meta attached). Best-effort: the saveCallLog write swallows its own + * errors. The per-request context (provider/model/ids/combo/etc.) is threaded via `ctx` so the 16 + * call sites in the handler stay byte-identical; behaviour is unchanged. + */ + +import { extractProviderWarnings } from "@/lib/compliance/providerAudit"; +import { logAuditEvent } from "@/lib/compliance"; +import { saveCallLog } from "@/lib/usageDb"; +import { cloneBoundedChatLogPayload, truncateForLog } from "./logTruncation.ts"; +import { attachLogMeta } from "./cacheUsageMeta.ts"; + +export type PersistAttemptLogsArgs = { + status: number; + tokens?: unknown; + responseBody?: unknown; + error?: string | null; + providerRequest?: unknown; + providerResponse?: unknown; + clientResponse?: unknown; + claudeCacheMeta?: Record; + claudeCacheUsageMeta?: Record; + cacheSource?: "upstream" | "semantic"; +}; + +export type PersistAttemptLogsContext = { + provider: string | null | undefined; + connectionId: string | null | undefined; + model: string | null | undefined; + skillRequestId: string; + detailedLoggingEnabled: boolean; + reqLogger: + | { getPipelinePayloads?: () => Record | undefined } + | null + | undefined; + pendingRequestId: unknown; + clientRawRequest: { endpoint?: string } | null | undefined; + requestedModel: unknown; + credentials: { connectionId?: string } | null | undefined; + startTime: number; + body: unknown; + sourceFormat: unknown; + targetFormat: unknown; + comboName: unknown; + comboStepId: unknown; + comboExecutionKey: unknown; + tokensCompressed: unknown; + apiKeyInfo: { id?: string | null; name?: string | null } | null | undefined; + noLogEnabled: unknown; +}; + +export function persistAttemptLogs(args: PersistAttemptLogsArgs, ctx: PersistAttemptLogsContext) { + const { + status, + tokens, + responseBody, + error, + providerRequest, + providerResponse, + clientResponse, + claudeCacheMeta, + claudeCacheUsageMeta, + cacheSource, + } = args; + const { + provider, + connectionId, + model, + skillRequestId, + detailedLoggingEnabled, + reqLogger, + pendingRequestId, + clientRawRequest, + requestedModel, + credentials, + startTime, + body, + sourceFormat, + targetFormat, + comboName, + comboStepId, + comboExecutionKey, + tokensCompressed, + apiKeyInfo, + noLogEnabled, + } = ctx; + + const providerWarnings = extractProviderWarnings(providerResponse, clientResponse, responseBody); + if (providerWarnings.length > 0) { + logAuditEvent({ + action: "provider.warning", + actor: "system", + target: [provider, connectionId].filter(Boolean).join(":") || provider || model, + resourceType: "provider_warning", + status: "warning", + requestId: skillRequestId, + details: { + provider, + model, + connectionId, + httpStatus: status, + warnings: providerWarnings, + }, + }); + } + + const pipelinePayloads = detailedLoggingEnabled + ? (reqLogger?.getPipelinePayloads?.() ?? {}) + : null; + + if (pipelinePayloads) { + if (providerRequest !== undefined && !pipelinePayloads.providerRequest) { + pipelinePayloads.providerRequest = providerRequest as Record; + } + if (providerResponse !== undefined && !pipelinePayloads.providerResponse) { + pipelinePayloads.providerResponse = providerResponse as Record; + } + if (clientResponse !== undefined) { + pipelinePayloads.clientResponse = clientResponse as Record; + } + if (error) { + pipelinePayloads.error = { + ...(typeof pipelinePayloads.error === "object" && pipelinePayloads.error + ? (pipelinePayloads.error as Record) + : {}), + message: error, + }; + } + } + + saveCallLog({ + id: pendingRequestId, + method: "POST", + path: clientRawRequest?.endpoint || "/v1/chat/completions", + status, + model, + requestedModel, + provider, + connectionId: connectionId || credentials?.connectionId || undefined, + duration: Date.now() - startTime, + tokens: tokens || {}, + requestBody: cloneBoundedChatLogPayload( + attachLogMeta(truncateForLog(body as Record), { + claudePromptCache: claudeCacheMeta, + }) + ), + responseBody: cloneBoundedChatLogPayload( + attachLogMeta(truncateForLog(responseBody as Record), { + claudePromptCache: claudeCacheMeta + ? { + applied: claudeCacheMeta.applied, + totalBreakpoints: claudeCacheMeta.totalBreakpoints, + anthropicBeta: claudeCacheMeta.anthropicBeta, + } + : null, + claudePromptCacheUsage: claudeCacheUsageMeta, + }) + ), + error: error || null, + sourceFormat, + targetFormat, + comboName, + comboStepId, + comboExecutionKey, + tokensCompressed, + cacheSource: cacheSource === "semantic" ? "semantic" : "upstream", + apiKeyId: apiKeyInfo?.id || null, + apiKeyName: apiKeyInfo?.name || null, + noLog: noLogEnabled, + pipelinePayloads, + }).catch(() => {}); +} diff --git a/tests/unit/chatcore-attempt-logging.test.ts b/tests/unit/chatcore-attempt-logging.test.ts new file mode 100644 index 0000000000..d4871162f3 --- /dev/null +++ b/tests/unit/chatcore-attempt-logging.test.ts @@ -0,0 +1,99 @@ +// tests/unit/chatcore-attempt-logging.test.ts +// Characterization of persistAttemptLogs — the per-attempt call-log persistence extracted from +// handleChatCore (chatCore god-file decomposition, #3501). Uses a real temp DB and polls the +// persisted row (saveCallLog is async + fire-and-forget). Locks: the field mapping, the +// cacheSource semantic/upstream normalization, the connectionId → credentials.connectionId +// fallback, and error persistence. +import { test, before, after } from "node:test"; +import assert from "node:assert/strict"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +const testDataDir = fs.mkdtempSync(path.join(os.tmpdir(), "omni-attempt-logging-test-")); +process.env.DATA_DIR = testDataDir; + +const coreDb = await import("../../src/lib/db/core.ts"); +const { getCallLogById } = await import("../../src/lib/usage/callLogs.ts"); +const { persistAttemptLogs } = await import("../../open-sse/handlers/chatCore/attemptLogging.ts"); + +function baseCtx(overrides: Record = {}) { + return { + provider: "openai", + connectionId: "conn-1", + model: "gpt-x", + skillRequestId: "skill-1", + detailedLoggingEnabled: false, + reqLogger: null, + pendingRequestId: "REPLACE", + clientRawRequest: { endpoint: "/v1/chat/completions" }, + requestedModel: "gpt-x-requested", + credentials: { connectionId: "cred-conn" }, + startTime: Date.now(), + body: { messages: [{ role: "user", content: "hi" }] }, + sourceFormat: "openai", + targetFormat: "openai", + comboName: null, + comboStepId: null, + comboExecutionKey: null, + tokensCompressed: 0, + apiKeyInfo: { id: "key-1", name: "Key One" }, + noLogEnabled: false, + ...overrides, + } as Parameters[1]; +} + +async function pollForCallLog(id: string, tries = 120) { + for (let i = 0; i < tries; i++) { + const row = await getCallLogById(id); + if (row) return row as Record; + await new Promise((r) => setTimeout(r, 20)); + } + return null; +} + +before(async () => { + await coreDb.ensureDbInitialized(); +}); + +after(() => { + coreDb.resetDbInstance(); + fs.rmSync(testDataDir, { recursive: true, force: true }); +}); + +test("persists a call log row with the mapped fields (default cacheSource=upstream)", async () => { + const id = "attempt-basic-1"; + persistAttemptLogs({ status: 200, tokens: { input: 1, output: 2 } }, baseCtx({ pendingRequestId: id })); + const row = await pollForCallLog(id); + assert.ok(row, "call log row should be persisted"); + assert.equal(row.status, 200); + assert.equal(row.model, "gpt-x"); + assert.equal(row.provider, "openai"); + assert.equal(row.requestedModel, "gpt-x-requested"); + assert.equal(row.connectionId, "conn-1"); + assert.equal(row.cacheSource, "upstream"); +}); + +test("cacheSource 'semantic' is preserved", async () => { + const id = "attempt-semantic-1"; + persistAttemptLogs( + { status: 200, cacheSource: "semantic" }, + baseCtx({ pendingRequestId: id }) + ); + const row = await pollForCallLog(id); + assert.ok(row); + assert.equal(row.cacheSource, "semantic"); +}); + +test("connectionId falls back to credentials.connectionId when null, and error is persisted", async () => { + const id = "attempt-fallback-1"; + persistAttemptLogs( + { status: 502, error: "upstream boom" }, + baseCtx({ pendingRequestId: id, connectionId: null }) + ); + const row = await pollForCallLog(id); + assert.ok(row); + assert.equal(row.connectionId, "cred-conn"); + assert.equal(row.status, 502); + assert.match(String(row.error ?? ""), /upstream boom/); +});