diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 539d38b238..dc9146bd3a 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -2,7 +2,7 @@ 🌐 **Languages:** 🇺🇸 [English](ARCHITECTURE.md) | 🇧🇷 [Português (Brasil)](i18n/pt-BR/ARCHITECTURE.md) | 🇪🇸 [Español](i18n/es/ARCHITECTURE.md) | 🇫🇷 [Français](i18n/fr/ARCHITECTURE.md) | 🇮🇹 [Italiano](i18n/it/ARCHITECTURE.md) | 🇷🇺 [Русский](i18n/ru/ARCHITECTURE.md) | 🇨🇳 [中文 (简体)](i18n/zh-CN/ARCHITECTURE.md) | 🇩🇪 [Deutsch](i18n/de/ARCHITECTURE.md) | 🇮🇳 [हिन्दी](i18n/in/ARCHITECTURE.md) | 🇹🇭 [ไทย](i18n/th/ARCHITECTURE.md) | 🇺🇦 [Українська](i18n/uk-UA/ARCHITECTURE.md) | 🇸🇦 [العربية](i18n/ar/ARCHITECTURE.md) | 🇯🇵 [日本語](i18n/ja/ARCHITECTURE.md) | 🇻🇳 [Tiếng Việt](i18n/vi/ARCHITECTURE.md) | 🇧🇬 [Български](i18n/bg/ARCHITECTURE.md) | 🇩🇰 [Dansk](i18n/da/ARCHITECTURE.md) | 🇫🇮 [Suomi](i18n/fi/ARCHITECTURE.md) | 🇮🇱 [עברית](i18n/he/ARCHITECTURE.md) | 🇭🇺 [Magyar](i18n/hu/ARCHITECTURE.md) | 🇮🇩 [Bahasa Indonesia](i18n/id/ARCHITECTURE.md) | 🇰🇷 [한국어](i18n/ko/ARCHITECTURE.md) | 🇲🇾 [Bahasa Melayu](i18n/ms/ARCHITECTURE.md) | 🇳🇱 [Nederlands](i18n/nl/ARCHITECTURE.md) | 🇳🇴 [Norsk](i18n/no/ARCHITECTURE.md) | 🇵🇹 [Português (Portugal)](i18n/pt/ARCHITECTURE.md) | 🇷🇴 [Română](i18n/ro/ARCHITECTURE.md) | 🇵🇱 [Polski](i18n/pl/ARCHITECTURE.md) | 🇸🇰 [Slovenčina](i18n/sk/ARCHITECTURE.md) | 🇸🇪 [Svenska](i18n/sv/ARCHITECTURE.md) | 🇵🇭 [Filipino](i18n/phi/ARCHITECTURE.md) | 🇨🇿 [Čeština](i18n/cs/ARCHITECTURE.md) -_Last updated: 2026-03-24_ +_Last updated: 2026-03-28_ ## Executive Summary @@ -756,10 +756,18 @@ Runtime visibility sources: - console logs from `src/sse/utils/logger.ts` - per-request usage aggregates in SQLite (`usage_history`, `call_logs`, `proxy_logs`) +- four-stage detailed payload captures in SQLite (`request_detail_logs`) when `settings.detailed_logs_enabled=true` - textual request status log in `log.txt` (optional/compat) - optional deep request/translation logs under `logs/` when `ENABLE_REQUEST_LOGS=true` - dashboard usage endpoints (`/api/usage/*`) for UI consumption +Detailed request payload capture stores up to four JSON payload stages per routed call: + +- raw request received from the client +- translated request actually sent upstream +- provider response reconstructed as JSON (including streamed event sequences when applicable) +- final client response returned by OmniRoute + ## Security-Sensitive Boundaries - JWT secret (`JWT_SECRET`) secures dashboard session cookie verification/signing diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index 3e09795434..2a3ee9c03e 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -14,10 +14,16 @@ import { createRequestLogger } from "../utils/requestLogger.ts"; import { getModelTargetFormat, PROVIDER_ID_TO_ALIAS } from "../config/providerModels.ts"; import { resolveModelAlias } from "../services/modelDeprecation.ts"; import { getUnsupportedParams } from "../config/providerRegistry.ts"; -import { createErrorResult, parseUpstreamError, formatProviderError } from "../utils/error.ts"; +import { + buildErrorBody, + createErrorResult, + parseUpstreamError, + formatProviderError, +} from "../utils/error.ts"; import { HTTP_STATUS, PROVIDER_MAX_TOKENS } from "../config/constants.ts"; import { classifyProviderError, PROVIDER_ERROR_TYPES } from "../services/errorClassifier.ts"; import { updateProviderConnection } from "@/lib/db/providers"; +import { isDetailedLoggingEnabled, saveRequestDetailLog } from "@/lib/db/detailedLogs"; import { logAuditEvent } from "@/lib/compliance"; import { handleBypassRequest } from "../utils/bypassHandler.ts"; import { @@ -72,6 +78,8 @@ import { EMERGENCY_FALLBACK_CONFIG, } from "../services/emergencyFallback.ts"; import { resolveStreamFlag, stripMarkdownCodeFence } from "../utils/aiSdkCompat.ts"; +import { generateRequestId } from "@/shared/utils/requestId"; +import { normalizePayloadForLog } from "@/lib/logPayloads"; export function shouldUseNativeCodexPassthrough({ provider, @@ -391,7 +399,8 @@ export async function handleChatCore({ credentials.providerSpecificData = nextProviderData; } catch (err) { - log?.debug?.("CODEX", `Failed to persist codex quota state: ${err?.message || err}`); + const errMessage = err instanceof Error ? err.message : String(err); + log?.debug?.("CODEX", `Failed to persist codex quota state: ${errMessage}`); } }; @@ -486,6 +495,88 @@ export async function handleChatCore({ const alias = PROVIDER_ID_TO_ALIAS[provider] || provider; const modelTargetFormat = getModelTargetFormat(alias, resolvedModel); const targetFormat = modelTargetFormat || getTargetFormat(provider); + const noLogEnabled = apiKeyInfo?.noLog === true; + const detailedLoggingEnabled = !noLogEnabled && (await isDetailedLoggingEnabled()); + const persistAttemptLogs = ({ + status, + tokens, + responseBody, + error, + providerRequest, + providerResponse, + clientResponse, + claudeCacheMeta, + claudeCacheUsageMeta, + }: { + status: number; + tokens?: unknown; + responseBody?: unknown; + error?: string | null; + providerRequest?: unknown; + providerResponse?: unknown; + clientResponse?: unknown; + claudeCacheMeta?: any; + claudeCacheUsageMeta?: any; + }) => { + const callLogId = generateRequestId(); + + saveCallLog({ + id: callLogId, + method: "POST", + path: clientRawRequest?.endpoint || "/v1/chat/completions", + status, + model, + requestedModel, + provider, + connectionId, + duration: Date.now() - startTime, + tokens: tokens || {}, + requestBody: attachLogMeta(body, { + claudePromptCache: claudeCacheMeta, + }), + responseBody: attachLogMeta(responseBody ?? undefined, { + claudePromptCache: claudeCacheMeta + ? { + applied: claudeCacheMeta.applied, + totalBreakpoints: claudeCacheMeta.totalBreakpoints, + anthropicBeta: claudeCacheMeta.anthropicBeta, + } + : null, + claudePromptCacheUsage: claudeCacheUsageMeta, + }), + error: error || null, + sourceFormat, + targetFormat, + comboName, + apiKeyId: apiKeyInfo?.id || null, + apiKeyName: apiKeyInfo?.name || null, + noLog: noLogEnabled, + }).catch(() => {}); + + if (!detailedLoggingEnabled) { + return; + } + + try { + saveRequestDetailLog({ + call_log_id: callLogId, + client_request: clientRawRequest?.body ?? body, + translated_request: providerRequest ?? null, + provider_response: providerResponse ?? null, + client_response: clientResponse ?? null, + provider, + model, + source_format: sourceFormat, + target_format: targetFormat, + duration_ms: Date.now() - startTime, + api_key_id: apiKeyInfo?.id || null, + no_log: noLogEnabled, + }); + } catch (err) { + const errMessage = err instanceof Error ? err.message : String(err); + log?.debug?.("DETAIL_LOG", `Failed to save detailed log: ${errMessage}`); + } + }; // 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)) @@ -919,40 +1010,34 @@ export async function handleChatCore({ ); } catch (error) { trackPendingRequest(model, provider, connectionId, false); + const failureStatus = error.name === "AbortError" ? 499 : HTTP_STATUS.BAD_GATEWAY; + const failureMessage = + error.name === "AbortError" + ? "Request aborted" + : formatProviderError(error, provider, model, HTTP_STATUS.BAD_GATEWAY); appendRequestLog({ model, provider, connectionId, - status: `FAILED ${error.name === "AbortError" ? 499 : HTTP_STATUS.BAD_GATEWAY}`, - }).catch(() => {}); - saveCallLog({ - method: "POST", - path: clientRawRequest?.endpoint || "/v1/chat/completions", - status: error.name === "AbortError" ? 499 : HTTP_STATUS.BAD_GATEWAY, - model, - requestedModel, - provider, - connectionId, - duration: Date.now() - startTime, - requestBody: attachLogMeta(body, { - claudePromptCache: claudePromptCacheLogMeta, - }), - error: error.message, - sourceFormat, - targetFormat, - comboName, - apiKeyId: apiKeyInfo?.id || null, - apiKeyName: apiKeyInfo?.name || null, - noLog: apiKeyInfo?.noLog === true, + status: `FAILED ${failureStatus}`, }).catch(() => {}); + persistAttemptLogs({ + status: failureStatus, + error: failureMessage, + providerRequest: finalBody || translatedBody, + clientResponse: buildErrorBody(failureStatus, failureMessage), + claudeCacheMeta: claudePromptCacheLogMeta, + }); if (error.name === "AbortError") { streamController.handleError(error); return createErrorResult(499, "Request aborted"); } - persistFailureUsage(HTTP_STATUS.BAD_GATEWAY, error?.name || "upstream_error"); - const errMsg = formatProviderError(error, provider, model, HTTP_STATUS.BAD_GATEWAY); - console.log(`${COLORS.red}[ERROR] ${errMsg}${COLORS.reset}`); - return createErrorResult(HTTP_STATUS.BAD_GATEWAY, errMsg); + persistFailureUsage( + HTTP_STATUS.BAD_GATEWAY, + error instanceof Error && error.name ? error.name : "upstream_error" + ); + console.log(`${COLORS.red}[ERROR] ${failureMessage}${COLORS.reset}`); + return createErrorResult(HTTP_STATUS.BAD_GATEWAY, failureMessage); } // Handle 401/403 - try token refresh using executor @@ -998,8 +1083,11 @@ export async function handleChatCore({ if (retryResult.response.ok) { providerResponse = retryResult.response; providerUrl = retryResult.url; + providerHeaders = retryResult.headers; + finalBody = retryResult.transformedBody; + reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); } - } catch (retryError) { + } catch { log?.warn?.("TOKEN", `${provider.toUpperCase()} | retry after refresh failed`); } } else { @@ -1012,10 +1100,12 @@ export async function handleChatCore({ // Check provider response - return error info for fallback handling if (!providerResponse.ok) { trackPendingRequest(model, provider, connectionId, false); - const { statusCode, message, retryAfterMs } = await parseUpstreamError( - providerResponse, - provider - ); + const { + statusCode, + message, + retryAfterMs, + responseBody: upstreamErrorBody, + } = await parseUpstreamError(providerResponse, provider); // T06/T10/T36: classify provider errors and persist terminal account states. const errorType = classifyProviderError(statusCode, message); @@ -1067,26 +1157,7 @@ export async function handleChatCore({ appendRequestLog({ model, provider, connectionId, status: `FAILED ${statusCode}` }).catch( () => {} ); - saveCallLog({ - method: "POST", - path: clientRawRequest?.endpoint || "/v1/chat/completions", - status: statusCode, - model, - requestedModel, - provider, - connectionId, - duration: Date.now() - startTime, - requestBody: attachLogMeta(body, { - claudePromptCache: claudePromptCacheLogMeta, - }), - error: message, - sourceFormat, - targetFormat, - comboName, - apiKeyId: apiKeyInfo?.id || null, - apiKeyName: apiKeyInfo?.name || null, - noLog: apiKeyInfo?.noLog === true, - }).catch(() => {}); + const errMsg = formatProviderError(new Error(message), provider, model, statusCode); console.log(`${COLORS.red}[ERROR] ${errMsg}${COLORS.reset}`); @@ -1098,6 +1169,12 @@ export async function handleChatCore({ // Log error with full request body for debugging reqLogger.logError(new Error(message), finalBody || translatedBody); + reqLogger.logProviderResponse( + providerResponse.status, + providerResponse.statusText, + providerResponse.headers, + upstreamErrorBody + ); // Update rate limiter from error response headers updateFromHeaders(provider, connectionId, providerResponse.headers, statusCode, model); @@ -1121,24 +1198,53 @@ export async function handleChatCore({ providerUrl = fallbackResult.url; providerHeaders = fallbackResult.headers; finalBody = fallbackResult.transformedBody; + reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); // Continue processing with the fallback response — skip error return log?.info?.("MODEL_FALLBACK", `Serving ${nextModel} as fallback for ${model}`); // Jump to streaming/non-streaming handling below // We fall through by NOT returning here } else { // Fallback also failed — return original error + persistAttemptLogs({ + status: statusCode, + error: errMsg, + providerRequest: finalBody || translatedBody, + providerResponse: upstreamErrorBody, + clientResponse: buildErrorBody(statusCode, errMsg), + }); persistFailureUsage(statusCode, "model_unavailable"); return createErrorResult(statusCode, errMsg, retryAfterMs); } } catch { + persistAttemptLogs({ + status: statusCode, + error: errMsg, + providerRequest: finalBody || translatedBody, + providerResponse: upstreamErrorBody, + clientResponse: buildErrorBody(statusCode, errMsg), + }); persistFailureUsage(statusCode, "model_unavailable"); return createErrorResult(statusCode, errMsg, retryAfterMs); } } else { + persistAttemptLogs({ + status: statusCode, + error: errMsg, + providerRequest: finalBody || translatedBody, + providerResponse: upstreamErrorBody, + clientResponse: buildErrorBody(statusCode, errMsg), + }); persistFailureUsage(statusCode, "model_unavailable"); return createErrorResult(statusCode, errMsg, retryAfterMs); } } else { + persistAttemptLogs({ + status: statusCode, + error: errMsg, + providerRequest: finalBody || translatedBody, + providerResponse: upstreamErrorBody, + clientResponse: buildErrorBody(statusCode, errMsg), + }); persistFailureUsage(statusCode, `upstream_${statusCode}`); return createErrorResult(statusCode, errMsg, retryAfterMs); } @@ -1183,6 +1289,10 @@ export async function handleChatCore({ }); if (fbResult.response.ok) { providerResponse = fbResult.response; + providerUrl = fbResult.url; + providerHeaders = fbResult.headers; + finalBody = fbResult.transformedBody; + reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); log?.info?.( "EMERGENCY_FALLBACK", `Serving ${fbDecision.provider}/${fbDecision.model} as budget fallback for ${provider}/${model}` @@ -1195,7 +1305,8 @@ export async function handleChatCore({ ); } } catch (fbErr) { - log?.warn?.("EMERGENCY_FALLBACK", `Emergency fallback error: ${fbErr?.message}`); + const errMessage = fbErr instanceof Error ? fbErr.message : String(fbErr); + log?.warn?.("EMERGENCY_FALLBACK", `Emergency fallback error: ${errMessage}`); } } } @@ -1208,6 +1319,7 @@ export async function handleChatCore({ const contentType = (providerResponse.headers.get("content-type") || "").toLowerCase(); let responseBody; const rawBody = await providerResponse.text(); + const normalizedProviderPayload = normalizePayloadForLog(rawBody); const looksLikeSSE = contentType.includes("text/event-stream") || /(^|\n)\s*(event|data):/m.test(rawBody); @@ -1225,11 +1337,16 @@ export async function handleChatCore({ connectionId, status: `FAILED ${HTTP_STATUS.BAD_GATEWAY}`, }).catch(() => {}); + const invalidSseMessage = "Invalid SSE response for non-streaming request"; + persistAttemptLogs({ + status: HTTP_STATUS.BAD_GATEWAY, + error: invalidSseMessage, + providerRequest: finalBody || translatedBody, + providerResponse: normalizedProviderPayload, + clientResponse: buildErrorBody(HTTP_STATUS.BAD_GATEWAY, invalidSseMessage), + }); persistFailureUsage(HTTP_STATUS.BAD_GATEWAY, "invalid_sse_payload"); - return createErrorResult( - HTTP_STATUS.BAD_GATEWAY, - "Invalid SSE response for non-streaming request" - ); + return createErrorResult(HTTP_STATUS.BAD_GATEWAY, invalidSseMessage); } responseBody = parsedFromSSE; @@ -1243,14 +1360,34 @@ export async function handleChatCore({ connectionId, status: `FAILED ${HTTP_STATUS.BAD_GATEWAY}`, }).catch(() => {}); + const invalidJsonMessage = "Invalid JSON response from provider"; + persistAttemptLogs({ + status: HTTP_STATUS.BAD_GATEWAY, + error: invalidJsonMessage, + providerRequest: finalBody || translatedBody, + providerResponse: normalizedProviderPayload, + clientResponse: buildErrorBody(HTTP_STATUS.BAD_GATEWAY, invalidJsonMessage), + }); persistFailureUsage(HTTP_STATUS.BAD_GATEWAY, "invalid_json_payload"); - return createErrorResult(HTTP_STATUS.BAD_GATEWAY, "Invalid JSON response from provider"); + return createErrorResult(HTTP_STATUS.BAD_GATEWAY, invalidJsonMessage); } } if (sourceFormat === FORMATS.CLAUDE && targetFormat === FORMATS.CLAUDE) { responseBody = restoreClaudePassthroughToolNames(responseBody, toolNameMap); } + reqLogger.logProviderResponse( + providerResponse.status, + providerResponse.statusText, + providerResponse.headers, + looksLikeSSE + ? { + _streamed: true, + _format: "sse-json", + summary: responseBody, + } + : responseBody + ); // Notify success - caller can clear error status if needed if (onRequestSuccess) { @@ -1265,36 +1402,6 @@ export async function handleChatCore({ // Save structured call log with full payloads const cacheUsageLogMeta = buildCacheUsageLogMeta(usage); - saveCallLog({ - method: "POST", - path: clientRawRequest?.endpoint || "/v1/chat/completions", - status: 200, - model, - requestedModel, - provider, - connectionId, - duration: Date.now() - startTime, - tokens: usage, - requestBody: attachLogMeta(body, { - claudePromptCache: claudePromptCacheLogMeta, - }), - responseBody: attachLogMeta(responseBody, { - claudePromptCache: claudePromptCacheLogMeta - ? { - applied: claudePromptCacheLogMeta.applied, - totalBreakpoints: claudePromptCacheLogMeta.totalBreakpoints, - anthropicBeta: claudePromptCacheLogMeta.anthropicBeta, - } - : null, - claudePromptCacheUsage: cacheUsageLogMeta, - }), - sourceFormat, - targetFormat, - comboName, - apiKeyId: apiKeyInfo?.id || null, - apiKeyName: apiKeyInfo?.name || null, - noLog: apiKeyInfo?.noLog === true, - }).catch(() => {}); if (usage && typeof usage === "object") { const msg = `[${new Date().toLocaleTimeString("en-US", { hour12: false, hour: "2-digit", minute: "2-digit" })}] 📊 [USAGE] ${provider.toUpperCase()} | in=${getLoggedInputTokens(usage)} | out=${getLoggedOutputTokens(usage)}${connectionId ? ` | account=${connectionId.slice(0, 8)}...` : ""}`; console.log(`${COLORS.green}${msg}${COLORS.reset}`); @@ -1387,6 +1494,23 @@ export async function handleChatCore({ // ── Phase 9.2: Save for idempotency ── saveIdempotency(idempotencyKey, translatedResponse, 200); + reqLogger.logConvertedResponse(translatedResponse); + persistAttemptLogs({ + status: 200, + tokens: usage, + responseBody, + providerRequest: finalBody || translatedBody, + providerResponse: looksLikeSSE + ? { + _streamed: true, + _format: "sse-json", + summary: responseBody, + } + : responseBody, + clientResponse: translatedResponse, + claudeCacheMeta: claudePromptCacheLogMeta, + claudeCacheUsageMeta: cacheUsageLogMeta, + }); return { success: true, @@ -1422,38 +1546,20 @@ export async function handleChatCore({ status: streamStatus, usage: streamUsage, responseBody: streamResponseBody, + providerPayload, + clientPayload, }) => { const cacheUsageLogMeta = buildCacheUsageLogMeta(streamUsage); - saveCallLog({ - method: "POST", - path: clientRawRequest?.endpoint || "/v1/chat/completions", + persistAttemptLogs({ status: streamStatus || 200, - model, - requestedModel, - provider, - connectionId, - duration: Date.now() - startTime, tokens: streamUsage || {}, - requestBody: attachLogMeta(body, { - claudePromptCache: claudePromptCacheLogMeta, - }), - responseBody: attachLogMeta(streamResponseBody ?? undefined, { - claudePromptCache: claudePromptCacheLogMeta - ? { - applied: claudePromptCacheLogMeta.applied, - totalBreakpoints: claudePromptCacheLogMeta.totalBreakpoints, - anthropicBeta: claudePromptCacheLogMeta.anthropicBeta, - } - : null, - claudePromptCacheUsage: cacheUsageLogMeta, - }), - sourceFormat, - targetFormat, - comboName, - apiKeyId: apiKeyInfo?.id || null, - apiKeyName: apiKeyInfo?.name || null, - noLog: apiKeyInfo?.noLog === true, - }).catch(() => {}); + responseBody: streamResponseBody ?? undefined, + providerRequest: finalBody || translatedBody, + providerResponse: providerPayload, + clientResponse: clientPayload ?? streamResponseBody ?? undefined, + claudeCacheMeta: claudePromptCacheLogMeta, + claudeCacheUsageMeta: cacheUsageLogMeta, + }); if (apiKeyInfo?.id && streamUsage) { calculateCost(provider, model, streamUsage) diff --git a/open-sse/utils/error.ts b/open-sse/utils/error.ts index 726f757eee..367d5b2730 100644 --- a/open-sse/utils/error.ts +++ b/open-sse/utils/error.ts @@ -1,5 +1,6 @@ import { getCorsOrigin } from "./cors.ts"; import { ERROR_TYPES, DEFAULT_ERROR_MESSAGES } from "../config/constants.ts"; +import { normalizePayloadForLog } from "@/lib/logPayloads"; /** * Build OpenAI-compatible error response body @@ -91,14 +92,16 @@ export function parseAntigravityRetryTime(message) { * Parse upstream provider error response * @param {Response} response - Fetch response from provider * @param {string} provider - Provider name (for Antigravity-specific parsing) - * @returns {Promise<{statusCode: number, message: string, retryAfterMs: number|null}>} + * @returns {Promise<{statusCode: number, message: string, retryAfterMs: number|null, responseBody: unknown}>} */ export async function parseUpstreamError(response, provider = null) { let message = ""; let retryAfterMs = null; + let responseBody = null; try { const text = await response.text(); + responseBody = normalizePayloadForLog(text); // Try parse as JSON try { @@ -109,6 +112,7 @@ export async function parseUpstreamError(response, provider = null) { } } catch { message = `Upstream error: ${response.status}`; + responseBody = { _rawText: message }; } const messageStr = typeof message === "string" ? message : JSON.stringify(message); @@ -122,6 +126,7 @@ export async function parseUpstreamError(response, provider = null) { statusCode: response.status, message: messageStr, retryAfterMs, + responseBody, }; } diff --git a/open-sse/utils/stream.ts b/open-sse/utils/stream.ts index 54e03f352a..8a2fe511f5 100644 --- a/open-sse/utils/stream.ts +++ b/open-sse/utils/stream.ts @@ -11,6 +11,7 @@ import { COLORS, } from "./usageTracking.ts"; import { parseSSELine, hasValuableContent, fixInvalidId, formatSSE } from "./streamHelpers.ts"; +import { createStructuredSSECollector } from "./streamPayloadCollector.ts"; import { STREAM_IDLE_TIMEOUT_MS, HTTP_STATUS } from "../config/constants.ts"; import { sanitizeStreamingChunk, @@ -32,6 +33,8 @@ type StreamCompletePayload = { usage: unknown; /** Minimal response body for call log (streaming: usage + note; non-streaming not used) */ responseBody?: unknown; + providerPayload?: unknown; + clientPayload?: unknown; }; type StreamOptions = { @@ -158,6 +161,12 @@ export function createSSEStream(options: StreamOptions = {}) { // Guard against duplicate [DONE] events — ensures exactly one per stream let doneSent = false; + const providerPayloadCollector = createStructuredSSECollector({ + stage: "provider_response", + }); + const clientPayloadCollector = createStructuredSSECollector({ + stage: "client_response", + }); // Per-stream instances to avoid shared state with concurrent streams const decoder = new TextDecoder(); @@ -212,6 +221,17 @@ export function createSSEStream(options: StreamOptions = {}) { if (mode === STREAM_MODE.PASSTHROUGH) { let output; let injectedUsage = false; + let clientPayload: unknown = null; + + if (trimmed.startsWith("data:")) { + const providerPayload = parseSSELine(trimmed); + if (providerPayload) { + providerPayloadCollector.push(providerPayload); + if ((providerPayload as { done?: unknown }).done === true) { + clientPayloadCollector.push(providerPayload); + } + } + } if (trimmed.startsWith("data:") && trimmed.slice(5).trim() !== "[DONE]") { try { @@ -380,6 +400,8 @@ export function createSSEStream(options: StreamOptions = {}) { injectedUsage = true; } } + + clientPayload = parsed; } catch {} } @@ -391,6 +413,10 @@ export function createSSEStream(options: StreamOptions = {}) { } } + if (clientPayload) { + clientPayloadCollector.push(clientPayload); + } + reqLogger?.appendConvertedChunk?.(output); controller.enqueue(encoder.encode(output)); continue; @@ -401,10 +427,12 @@ export function createSSEStream(options: StreamOptions = {}) { const parsed = parseSSELine(trimmed); if (!parsed) continue; + providerPayloadCollector.push(parsed); if (parsed && parsed.done) { if (!doneSent) { doneSent = true; + clientPayloadCollector.push({ done: true }); const output = "data: [DONE]\n\n"; reqLogger?.appendConvertedChunk?.(output); controller.enqueue(encoder.encode(output)); @@ -524,6 +552,7 @@ export function createSSEStream(options: StreamOptions = {}) { } const output = formatSSE(item, sourceFormat); + clientPayloadCollector.push(item); reqLogger?.appendConvertedChunk?.(output); controller.enqueue(encoder.encode(output)); } @@ -551,6 +580,11 @@ export function createSSEStream(options: StreamOptions = {}) { if (buffer.startsWith("data:") && !buffer.startsWith("data: ")) { output = "data: " + buffer.slice(5); } + const bufferedPayload = parseSSELine(buffer.trim()); + if (bufferedPayload) { + providerPayloadCollector.push(bufferedPayload); + clientPayloadCollector.push(bufferedPayload); + } reqLogger?.appendConvertedChunk?.(output); controller.enqueue(encoder.encode(output)); } @@ -601,7 +635,13 @@ export function createSSEStream(options: StreamOptions = {}) { }, _streamed: true, }; - onComplete({ status: 200, usage, responseBody }); + onComplete({ + status: 200, + usage, + responseBody, + providerPayload: providerPayloadCollector.build(), + clientPayload: clientPayloadCollector.build(responseBody), + }); } catch {} } return; @@ -611,6 +651,7 @@ export function createSSEStream(options: StreamOptions = {}) { if (buffer.trim()) { const parsed = parseSSELine(buffer.trim()); if (parsed && !parsed.done) { + providerPayloadCollector.push(parsed); // Extract usage from remaining buffer — if the usage-bearing event // (e.g. response.completed) is the last SSE line, it ends up here // in the flush handler where extractUsage was not called. @@ -647,6 +688,7 @@ export function createSSEStream(options: StreamOptions = {}) { if (translated?.length > 0) { for (const item of translated) { const output = formatSSE(item, sourceFormat); + clientPayloadCollector.push(item); reqLogger?.appendConvertedChunk?.(output); controller.enqueue(encoder.encode(output)); } @@ -666,6 +708,7 @@ export function createSSEStream(options: StreamOptions = {}) { if (flushed?.length > 0) { for (const item of flushed) { const output = formatSSE(item, sourceFormat); + clientPayloadCollector.push(item); reqLogger?.appendConvertedChunk?.(output); controller.enqueue(encoder.encode(output)); } @@ -684,6 +727,7 @@ export function createSSEStream(options: StreamOptions = {}) { // Send [DONE] (only if not already sent during transform) if (!doneSent) { doneSent = true; + clientPayloadCollector.push({ done: true }); const doneOutput = "data: [DONE]\n\n"; reqLogger?.appendConvertedChunk?.(doneOutput); controller.enqueue(encoder.encode(doneOutput)); @@ -747,7 +791,13 @@ export function createSSEStream(options: StreamOptions = {}) { }, _streamed: true, }; - onComplete({ status: 200, usage: state?.usage, responseBody }); + onComplete({ + status: 200, + usage: state?.usage, + responseBody, + providerPayload: providerPayloadCollector.build(), + clientPayload: clientPayloadCollector.build(responseBody), + }); } catch {} } } catch (error) { diff --git a/open-sse/utils/streamPayloadCollector.ts b/open-sse/utils/streamPayloadCollector.ts new file mode 100644 index 0000000000..c5d1446c22 --- /dev/null +++ b/open-sse/utils/streamPayloadCollector.ts @@ -0,0 +1,72 @@ +import { cloneLogPayload } from "@/lib/logPayloads"; + +type StructuredSSEEvent = { + index: number; + event?: string; + data: unknown; +}; + +type CollectorOptions = { + maxEvents?: number; + maxBytes?: number; + stage?: string; +}; + +function getEventName(payload: unknown): string | undefined { + if (!payload || typeof payload !== "object" || Array.isArray(payload)) return undefined; + + if (typeof (payload as { event?: unknown }).event === "string") { + return (payload as { event: string }).event; + } + if (typeof (payload as { type?: unknown }).type === "string") { + return (payload as { type: string }).type; + } + if ((payload as { done?: unknown }).done === true) { + return "[DONE]"; + } + return undefined; +} + +export function createStructuredSSECollector(options: CollectorOptions = {}) { + const { maxEvents = 200, maxBytes = 49152, stage } = options; + const events: StructuredSSEEvent[] = []; + let usedBytes = 0; + let droppedEvents = 0; + + return { + push(payload: unknown, explicitEvent?: string) { + if (payload === null || payload === undefined) return; + + const event: StructuredSSEEvent = { + index: events.length + droppedEvents, + data: cloneLogPayload(payload), + }; + + const eventName = explicitEvent || getEventName(payload); + if (eventName) { + event.event = eventName; + } + + const serializedSize = JSON.stringify(event).length; + if (events.length >= maxEvents || usedBytes + serializedSize > maxBytes) { + droppedEvents += 1; + return; + } + + usedBytes += serializedSize; + events.push(event); + }, + + build(summary?: unknown) { + return { + _streamed: true, + _format: "sse-json", + ...(stage ? { _stage: stage } : {}), + _eventCount: events.length + droppedEvents, + ...(droppedEvents > 0 ? { _truncated: true, _droppedEvents: droppedEvents } : {}), + events, + ...(summary === undefined ? {} : { summary: cloneLogPayload(summary) }), + }; + }, + }; +} diff --git a/src/app/api/logs/detail/route.ts b/src/app/api/logs/detail/route.ts index 6e4cd5e6db..e2734b5cd3 100644 --- a/src/app/api/logs/detail/route.ts +++ b/src/app/api/logs/detail/route.ts @@ -1,10 +1,9 @@ /** - * GET /api/logs/detail — List detailed request logs - * GET /api/logs/detail/:id — Get specific detailed log - * POST /api/logs/detail/toggle — Enable/disable detailed logging + * GET /api/logs/detail — List detailed request logs + current enabled flag + * POST /api/logs/detail — Enable/disable detailed logging */ import { NextRequest, NextResponse } from "next/server"; -import { isAuthenticated } from "@/shared/utils/apiAuth"; +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; import { getRequestDetailLogs, getRequestDetailLogCount, @@ -15,9 +14,8 @@ import { updateSettings } from "@/lib/db/settings"; export const dynamic = "force-dynamic"; export async function GET(req: NextRequest) { - if (!isAuthenticated(req)) { - return NextResponse.json({ error: "Unauthorized" }, { status: 401 }); - } + const authError = await requireManagementAuth(req); + if (authError) return authError; const url = new URL(req.url); const limit = Math.min(Number(url.searchParams.get("limit") ?? 50), 200); @@ -31,9 +29,8 @@ export async function GET(req: NextRequest) { } export async function POST(req: NextRequest) { - if (!isAuthenticated(req)) { - return NextResponse.json({ error: "Unauthorized" }, { status: 401 }); - } + const authError = await requireManagementAuth(req); + if (authError) return authError; const body = await req.json(); const enabled = body.enabled === true || body.enabled === "1"; diff --git a/src/app/api/usage/call-logs/[id]/route.ts b/src/app/api/usage/call-logs/[id]/route.ts index f975b15f47..ee9a864b80 100644 --- a/src/app/api/usage/call-logs/[id]/route.ts +++ b/src/app/api/usage/call-logs/[id]/route.ts @@ -1,8 +1,12 @@ import { NextResponse } from "next/server"; +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; import { getCallLogById } from "@/lib/usageDb"; export async function GET(request, { params }) { try { + const authError = await requireManagementAuth(request); + if (authError) return authError; + const { id } = await params; const log = await getCallLogById(id); diff --git a/src/app/api/usage/call-logs/route.ts b/src/app/api/usage/call-logs/route.ts index 23b62ca24e..671e5d2e5e 100644 --- a/src/app/api/usage/call-logs/route.ts +++ b/src/app/api/usage/call-logs/route.ts @@ -1,8 +1,12 @@ import { NextResponse } from "next/server"; +import { requireManagementAuth } from "@/lib/api/requireManagementAuth"; import { getCallLogs } from "@/lib/usageDb"; export async function GET(request: Request) { try { + const authError = await requireManagementAuth(request); + if (authError) return authError; + const { searchParams } = new URL(request.url); const filter: Record = {}; diff --git a/src/app/api/v1/completions/route.ts b/src/app/api/v1/completions/route.ts index 69d79a2777..86eafda624 100644 --- a/src/app/api/v1/completions/route.ts +++ b/src/app/api/v1/completions/route.ts @@ -1,5 +1,5 @@ import { CORS_ORIGIN, CORS_HEADERS } from "@/shared/utils/cors"; -import { handleChat } from "@/sse/handlers/chat"; +import { buildClientRawRequest, handleChat } from "@/sse/handlers/chat"; import { initTranslators } from "@omniroute/open-sse/translator/index.ts"; import { createInjectionGuard } from "@/middleware/promptInjectionGuard"; @@ -75,7 +75,7 @@ export async function POST(request: Request) { headers: request.headers, body: JSON.stringify(normalized), }); - return await handleChat(newRequest); + return await handleChat(newRequest, buildClientRawRequest(request, body)); } } } catch (error) { diff --git a/src/app/api/v1/providers/[provider]/chat/completions/route.ts b/src/app/api/v1/providers/[provider]/chat/completions/route.ts index c111e2e3e6..8bf8924922 100644 --- a/src/app/api/v1/providers/[provider]/chat/completions/route.ts +++ b/src/app/api/v1/providers/[provider]/chat/completions/route.ts @@ -1,5 +1,5 @@ import { CORS_ORIGIN } from "@/shared/utils/cors"; -import { handleChat } from "@/sse/handlers/chat"; +import { buildClientRawRequest, handleChat } from "@/sse/handlers/chat"; import { initTranslators } from "@omniroute/open-sse/translator/index.ts"; import { errorResponse } from "@omniroute/open-sse/utils/error.ts"; import { HTTP_STATUS } from "@omniroute/open-sse/config/constants.ts"; @@ -91,5 +91,5 @@ export async function POST(request, { params }) { body: JSON.stringify(body), }); - return await handleChat(newRequest); + return await handleChat(newRequest, buildClientRawRequest(request, rawBody)); } diff --git a/src/app/api/v1beta/models/[...path]/route.ts b/src/app/api/v1beta/models/[...path]/route.ts index 9b60596e78..96c76a8417 100644 --- a/src/app/api/v1beta/models/[...path]/route.ts +++ b/src/app/api/v1beta/models/[...path]/route.ts @@ -1,5 +1,5 @@ import { CORS_ORIGIN } from "@/shared/utils/cors"; -import { handleChat } from "@/sse/handlers/chat"; +import { buildClientRawRequest, handleChat } from "@/sse/handlers/chat"; import { initTranslators } from "@omniroute/open-sse/translator/index.ts"; import { v1betaGeminiGenerateSchema } from "@/shared/validation/schemas"; import { isValidationFailure, validateBody } from "@/shared/validation/helpers"; @@ -87,7 +87,7 @@ export async function POST(request, { params }) { body: JSON.stringify(convertedBody), }); - return await handleChat(newRequest); + return await handleChat(newRequest, buildClientRawRequest(request, rawBody)); } catch (error) { console.log("Error handling Gemini request:", error); return Response.json({ error: { message: error.message, code: 500 } }, { status: 500 }); diff --git a/src/lib/db/detailedLogs.ts b/src/lib/db/detailedLogs.ts index a49e1d038b..0573090227 100644 --- a/src/lib/db/detailedLogs.ts +++ b/src/lib/db/detailedLogs.ts @@ -8,20 +8,28 @@ import { v4 as uuidv4 } from "uuid"; import { getDbInstance } from "./core"; import { getSettings } from "./settings"; +import { isNoLog } from "../compliance"; +import { + protectPayloadForLog, + serializePayloadForStorage, + parseStoredPayload, +} from "../logPayloads"; export interface RequestDetailLog { id?: string; call_log_id?: string | null; timestamp?: string; - client_request?: string | null; - translated_request?: string | null; - provider_response?: string | null; - client_response?: string | null; + client_request?: unknown | null; + translated_request?: unknown | null; + provider_response?: unknown | null; + client_response?: unknown | null; provider?: string | null; model?: string | null; source_format?: string | null; target_format?: string | null; duration_ms?: number; + api_key_id?: string | null; + no_log?: boolean; } /** Returns true if detailed logging is enabled in settings */ @@ -37,16 +45,14 @@ export async function isDetailedLoggingEnabled(): Promise { /** Save a detailed log entry — caller must verify isDetailedLoggingEnabled() first */ export function saveRequestDetailLog(entry: RequestDetailLog): void { + const noLogEnabled = + Boolean(entry.no_log) || (entry.api_key_id ? isNoLog(entry.api_key_id) : false); + if (noLogEnabled) return; + const db = getDbInstance(); const id = entry.id ?? uuidv4(); const timestamp = entry.timestamp ?? new Date().toISOString(); - // Trim large bodies to avoid excessive disk usage (max 64KB each) - const trim = (s: string | null | undefined, max = 65536): string | null => { - if (!s) return null; - return s.length > max ? s.slice(0, max) + "…[truncated]" : s; - }; - db.prepare( ` INSERT INTO request_detail_logs @@ -58,10 +64,10 @@ export function saveRequestDetailLog(entry: RequestDetailLog): void { id, entry.call_log_id ?? null, timestamp, - trim(entry.client_request), - trim(entry.translated_request), - trim(entry.provider_response), - trim(entry.client_response), + serializePayloadForStorage(protectPayloadForLog(entry.client_request)), + serializePayloadForStorage(protectPayloadForLog(entry.translated_request)), + serializePayloadForStorage(protectPayloadForLog(entry.provider_response)), + serializePayloadForStorage(protectPayloadForLog(entry.client_response)), entry.provider ?? null, entry.model ?? null, entry.source_format ?? null, @@ -73,7 +79,7 @@ export function saveRequestDetailLog(entry: RequestDetailLog): void { /** Fetch detailed logs (latest first) */ export function getRequestDetailLogs(limit = 50, offset = 0): RequestDetailLog[] { const db = getDbInstance(); - return db + const rows = db .prepare( ` SELECT * FROM request_detail_logs @@ -81,14 +87,34 @@ export function getRequestDetailLogs(limit = 50, offset = 0): RequestDetailLog[] LIMIT ? OFFSET ? ` ) - .all(limit, offset) as RequestDetailLog[]; + .all(limit, offset) as Array>; + + return rows.map(mapDetailedLogRow); } /** Get a single detailed log by ID */ export function getRequestDetailLogById(id: string): RequestDetailLog | null { const db = getDbInstance(); - return (db.prepare("SELECT * FROM request_detail_logs WHERE id = ?").get(id) ?? - null) as RequestDetailLog | null; + const row = db.prepare("SELECT * FROM request_detail_logs WHERE id = ?").get(id) as + | Record + | undefined; + return row ? mapDetailedLogRow(row) : null; +} + +/** Get the most recent detailed log for a call log ID */ +export function getRequestDetailLogByCallLogId(callLogId: string): RequestDetailLog | null { + const db = getDbInstance(); + const row = db + .prepare( + ` + SELECT * FROM request_detail_logs + WHERE call_log_id = ? + ORDER BY timestamp DESC + LIMIT 1 + ` + ) + .get(callLogId) as Record | undefined; + return row ? mapDetailedLogRow(row) : null; } /** Get total count of detailed logs */ @@ -99,3 +125,20 @@ export function getRequestDetailLogCount(): number { }; return row?.cnt ?? 0; } + +function mapDetailedLogRow(row: Record): RequestDetailLog { + return { + id: typeof row.id === "string" ? row.id : undefined, + call_log_id: typeof row.call_log_id === "string" ? row.call_log_id : null, + timestamp: typeof row.timestamp === "string" ? row.timestamp : undefined, + client_request: parseStoredPayload(row.client_request), + translated_request: parseStoredPayload(row.translated_request), + provider_response: parseStoredPayload(row.provider_response), + client_response: parseStoredPayload(row.client_response), + provider: typeof row.provider === "string" ? row.provider : null, + model: typeof row.model === "string" ? row.model : null, + source_format: typeof row.source_format === "string" ? row.source_format : null, + target_format: typeof row.target_format === "string" ? row.target_format : null, + duration_ms: typeof row.duration_ms === "number" ? row.duration_ms : 0, + }; +} diff --git a/src/lib/logPayloads.ts b/src/lib/logPayloads.ts new file mode 100644 index 0000000000..1582b2fe77 --- /dev/null +++ b/src/lib/logPayloads.ts @@ -0,0 +1,109 @@ +import { sanitizePII } from "./piiSanitizer"; + +const SENSITIVE_KEYS = new Set([ + "api_key", + "apiKey", + "api-key", + "authorization", + "Authorization", + "x-api-key", + "X-Api-Key", + "access_token", + "accessToken", + "refresh_token", + "refreshToken", + "password", + "secret", + "token", +]); + +type JsonRecord = Record; + +export function cloneLogPayload(value: T): T { + if (value === null || value === undefined) return value; + if (typeof globalThis.structuredClone === "function") { + return globalThis.structuredClone(value); + } + return JSON.parse(JSON.stringify(value)) as T; +} + +export function normalizePayloadForLog(payload: unknown): unknown { + if (typeof payload !== "string") return payload; + + const trimmed = payload.trim(); + if (!trimmed) return ""; + + try { + return JSON.parse(trimmed); + } catch { + return { _rawText: payload }; + } +} + +export function redactPayload(payload: unknown): unknown { + if (!payload || typeof payload !== "object") return payload; + if (Array.isArray(payload)) return payload.map(redactPayload); + + const redacted: JsonRecord = {}; + for (const [key, value] of Object.entries(payload)) { + if (SENSITIVE_KEYS.has(key)) { + redacted[key] = "[REDACTED]"; + } else if (typeof value === "string" && value.startsWith("Bearer ")) { + redacted[key] = "Bearer [REDACTED]"; + } else if (typeof value === "object" && value !== null) { + redacted[key] = redactPayload(value); + } else { + redacted[key] = value; + } + } + return redacted; +} + +export function sanitizePayloadPII(payload: unknown): unknown { + if (typeof payload === "string") { + return sanitizePII(payload).text; + } + if (Array.isArray(payload)) { + return payload.map(sanitizePayloadPII); + } + if (!payload || typeof payload !== "object") { + return payload; + } + + const sanitized: JsonRecord = {}; + for (const [key, value] of Object.entries(payload)) { + sanitized[key] = sanitizePayloadPII(value); + } + return sanitized; +} + +export function protectPayloadForLog(payload: unknown): unknown { + if (payload === null || payload === undefined) return null; + const normalized = normalizePayloadForLog(payload); + const piiSanitized = sanitizePayloadPII(normalized); + return redactPayload(piiSanitized); +} + +export function serializePayloadForStorage(payload: unknown, maxLength = 65536): string | null { + if (payload === null || payload === undefined) return null; + + const exact = JSON.stringify(payload); + if (exact.length <= maxLength) { + return exact; + } + + return JSON.stringify({ + _truncated: true, + _originalSize: exact.length, + _preview: exact.slice(0, maxLength), + }); +} + +export function parseStoredPayload(value: unknown): unknown | null { + if (typeof value !== "string" || value.trim().length === 0) return null; + try { + return JSON.parse(value); + } catch { + return { _rawText: value }; + } +} diff --git a/src/lib/usage/callLogs.ts b/src/lib/usage/callLogs.ts index 35dbfedd68..7b1d6aa870 100644 --- a/src/lib/usage/callLogs.ts +++ b/src/lib/usage/callLogs.ts @@ -11,10 +11,12 @@ import path from "path"; import fs from "fs"; import { getDbInstance } from "../db/core"; import { getSettings } from "../db/settings"; +import { getRequestDetailLogByCallLogId } from "../db/detailedLogs"; import { shouldPersistToDisk, CALL_LOGS_DIR } from "./migrations"; import { getLoggedInputTokens, getLoggedOutputTokens } from "./tokenAccounting"; import { isNoLog } from "../compliance"; import { sanitizePII } from "../piiSanitizer"; +import { protectPayloadForLog, parseStoredPayload } from "../logPayloads"; type JsonRecord = Record; @@ -35,15 +37,6 @@ function toStringOrNull(value: unknown): string | null { return typeof value === "string" ? value : null; } -function parseJsonString(value: unknown): unknown | null { - if (typeof value !== "string" || value.trim().length === 0) return null; - try { - return JSON.parse(value); - } catch { - return null; - } -} - function hasTruncatedFlag(value: unknown): boolean { if (!value || typeof value !== "object" || Array.isArray(value)) return false; return (value as Record)._truncated === true; @@ -108,80 +101,6 @@ export function invalidateCallLogsMaxCache(): void { expiresAt: 0, }; } - -/** Fields that should always be redacted from logged payloads */ -const SENSITIVE_KEYS = new Set([ - "api_key", - "apiKey", - "api-key", - "authorization", - "Authorization", - "x-api-key", - "X-Api-Key", - "access_token", - "accessToken", - "refresh_token", - "refreshToken", - "password", - "secret", - "token", -]); - -/** - * Redact sensitive fields from a payload before persistence. - */ -function redactPayload(obj: any): any { - if (!obj || typeof obj !== "object") return obj; - if (Array.isArray(obj)) return obj.map(redactPayload); - - const redacted: Record = {}; - for (const [key, value] of Object.entries(obj)) { - if (SENSITIVE_KEYS.has(key)) { - redacted[key] = "[REDACTED]"; - } else if (typeof value === "string" && value.startsWith("Bearer ")) { - redacted[key] = "Bearer [REDACTED]"; - } else if (typeof value === "object" && value !== null) { - redacted[key] = redactPayload(value); - } else { - redacted[key] = value; - } - } - return redacted; -} - -/** - * Recursively sanitize PII from string fields in a payload. - * Uses lib/piiSanitizer config flags to determine if redaction is enabled. - */ -function sanitizePayloadPII(obj: any): any { - if (typeof obj === "string") { - return sanitizePII(obj).text; - } - if (Array.isArray(obj)) { - return obj.map(sanitizePayloadPII); - } - if (!obj || typeof obj !== "object") { - return obj; - } - - const sanitized: Record = {}; - for (const [key, value] of Object.entries(obj)) { - sanitized[key] = sanitizePayloadPII(value); - } - return sanitized; -} - -/** - * Apply payload protection chain before persistence. - * 1) Optional PII sanitization - * 2) Mandatory key/token redaction - */ -function protectPayloadForLog(payload: any): any { - if (!payload || !shouldLogPayloadInDb) return null; - const piiSanitized = sanitizePayloadPII(payload); - return redactPayload(piiSanitized); -} - let logIdCounter = 0; function generateLogId() { logIdCounter++; @@ -198,8 +117,10 @@ export async function saveCallLog(entry: any) { const apiKeyId = entry.apiKeyId || null; const noLogEnabled = Boolean(entry.noLog) || (apiKeyId ? isNoLog(apiKeyId) : false); - const protectedRequestBody = noLogEnabled ? null : protectPayloadForLog(entry.requestBody); - const protectedResponseBody = noLogEnabled ? null : protectPayloadForLog(entry.responseBody); + const protectedRequestBody = + noLogEnabled || !shouldLogPayloadInDb ? null : protectPayloadForLog(entry.requestBody); + const protectedResponseBody = + noLogEnabled || !shouldLogPayloadInDb ? null : protectPayloadForLog(entry.responseBody); // Resolve account name let account = entry.connectionId ? entry.connectionId.slice(0, 8) : "-"; @@ -227,7 +148,7 @@ export async function saveCallLog(entry: any) { }; const logEntry = { - id: generateLogId(), + id: typeof entry.id === "string" && entry.id.length > 0 ? entry.id : generateLogId(), timestamp: new Date().toISOString(), method: entry.method || "POST", path: entry.path || "/v1/chat/completions", @@ -470,8 +391,8 @@ export async function getCallLogById(id: string) { apiKeyId: toStringOrNull(entryRow.api_key_id), apiKeyName: toStringOrNull(entryRow.api_key_name), comboName: toStringOrNull(entryRow.combo_name), - requestBody: parseJsonString(entryRow.request_body), - responseBody: parseJsonString(entryRow.response_body), + requestBody: parseStoredPayload(entryRow.request_body), + responseBody: parseStoredPayload(entryRow.response_body), error: toStringOrNull(entryRow.error), }; @@ -492,7 +413,20 @@ export async function getCallLogById(id: string) { } } - return entry; + const detailed = getRequestDetailLogByCallLogId(id); + if (!detailed) { + return entry; + } + + return { + ...entry, + pipelinePayloads: { + clientRequest: detailed.client_request ?? null, + providerRequest: detailed.translated_request ?? null, + providerResponse: detailed.provider_response ?? null, + clientResponse: detailed.client_response ?? null, + }, + }; } /** diff --git a/src/shared/components/RequestLoggerDetail.tsx b/src/shared/components/RequestLoggerDetail.tsx index e6274c78df..e4cf1fbfa0 100644 --- a/src/shared/components/RequestLoggerDetail.tsx +++ b/src/shared/components/RequestLoggerDetail.tsx @@ -80,8 +80,42 @@ export default function RequestLoggerDetail({ log, detail, loading, onClose, onC } }; - const requestJson = detail?.requestBody ? JSON.stringify(detail.requestBody, null, 2) : null; - const responseJson = detail?.responseBody ? JSON.stringify(detail.responseBody, null, 2) : null; + const toPrettyJson = (payload) => { + if (payload === null || payload === undefined) return null; + try { + return JSON.stringify(payload, null, 2); + } catch { + return String(payload); + } + }; + + const pipelinePayloads = detail?.pipelinePayloads || null; + const payloadSections = pipelinePayloads + ? [ + { + key: "client-request", + title: "Client Request", + json: toPrettyJson(pipelinePayloads.clientRequest), + }, + { + key: "provider-request", + title: "Provider Request", + json: toPrettyJson(pipelinePayloads.providerRequest), + }, + { + key: "provider-response", + title: "Provider Response", + json: toPrettyJson(pipelinePayloads.providerResponse), + }, + { + key: "client-response", + title: "Client Response", + json: toPrettyJson(pipelinePayloads.clientResponse), + }, + ].filter((section) => section.json) + : []; + const requestJson = detail?.requestBody ? toPrettyJson(detail.requestBody) : null; + const responseJson = detail?.responseBody ? toPrettyJson(detail.responseBody) : null; return (
) : ( <> - {/* Response Payload (返回) — show first */} - {responseJson && ( + {payloadSections.length > 0 && + payloadSections.map((section) => ( + onCopy(section.json)} + /> + ))} + + {payloadSections.length === 0 && responseJson && ( onCopy(responseJson)} /> )} - {/* Request Payload (请求) */} - {requestJson && ( + {payloadSections.length === 0 && requestJson && ( onCopy(requestJson)} /> )} - {!requestJson && !responseJson && !loading && ( + {payloadSections.length === 0 && !requestJson && !responseJson && !loading && (
info

No payload data available for this log entry.

- Request/response bodies are only captured for non-streaming calls or when - streaming completes normally. + Enable detailed logging first if you want the four-stage client/provider payload + view for new requests.

)} diff --git a/src/shared/components/RequestLoggerV2.tsx b/src/shared/components/RequestLoggerV2.tsx index cd34ebf418..e8a1ca366d 100644 --- a/src/shared/components/RequestLoggerV2.tsx +++ b/src/shared/components/RequestLoggerV2.tsx @@ -93,6 +93,9 @@ export default function RequestLoggerV2() { const [selectedLog, setSelectedLog] = useState(null); const [detailLoading, setDetailLoading] = useState(false); const [detailData, setDetailData] = useState(null); + const [detailLoggingEnabled, setDetailLoggingEnabled] = useState(false); + const [detailLoggingLoading, setDetailLoggingLoading] = useState(false); + const [detailLoggingReady, setDetailLoggingReady] = useState(false); const intervalRef = useRef(null); const hasLoadedRef = useRef(false); const [providerNodes, setProviderNodes] = useState([]); @@ -161,6 +164,20 @@ export default function RequestLoggerV2() { .catch(() => {}); }, []); + useEffect(() => { + fetch("/api/logs/detail?limit=1") + .then(async (res) => { + if (!res.ok) return null; + return await res.json(); + }) + .then((data) => { + if (!data) return; + setDetailLoggingEnabled(data.enabled === true); + setDetailLoggingReady(true); + }) + .catch(() => {}); + }, []); + // Auto-refresh useEffect(() => { if (intervalRef.current) clearInterval(intervalRef.current); @@ -232,6 +249,25 @@ export default function RequestLoggerV2() { setDetailData(null); }; + const toggleDetailLogging = async () => { + setDetailLoggingLoading(true); + try { + const nextEnabled = !detailLoggingEnabled; + const res = await fetch("/api/logs/detail", { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ enabled: nextEnabled }), + }); + if (!res.ok) throw new Error("Failed to update detailed logging"); + setDetailLoggingEnabled(nextEnabled); + setDetailLoggingReady(true); + } catch (error) { + console.error("Failed to toggle detailed logging:", error); + } finally { + setDetailLoggingLoading(false); + } + }; + // Unique accounts and providers for dropdowns const uniqueAccounts = [...new Set(logs.map((l) => l.account).filter((a) => a && a !== "-"))]; @@ -271,6 +307,33 @@ export default function RequestLoggerV2() { {recording ? "Recording" : "Paused"} + + + {detailLoggingReady && ( + + New requests will {detailLoggingEnabled ? "" : "not "}capture client/provider pipeline + payloads. + + )} + {/* Search */}
diff --git a/src/sse/handlers/chat.ts b/src/sse/handlers/chat.ts index b4c0ca889a..eec303ffa5 100644 --- a/src/sse/handlers/chat.ts +++ b/src/sse/handlers/chat.ts @@ -44,6 +44,7 @@ import { RequestTelemetry, recordTelemetry } from "../../shared/utils/requestTel import { generateRequestId } from "../../shared/utils/requestId"; import { logAuditEvent } from "../../lib/compliance/index"; import { enforceApiKeyPolicy } from "../../shared/utils/apiKeyPolicy"; +import { cloneLogPayload } from "@/lib/logPayloads"; import { applyTaskAwareRouting, getTaskRoutingConfig, @@ -81,6 +82,13 @@ export async function handleChat(request: any, clientRawRequest: any = null) { return errorResponse(HTTP_STATUS.BAD_REQUEST, "Invalid JSON body"); } + const rawClientBody = cloneLogPayload(body); + + // Build clientRawRequest for logging (if not provided) + if (!clientRawRequest) { + clientRawRequest = buildClientRawRequest(request, rawClientBody); + } + // FASE-01: Input sanitization — prompt injection detection & PII redaction telemetry.startPhase("validate"); const sanitizeResult = sanitizeRequest(body, log as any); @@ -113,16 +121,6 @@ export async function handleChat(request: any, clientRawRequest: any = null) { ); } - // Build clientRawRequest for logging (if not provided) - if (!clientRawRequest) { - const url = new URL(request.url); - clientRawRequest = { - endpoint: url.pathname, - body, - headers: Object.fromEntries(request.headers.entries()), - }; - } - // Log request endpoint and model const url = new URL(request.url); const modelStr = body.model; @@ -344,6 +342,15 @@ export async function handleChat(request: any, clientRawRequest: any = null) { return withSessionHeader(response, sessionId); } +export function buildClientRawRequest(request: Request, body: unknown) { + const url = new URL(request.url); + return { + endpoint: url.pathname, + body: cloneLogPayload(body), + headers: Object.fromEntries(request.headers.entries()), + }; +} + /** * Handle single model chat request * diff --git a/tests/unit/request-log-payloads.test.mjs b/tests/unit/request-log-payloads.test.mjs new file mode 100644 index 0000000000..d06155456c --- /dev/null +++ b/tests/unit/request-log-payloads.test.mjs @@ -0,0 +1,65 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +const { + normalizePayloadForLog, + protectPayloadForLog, + serializePayloadForStorage, + parseStoredPayload, +} = await import("../../src/lib/logPayloads.ts"); +const { createStructuredSSECollector } = + await import("../../open-sse/utils/streamPayloadCollector.ts"); + +test("normalizes JSON strings before log protection and redacts sensitive keys", () => { + const protectedPayload = protectPayloadForLog( + JSON.stringify({ + authorization: "Bearer secret-token-value", + nested: { + apiKey: "top-secret-key", + }, + }) + ); + + assert.deepEqual(protectedPayload, { + authorization: "[REDACTED]", + nested: { + apiKey: "[REDACTED]", + }, + }); +}); + +test("wraps raw text payloads in JSON-safe objects", () => { + const normalized = normalizePayloadForLog("event: ping\ndata: plain-text\n\n"); + + assert.deepEqual(normalized, { + _rawText: "event: ping\ndata: plain-text\n\n", + }); +}); + +test("serializes truncated payloads as valid JSON objects", () => { + const stored = serializePayloadForStorage({ text: "x".repeat(200) }, 80); + const parsed = parseStoredPayload(stored); + + assert.equal(parsed._truncated, true); + assert.equal(parsed._originalSize > 80, true); + assert.equal(typeof parsed._preview, "string"); +}); + +test("structured SSE collector preserves event order and marks truncation", () => { + const collector = createStructuredSSECollector({ maxEvents: 2, maxBytes: 200 }); + + collector.push({ type: "response.created", id: "r1" }); + collector.push({ type: "response.output_text.delta", delta: "hi" }); + collector.push({ type: "response.completed" }); + + const payload = collector.build({ done: true }); + + assert.equal(payload._streamed, true); + assert.equal(payload._eventCount, 3); + assert.equal(payload._truncated, true); + assert.equal(payload._droppedEvents, 1); + assert.equal(payload.events.length, 2); + assert.equal(payload.events[0].event, "response.created"); + assert.equal(payload.events[1].event, "response.output_text.delta"); + assert.deepEqual(payload.summary, { done: true }); +});