diff --git a/docs-site/src/content/docs/reference/cli.md b/docs-site/src/content/docs/reference/cli.md index 04bcd327..4abfab58 100644 --- a/docs-site/src/content/docs/reference/cli.md +++ b/docs-site/src/content/docs/reference/cli.md @@ -511,6 +511,32 @@ With no scope, `ocx debug` prints usage and, when the proxy is stopped, the next defaults. Provider debug defaults from `OCX_DEBUG=1` (legacy `OCX_DEBUG_FRAMES=1` also works); usage debug defaults from `OPENCODEX_USAGE_DEBUG=1`. +## Request traces + +`usage.jsonl` stays a compact telemetry log. Full prompts and responses go to a separate local store, +`trace.sqlite` in the opencodex config directory, linked to usage rows by `traceId` (the request id). +Tracing is **off by default** and is configured by environment variable at proxy start: + +| Variable | Default | Meaning | +| -------------------------- | -------- | ----------------------------------------------------------------------- | +| `OCX_TRACE` | `off` | `off`, `metadata`, `redacted`, or `full`. | +| `OCX_TRACE_TTL_HOURS` | `24` | Retention for stored traces (1–720). | +| `OCX_TRACE_MAX_BODY_BYTES` | `524288` | Cap per stored body. Hashes and byte counts always cover the full body. | +| `OCX_TRACE_MAX_DB_MB` | `256` | Total store cap; oldest traces are evicted first. | +| `OCX_TRACE_SAMPLE` | `1` | Fraction of requests whose bodies are stored in `redacted`/`full` mode. | + +- `metadata` stores nothing but adds a payload-free `trace` object to each usage row: byte sizes, + message/tool-call/attachment counts, request/outbound/response hashes, and `systemHash`, `toolsHash` + and `prefixHash` of the final provider wire body. Comparing these between two turns of one conversation + shows which section broke a provider prompt-cache prefix. +- `redacted` also stores the inbound request, the final provider body, and the response, with credentials + and token-shaped values removed. Free-text secrets in prompts that do not match a known pattern are not + detectable, so treat the store as sensitive. +- `full` stores bodies verbatim. Use it deliberately and briefly. + +Outbound capture covers adapters that send through the shared upstream fetch helper; the inbound request +and response are captured for `/v1/responses`, `/v1/messages`, and `/v1/chat/completions`. + ## Updating ### `ocx update` diff --git a/src/server/index.ts b/src/server/index.ts index a35625d4..3a5ed9f9 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -1,4 +1,5 @@ import { markActivity } from "../lib/sidecar-tracker"; +import { beginTrace, runWithTrace } from "../trace/capture"; import { initServerSentry } from "../telemetry/sentry-server"; import { buildWarmupCompletionFrames, @@ -350,6 +351,12 @@ function attachLiveSidebandUpstream(ws: ServerWebSocket): void { // trackSseForRequestLog( // export function relaySseWithHeartbeat +/** + * Start the proxy, management API, and dashboard using the loaded configuration. + * Apply startup migrations and initialize background maintenance. `port` overrides + * the configured port (default 10100); 0 requests an available port. Return the + * listening Bun server. Synchronous initialization and listen errors propagate. + */ export function startServer(port?: number) { initServerSentry(); const config = runAlibabaRegionStartupMigration( @@ -1198,7 +1205,8 @@ export function startServer(port?: number) { const responsesWorkReq = responsesCacheProbe?.request ?? req; const start = Date.now(); const requestId = nextRequestLogId(start); - const logCtx = { model: "unknown", provider: "unknown" }; + const logCtx: RequestLogContext = { model: "unknown", provider: "unknown" }; + await beginTrace(logCtx, responsesWorkReq); let logged = false; const finalizeNativePassthroughLog = ( status: number, @@ -1211,7 +1219,7 @@ export function startServer(port?: number) { logged = true; addFinalRequestLog(requestId, start, logCtx, status, meta); }; - const response = await handleResponses( + const response = await runWithTrace(logCtx, () => handleResponses( responsesWorkReq, config, logCtx, @@ -1233,7 +1241,7 @@ export function startServer(port?: number) { }); }, }, - ); + )); responsesCacheProbe?.store(response); return withCors( responseWithDeferredRequestLog(response, requestId, start, logCtx), @@ -1341,15 +1349,16 @@ export function startServer(port?: number) { model: "unknown", provider: "unknown", }; + await beginTrace(logCtx, messagesWorkReq); // Logging is finalized inside handleClaudeMessages (Responses-vocab tap on the // pre-translation stream + native passthrough callbacks) — do not re-wrap the // translated Anthropic stream here. - const response = await handleClaudeMessages( + const response = await runWithTrace(logCtx, () => handleClaudeMessages( messagesWorkReq, config, logCtx, { requestId, start }, - ); + )); messagesCacheProbe?.store(response); return withCors(response, req, config); } @@ -1396,12 +1405,13 @@ export function startServer(port?: number) { model: "unknown", provider: "unknown", }; - const response = await handleChatCompletions( + await beginTrace(logCtx, chatWorkReq); + const response = await runWithTrace(logCtx, () => handleChatCompletions( chatWorkReq, config, logCtx, { requestId, start }, - ); + )); chatCacheProbe?.store(response); return withCors(response, req, config); } diff --git a/src/server/request-log.ts b/src/server/request-log.ts index 37db9d4d..d4c9dac3 100644 --- a/src/server/request-log.ts +++ b/src/server/request-log.ts @@ -11,6 +11,8 @@ import { readCodexCatalogPath } from "../codex/catalog"; import type { OcxConfig, OcxUsage } from "../types"; import type { AdapterRequest } from "../adapters/base"; import { redactSecretString } from "../lib/redact"; +import { appendTraceResponse, finalizeTrace, type TraceCapture } from "../trace/capture"; +import type { UsageTraceMeta } from "../trace/types"; import { providerAccountLabel, baseProviderLabel } from "../providers/label"; import { getAccountSet } from "../oauth/store"; import { @@ -109,6 +111,8 @@ export interface RequestLogContext { affinity?: "reused" | "new_bind" | "rebound" | "cleared"; transportPhase?: "pre_headers" | "mid_stream" | "terminal_sse"; terminalSource?: "upstream" | "synthetic"; + /** Internal per-request trace capture (src/trace); never persisted or serialized. */ + trace?: TraceCapture; } export interface RequestLogEntry { @@ -160,6 +164,10 @@ export interface RequestLogEntry { transportPhase?: "pre_headers" | "mid_stream" | "terminal_sse"; /** Whether the terminal came from a real upstream SSE event or a proxy synthetic tail. */ terminalSource?: "upstream" | "synthetic"; + /** Key into trace.sqlite when this request's bodies were stored. */ + traceId?: string; + /** Payload-free trace summary (sizes, counts, hashes). */ + trace?: UsageTraceMeta; } const requestLog: RequestLogEntry[] = []; @@ -316,6 +324,8 @@ export function requestLogEntryFromPersistedUsage(entry: PersistedUsageEntry): R ...(entry.usage ? { usage: entry.usage } : {}), ...(entry.totalTokens !== undefined ? { totalTokens: entry.totalTokens } : {}), ...(entry.attempts?.length ? { attempts: entry.attempts } : {}), + ...(entry.traceId ? { traceId: entry.traceId } : {}), + ...(entry.trace ? { trace: entry.trace } : {}), }; } @@ -432,6 +442,8 @@ export function addRequestLog(entry: RequestLogEntry) { ...(entry.totalTokens !== undefined ? { totalTokens: entry.totalTokens } : {}), ...(entry.attempts?.length ? { attempts: entry.attempts } : {}), ...failureDiagnostics, + ...(entry.traceId ? { traceId: entry.traceId } : {}), + ...(entry.trace ? { trace: entry.trace } : {}), }); } catch { /* request logging must never fail a user request */ @@ -672,6 +684,11 @@ export function usageFromResponsesPayload(usage: unknown): OcxUsage | undefined return undefined; } +/** + * Update response metadata and error diagnostics from a complete response body, + * and append it to any active trace. Invalid JSON skips metadata extraction but + * still reaches error and trace capture. Retain a bounded sample when debug is enabled. + */ export function inspectResponseLogJson(logCtx: RequestLogContext, text: string): void { try { applyResponseLogMetadata(logCtx, JSON.parse(text)); @@ -679,12 +696,19 @@ export function inspectResponseLogJson(logCtx: RequestLogContext, text: string): /* body may not be JSON; request log metadata is best-effort only */ } captureUpstreamError(logCtx, text); + appendTraceResponse(logCtx.trace, text); if (isUsageDebugEnabled() && logCtx.usageDebugBodyKind === undefined) { logCtx.usageDebugBodyKind = "json"; logCtx.usageDebugBodySample = truncateForDebug(text); } } +/** + * Inspect one SSE data payload for response metadata and error diagnostics, + * append it to any active trace, and accumulate a bounded sample when debug is + * enabled. Ignore null, empty, and [DONE] payloads; invalid JSON still reaches + * error and trace capture. + */ export function inspectResponseLogSsePayload(logCtx: RequestLogContext, payload: string | null): void { if (!payload || payload.trim() === "[DONE]") return; const debugEnabled = isUsageDebugEnabled(); @@ -695,6 +719,7 @@ export function inspectResponseLogSsePayload(logCtx: RequestLogContext, payload: /* SSE block payload may not be JSON; metadata inspection is best-effort */ } captureUpstreamError(logCtx, payload); + appendTraceResponse(logCtx.trace, payload); if (debugEnabled) { if (!sseAlreadyMarked) { logCtx.usageDebugBodyKind = "sse"; @@ -831,6 +856,9 @@ export function httpStatusForRequestLogTerminal( /** * Finalizes and records a request log entry with status, diagnostics, usage, routing, and retry-attempt data. * + * Finalizes any trace before storing the entry and always releases the bound + * provider account. Errors from `addLog` propagate after account release. + * * @param requestId - The request identifier * @param start - The request start timestamp in milliseconds * @param logCtx - Mutable metadata collected during the request @@ -904,6 +932,13 @@ export function addFinalRequestLog( const loggedUsage = aggregate?.usage ?? existing.usage; const usageStatus = aggregate?.status ?? existing.status; const totalTokens = aggregate?.totalTokens ?? existing.totalTokens; + const traced = finalizeTrace(requestId, logCtx, { + ...(logCtx.conversationId ? { conversationId: logCtx.conversationId } : {}), + provider, + model, + status: effectiveStatus, + timestamp: start, + }); addLog({ requestId, timestamp: start, @@ -942,6 +977,8 @@ export function addFinalRequestLog( ...(logCtx.affinity ? { affinity: logCtx.affinity } : {}), ...(logCtx.transportPhase ? { transportPhase: logCtx.transportPhase } : {}), ...(logCtx.terminalSource ? { terminalSource: logCtx.terminalSource } : {}), + ...(traced.traceId ? { traceId: traced.traceId } : {}), + ...(traced.trace ? { trace: traced.trace } : {}), }); if (isUsageDebugEnabled()) { appendUsageDebug({ diff --git a/src/server/responses/fetch-helpers.ts b/src/server/responses/fetch-helpers.ts index 6a54fe02..cb80d7cc 100644 --- a/src/server/responses/fetch-helpers.ts +++ b/src/server/responses/fetch-helpers.ts @@ -1,6 +1,7 @@ import type { Server } from "bun"; import type { OcxProviderConfig } from "../../types"; import type { WsData } from "../ws-bridge"; +import { noteOutboundRequestBody } from "../../trace/capture"; export function disableResponsesRequestTimeout(req: Request, server: Pick, "timeout"> | undefined): boolean { @@ -31,6 +32,14 @@ export function providerFetch(provider: OcxProviderConfig): typeof globalThis.fe +/** + * Fetch with caller cancellation and a timeout of `timeoutMs` milliseconds until + * the executor settles; the timer does not cover reading the response body. + * Record supported outbound bodies in the active trace. `preferIdentityEncoding` + * adds Accept-Encoding: identity only when the caller has not set that header. + * Return the response unchanged, including HTTP errors; header-construction and + * executor errors (including abort/timeout rejections) propagate. + */ export async function fetchWithHeaderTimeout( url: string, init: Omit, @@ -44,6 +53,7 @@ export async function fetchWithHeaderTimeout( if (!timeout.signal.aborted) timeout.abort(new DOMException("Timeout elapsed", "TimeoutError")); }, timeoutMs); const headers = new Headers(init.headers); + noteOutboundRequestBody(init.body); // Compressed SSE can be held until the decompressor has a complete block. Streaming calls // default to identity for low-latency frame delivery, while an explicit caller choice wins. if (preferIdentityEncoding && !headers.has("accept-encoding")) { diff --git a/src/trace/capture.ts b/src/trace/capture.ts new file mode 100644 index 00000000..0693fc23 --- /dev/null +++ b/src/trace/capture.ts @@ -0,0 +1,479 @@ +/** + * Per-request trace capture. A `TraceCapture` hangs off the request log context (`logCtx.trace`) and is + * also installed in an AsyncLocalStorage so the shared upstream fetch helper can record the final provider + * wire body without every adapter knowing about tracing. + * + * Nothing runs when the mode is `off`. Hashes and counts are computed over the full bodies; only the + * stored copies are capped, and in `redacted` mode secrets are stripped before anything is written. + */ + +import { AsyncLocalStorage } from "node:async_hooks"; +import { createHash, type Hash } from "node:crypto"; +import { redactSecretString, redactSecrets } from "../lib/redact"; +import { getTraceSettings } from "./settings"; +import { writeTrace } from "./store"; +import type { TraceMode, UsageTraceMeta } from "./types"; + +/** Largest inbound body we will buffer for tracing; bigger requests are counted but not read. */ +const MAX_INBOUND_READ_BYTES = 8 * 1024 * 1024; +const WALK_NODE_BUDGET = 20_000; +const WALK_MAX_DEPTH = 8; + +export interface TraceCapture { + mode: Exclude; + /** Bodies are stored (redacted/full mode and sampled in). */ + persist: boolean; + maxBodyBytes: number; + inbound?: string; + inboundBytes?: number; + outbound?: string; + outboundBytes?: number; + outboundCount: number; + response: string; + responseBytes: number; + responseStoredBytes: number; + responseTruncated: boolean; + responseHasher: Hash; + done: boolean; +} + +const traceStorage = new AsyncLocalStorage(); + +function sha(text: string): string { + return createHash("sha256").update(text).digest("hex").slice(0, 32); +} + +function byteLength(text: string): number { + return Buffer.byteLength(text, "utf8"); +} + +/** + * Keep a prefix within `maxBytes` UTF-8 bytes without splitting a character. + * Return the retained byte count and whether any input was omitted. + */ +function truncateUtf8( + text: string, + maxBytes: number, +): { text: string; bytes: number; truncated: boolean } { + const encoded = new TextEncoder().encode(text); + if (encoded.byteLength <= maxBytes) { + return { text, bytes: encoded.byteLength, truncated: false }; + } + + let end = Math.max(0, Math.min(maxBytes, encoded.byteLength)); + while (end > 0) { + try { + const decoded = new TextDecoder("utf-8", { fatal: true }).decode( + encoded.subarray(0, end), + ); + return { text: decoded, bytes: end, truncated: true }; + } catch { + end -= 1; + } + } + return { text: "", bytes: 0, truncated: true }; +} + +/** + * Create a capture from current settings, or undefined when tracing is off. + * Sampling controls body persistence in redacted/full mode; metadata is still + * collected for requests that are not sampled. + */ +export function createTraceCapture(): TraceCapture | undefined { + const settings = getTraceSettings(); + if (settings.mode === "off") return undefined; + const wantsBodies = settings.mode === "redacted" || settings.mode === "full"; + return { + mode: settings.mode, + persist: wantsBodies && Math.random() < settings.sample, + maxBodyBytes: settings.maxBodyBytes, + outboundCount: 0, + response: "", + responseBytes: 0, + responseStoredBytes: 0, + responseTruncated: false, + responseHasher: createHash("sha256"), + done: false, + }; +} + +/** + * Attach a capture and read a clone of `req`, leaving the original body readable. + * Do nothing when tracing is off. Omit the inbound body and byte count if the + * declared or observed size exceeds 8 MiB. Failures are swallowed and may leave + * a capture without inbound data. + */ +export async function beginTrace( + logCtx: { trace?: TraceCapture }, + req: Request, +): Promise { + try { + const trace = createTraceCapture(); + if (!trace) return; + logCtx.trace = trace; + const contentLength = req.headers.get("content-length"); + if (contentLength !== null) { + const declared = Number(contentLength); + if (Number.isFinite(declared) && declared > MAX_INBOUND_READ_BYTES) return; + } + + const reader = req.clone().body?.getReader(); + if (!reader) { + trace.inbound = ""; + trace.inboundBytes = 0; + return; + } + + const chunks: Uint8Array[] = []; + let bytes = 0; + while (true) { + const { done, value } = await reader.read(); + if (done) break; + bytes += value.byteLength; + if (bytes > MAX_INBOUND_READ_BYTES) { + void reader.cancel().catch(() => {}); + return; + } + chunks.push(value); + } + + const body = new Uint8Array(bytes); + let offset = 0; + for (const chunk of chunks) { + body.set(chunk, offset); + offset += chunk.byteLength; + } + trace.inbound = new TextDecoder().decode(body); + trace.inboundBytes = bytes; + } catch { + /* tracing is best-effort */ + } +} + +/** + * Run `fn` with `logCtx.trace` available to outbound capture in its async context. + * Without a capture, call `fn` directly. Return its result unchanged, including + * promises, and propagate its errors. + */ +export function runWithTrace( + logCtx: { trace?: TraceCapture }, + fn: () => T, +): T { + return logCtx.trace ? traceStorage.run(logCtx.trace, fn) : fn(); +} + +/** + * Record a string or UTF-8 Uint8Array body in the active, unfinished capture. + * The last accepted body wins and increments the outbound count, including + * retries. Ignore other body types, absent/finished captures, and capture errors. + */ +export function noteOutboundRequestBody(body: unknown): void { + const trace = traceStorage.getStore(); + if (!trace || trace.done) return; + try { + let text: string | undefined; + if (typeof body === "string") text = body; + else if (body instanceof Uint8Array) text = new TextDecoder().decode(body); + if (text === undefined) return; + trace.outbound = text; + trace.outboundBytes = byteLength(text); + trace.outboundCount += 1; + } catch { + /* ignore */ + } +} + +/** + * Hash and count one response payload (SSE data block or JSON body) before + * redaction or truncation. For sampled body capture, retain a UTF-8 byte-capped + * copy, redacting each payload in redacted mode. Inserted newlines count toward + * storage limits but not response bytes or hashes. Ignore empty payloads, + * absent/finished captures, and capture errors. + */ +export function appendTraceResponse( + trace: TraceCapture | undefined, + chunk: string, +): void { + if (!trace || trace.done || !chunk) return; + try { + trace.responseHasher.update(chunk); + trace.responseBytes += byteLength(chunk); + if (!trace.persist) return; + const storedChunk = trace.mode === "redacted" + ? redactBodyForStorage(chunk) + : chunk; + const separator = trace.response ? "\n" : ""; + const separatorBytes = separator ? 1 : 0; + const remaining = trace.maxBodyBytes - trace.responseStoredBytes; + if (remaining <= separatorBytes) { + trace.responseTruncated = true; + return; + } + const piece = truncateUtf8(storedChunk, remaining - separatorBytes); + if (!piece.text) { + trace.responseTruncated = true; + return; + } + trace.response += separator + piece.text; + trace.responseStoredBytes += separatorBytes + piece.bytes; + if (piece.truncated) trace.responseTruncated = true; + } catch { + /* ignore */ + } +} + +const ATTACHMENT_TYPES = new Set([ + "image", + "image_url", + "input_image", + "input_file", + "input_audio", + "document", + "file", +]); +const TOOL_CALL_TYPES = new Set([ + "function_call", + "custom_tool_call", + "tool_use", + "mcp_call", + "computer_call", +]); + +export interface BodyShape { + messageCount?: number; + toolCallCount: number; + toolDefCount?: number; + attachmentCount: number; +} + +function isRecord(value: unknown): value is Record { + return value !== null && typeof value === "object" && !Array.isArray(value); +} + +/** + * Count top-level messages/input and tool definitions in a parsed request body. + * Count recognized tool calls and attachments outside tool definitions, visiting + * up to 20,000 objects/arrays through depth 8; these counts may be partial. + * Non-object or array roots yield zero calls/attachments and no other counts. + */ +export function shapeOf(parsed: unknown): BodyShape { + const shape: BodyShape = { toolCallCount: 0, attachmentCount: 0 }; + if (!isRecord(parsed)) return shape; + if (Array.isArray(parsed.messages)) + shape.messageCount = parsed.messages.length; + else if (Array.isArray(parsed.input)) + shape.messageCount = parsed.input.length; + else if (typeof parsed.input === "string") shape.messageCount = 1; + if (Array.isArray(parsed.tools)) shape.toolDefCount = parsed.tools.length; + let budget = WALK_NODE_BUDGET; + const visit = (node: unknown, depth: number): void => { + if ( + budget <= 0 || + depth > WALK_MAX_DEPTH || + node === null || + typeof node !== "object" + ) + return; + budget -= 1; + if (Array.isArray(node)) { + for (const item of node) visit(item, depth + 1); + return; + } + const rec = node as Record; + const type = rec.type; + if (typeof type === "string") { + if (ATTACHMENT_TYPES.has(type)) shape.attachmentCount += 1; + if (TOOL_CALL_TYPES.has(type)) shape.toolCallCount += 1; + } + if (Array.isArray(rec.tool_calls)) + shape.toolCallCount += rec.tool_calls.length; + for (const key of Object.keys(rec)) { + if (key === "tools") continue; // tool definitions are counted separately + visit(rec[key], depth + 1); + } + }; + visit(parsed, 0); + return shape; +} + +export interface SectionHashes { + systemHash?: string; + toolsHash?: string; + prefixHash?: string; +} + +/** + * Hash JSON-serialized system/instructions, tools, and messages/input excluding + * the last turn, using the first 32 hexadecimal SHA-256 characters. The prefix + * hash covers only that turn array, not system or tools; it does not establish + * provider cache eligibility. Missing sections are omitted; non-object or array + * roots return {}. JSON serialization errors propagate to the caller. + */ +export function sectionHashes(parsed: unknown): SectionHashes { + if (!isRecord(parsed)) return {}; + const out: SectionHashes = {}; + const system = parsed.system ?? parsed.instructions; + if (system !== undefined) out.systemHash = sha(JSON.stringify(system)); + if (Array.isArray(parsed.tools)) + out.toolsHash = sha(JSON.stringify(parsed.tools)); + const turns = Array.isArray(parsed.messages) + ? parsed.messages + : Array.isArray(parsed.input) + ? parsed.input + : undefined; + if (turns && turns.length > 0) + out.prefixHash = sha(JSON.stringify(turns.slice(0, -1))); + return out; +} + +/** Parse JSON, returning undefined for missing or invalid text. */ +function tryParse(text: string | undefined): unknown { + if (text === undefined) return undefined; + try { + return JSON.parse(text); + } catch { + return undefined; + } +} + +/** + * Redact known sensitive keys and value patterns in parsed JSON, or known + * value patterns in non-JSON text. Secrets matching neither remain unchanged. + */ +function redactBodyForStorage(text: string): string { + const parsed = tryParse(text); + return parsed !== undefined + ? JSON.stringify(redactSecrets(parsed)) + : redactSecretString(text); +} + +/** + * Prepare an optional body for storage, redacting in redacted mode before + * truncating to `max` UTF-8 bytes. An absent body is not marked truncated. + */ +function storableBody( + text: string | undefined, + mode: Exclude, + max: number, +): { text?: string; truncated: boolean } { + if (text === undefined) return { truncated: false }; + const out = mode === "redacted" ? redactBodyForStorage(text) : text; + const bounded = truncateUtf8(out, max); + return { text: bounded.text, truncated: bounded.truncated }; +} + +export interface TraceFinalizeInfo { + conversationId?: string; + provider?: string; + model?: string; + status?: number; + timestamp: number; +} + +export interface TraceFinalizeResult { + /** Present only when bodies were stored, so the link from usage.jsonl is always resolvable. */ + traceId?: string; + trace?: UsageTraceMeta; +} + +/** + * Finalize a capture once, returning usage metadata and, on a successful sampled + * body write, `requestId` as `traceId`. `info.timestamp` is the request start in + * Unix milliseconds; retention is measured from the write time. Shape counts use + * parsed inbound data with an outbound fallback; section hashes prefer outbound. + * Release buffered bodies even on failure. Missing/finished captures and capture + * errors return {}; a failed store write retains metadata with `stored: false`. + */ +export function finalizeTrace( + requestId: string, + logCtx: { trace?: TraceCapture }, + info: TraceFinalizeInfo, +): TraceFinalizeResult { + const trace = logCtx.trace; + if (!trace || trace.done) return {}; + trace.done = true; + try { + const inboundParsed = tryParse(trace.inbound); + const outboundParsed = tryParse(trace.outbound); + const wire = outboundParsed ?? inboundParsed; + const shape = shapeOf(inboundParsed ?? outboundParsed); + const hashes = sectionHashes(wire); + const meta: UsageTraceMeta = { + mode: trace.mode, + stored: false, + ...(trace.inboundBytes !== undefined + ? { requestBytes: trace.inboundBytes } + : {}), + ...(trace.outboundBytes !== undefined + ? { outboundBytes: trace.outboundBytes } + : {}), + ...(trace.responseBytes > 0 + ? { responseBytes: trace.responseBytes } + : {}), + ...(shape.messageCount !== undefined + ? { messageCount: shape.messageCount } + : {}), + toolCallCount: shape.toolCallCount, + ...(shape.toolDefCount !== undefined + ? { toolDefCount: shape.toolDefCount } + : {}), + attachmentCount: shape.attachmentCount, + ...(trace.outboundCount > 0 + ? { outboundCount: trace.outboundCount } + : {}), + ...(trace.inbound !== undefined + ? { requestHash: sha(trace.inbound) } + : {}), + ...(trace.outbound !== undefined + ? { outboundHash: sha(trace.outbound) } + : {}), + ...(trace.responseBytes > 0 + ? { responseHash: trace.responseHasher.digest("hex").slice(0, 32) } + : {}), + ...hashes, + }; + if (trace.persist) { + const inbound = storableBody( + trace.inbound, + trace.mode, + trace.maxBodyBytes, + ); + const outbound = storableBody( + trace.outbound, + trace.mode, + trace.maxBodyBytes, + ); + const response = storableBody( + trace.response || undefined, + trace.mode, + trace.maxBodyBytes, + ); + meta.stored = writeTrace({ + traceId: requestId, + createdAt: info.timestamp, + mode: trace.mode, + ...(info.conversationId ? { conversationId: info.conversationId } : {}), + ...(info.provider ? { provider: info.provider } : {}), + ...(info.model ? { model: info.model } : {}), + ...(info.status !== undefined ? { status: info.status } : {}), + meta: { ...meta, stored: true }, + ...(inbound.text !== undefined ? { inbound: inbound.text } : {}), + ...(outbound.text !== undefined ? { outbound: outbound.text } : {}), + ...(response.text !== undefined ? { response: response.text } : {}), + truncated: + inbound.truncated || + outbound.truncated || + response.truncated || + trace.responseTruncated, + }); + } + return { ...(meta.stored ? { traceId: requestId } : {}), trace: meta }; + } catch { + return {}; + } finally { + // Release buffered bodies; the capture is single-use. + trace.inbound = undefined; + trace.outbound = undefined; + trace.response = ""; + } +} diff --git a/src/trace/settings.ts b/src/trace/settings.ts new file mode 100644 index 00000000..7474b297 --- /dev/null +++ b/src/trace/settings.ts @@ -0,0 +1,110 @@ +/** + * Trace capture settings. + * + * Modes: `off` (default, nothing read or stored), `metadata` (hashes/counts on the usage row only), + * `redacted` (bodies stored with secrets stripped), `full` (bodies stored verbatim, opt-in). + * Env: OCX_TRACE, OCX_TRACE_TTL_HOURS, OCX_TRACE_MAX_BODY_BYTES, OCX_TRACE_MAX_DB_MB, OCX_TRACE_SAMPLE. + */ + +import type { TraceMode } from "./types"; + +export type { TraceMode } from "./types"; + +export const TRACE_MODES: readonly TraceMode[] = [ + "off", + "metadata", + "redacted", + "full", +]; + +export const TRACE_ENV = { + mode: "OCX_TRACE", + ttlHours: "OCX_TRACE_TTL_HOURS", + maxBodyBytes: "OCX_TRACE_MAX_BODY_BYTES", + maxDbMb: "OCX_TRACE_MAX_DB_MB", + sample: "OCX_TRACE_SAMPLE", +} as const; + +export interface TraceSettings { + mode: TraceMode; + ttlHours: number; + /** Per-body cap for stored payloads; hashes and byte counts always cover the full body. */ + maxBodyBytes: number; + maxDbMb: number; + /** Fraction (0..1) of requests whose bodies are stored in redacted/full mode. */ + sample: number; +} + +const DEFAULTS: TraceSettings = { + mode: "off", + ttlHours: 24, + maxBodyBytes: 512 * 1024, + maxDbMb: 256, + sample: 1, +}; + +let override: Partial = {}; + +/** Parse a mode case-insensitively after trimming; return undefined for unrecognized values. */ +export function parseTraceMode(value: unknown): TraceMode | undefined { + if (typeof value !== "string") return undefined; + const v = value.trim().toLowerCase(); + return (TRACE_MODES as readonly string[]).includes(v) + ? (v as TraceMode) + : undefined; +} + +function envNumber(name: string): number | undefined { + const raw = process.env[name]; + if (raw === undefined || raw.trim() === "") return undefined; + const n = Number(raw); + return Number.isFinite(n) ? n : undefined; +} + +function clamp(n: number, min: number, max: number): number { + return Math.min(max, Math.max(min, n)); +} + +/** + * Resolve process overrides, then environment values, then defaults on each call. + * Clamp retention to 1–720 hours, body caps to 4 KiB–8 MiB, database caps to + * 16–4096 MiB, and sampling to 0–1; round body and database caps down. + */ +export function getTraceSettings(): TraceSettings { + const mode = + override.mode ?? + parseTraceMode(process.env[TRACE_ENV.mode]) ?? + DEFAULTS.mode; + const ttl = + override.ttlHours ?? envNumber(TRACE_ENV.ttlHours) ?? DEFAULTS.ttlHours; + const body = + override.maxBodyBytes ?? + envNumber(TRACE_ENV.maxBodyBytes) ?? + DEFAULTS.maxBodyBytes; + const db = + override.maxDbMb ?? envNumber(TRACE_ENV.maxDbMb) ?? DEFAULTS.maxDbMb; + const sample = + override.sample ?? envNumber(TRACE_ENV.sample) ?? DEFAULTS.sample; + return { + mode, + ttlHours: clamp(ttl, 1, 720), + maxBodyBytes: Math.floor(clamp(body, 4 * 1024, 8 * 1024 * 1024)), + maxDbMb: Math.floor(clamp(db, 16, 4096)), + sample: clamp(sample, 0, 1), + }; +} + +/** + * Merge process-wide overrides, retaining unspecified fields, and return the + * resolved settings. Clamping occurs when settings are read. + */ +export function setTraceSettings( + partial: Partial, +): TraceSettings { + override = { ...override, ...partial }; + return getTraceSettings(); +} + +export function resetTraceSettingsForTests(): void { + override = {}; +} diff --git a/src/trace/store.ts b/src/trace/store.ts new file mode 100644 index 00000000..b7ec0860 --- /dev/null +++ b/src/trace/store.ts @@ -0,0 +1,328 @@ +/** + * Local trace store (trace.sqlite): optional request/response payloads linked to usage.jsonl rows by + * `traceId` (= requestId). Bodies are gzip BLOBs with a TTL and a total size cap. Every entry point + * swallows its own failures: trace capture must never break the proxy. + */ + +import { Database } from "bun:sqlite"; +import { chmodSync, existsSync, mkdirSync, statSync } from "node:fs"; +import { join } from "node:path"; +import { getConfigDir } from "../config"; +import { recordOwnedConfigPath } from "../lib/config-ownership"; +import { getTraceSettings } from "./settings"; +import type { TraceMode, UsageTraceMeta } from "./types"; + +export interface TraceRow { + traceId: string; + createdAt: number; + mode: TraceMode; + conversationId?: string; + provider?: string; + model?: string; + status?: number; + meta: UsageTraceMeta; + inbound?: string; + outbound?: string; + response?: string; + truncated: boolean; +} + +export interface TraceRowSummary { + traceId: string; + createdAt: number; + expiresAt: number; + mode: TraceMode; + conversationId?: string; + provider?: string; + model?: string; + status?: number; + meta: UsageTraceMeta; + truncated: boolean; +} + +const PRUNE_EVERY_WRITES = 100; +const PRUNE_EVERY_MS = 60_000; + +let db: Database | null = null; +let dbPath = ""; +let writesSincePrune = 0; +let lastPruneAt = 0; + +/** Return the trace database path in the current config directory without opening it. */ +export function traceDbPath(): string { + return join(getConfigDir(), "trace.sqlite"); +} + +/** + * Open or reuse the store for the current config directory, creating its schema + * as needed and closing any handle for a previous directory. Permission changes + * are best-effort; directory creation and SQLite errors propagate to callers. + */ +function openStore(): Database { + const path = traceDbPath(); + if (db && dbPath === path) return db; + closeTraceStore(); + const dir = getConfigDir(); + recordOwnedConfigPath(dir, path); + mkdirSync(dir, { recursive: true, mode: 0o700 }); + try { + chmodSync(dir, 0o700); + } catch { + /* best-effort */ + } + const fresh = new Database(path, { create: true }); + // auto_vacuum only takes effect on a new file, before the first table exists. + fresh.exec("PRAGMA auto_vacuum = INCREMENTAL"); + fresh.exec("PRAGMA journal_mode = WAL"); + fresh.exec("PRAGMA synchronous = NORMAL"); + fresh.exec(` + CREATE TABLE IF NOT EXISTS traces ( + trace_id TEXT PRIMARY KEY, + created_at INTEGER NOT NULL, + expires_at INTEGER NOT NULL, + mode TEXT NOT NULL, + conversation_id TEXT, + provider TEXT, + model TEXT, + status INTEGER, + meta TEXT NOT NULL, + inbound BLOB, + outbound BLOB, + response BLOB, + truncated INTEGER NOT NULL DEFAULT 0 + ); + CREATE INDEX IF NOT EXISTS traces_expires ON traces(expires_at); + CREATE INDEX IF NOT EXISTS traces_conversation ON traces(conversation_id, created_at); + `); + try { + chmodSync(path, 0o600); + } catch { + /* best-effort */ + } + db = fresh; + dbPath = path; + return fresh; +} + +/** Close the cached database and reset pruning counters, ignoring close errors. */ +export function closeTraceStore(): void { + if (db) { + try { + db.close(); + } catch { + /* ignore */ + } + } + db = null; + dbPath = ""; + writesSincePrune = 0; + lastPruneAt = 0; +} + +function pack(text: string | undefined): Uint8Array | null { + if (text === undefined) return null; + return Bun.gzipSync(new TextEncoder().encode(text)); +} + +/** + * Decode a gzip byte buffer as UTF-8, or return undefined for a non-buffer value. + * Decompression errors propagate to the caller. + */ +function unpack(blob: unknown): string | undefined { + if (!(blob instanceof Uint8Array)) return undefined; + // bun:sqlite may surface a SharedArrayBuffer-backed view, while gunzipSync + // requires an owned ArrayBuffer-backed Uint8Array. Copy at this boundary. + return new TextDecoder().decode(Bun.gunzipSync(Uint8Array.from(blob))); +} + +/** Sum database, WAL, and shared-memory file sizes in bytes, ignoring stat failures. */ +function dbBytes(path: string): number { + let total = 0; + for (const suffix of ["", "-wal", "-shm"]) { + try { + total += statSync(path + suffix).size; + } catch { + /* missing sidecar */ + } + } + return total; +} + +function deleteExpiredRows(store: Database, now: number): number { + return store + .query("DELETE FROM traces WHERE expires_at <= ?") + .run(now).changes; +} + +function compactAfterDelete(store: Database): void { + store.exec("PRAGMA wal_checkpoint(TRUNCATE)"); + store.exec("PRAGMA incremental_vacuum"); +} + +/** + * Delete rows expiring at or before `now` (Unix milliseconds), then evict oldest + * rows toward the configured database size cap, including WAL and shared memory. + * Eviction is bounded to 200 batches of 50 rows, so the cap may remain exceeded. + * May create the store. Return rows deleted, or 0 on error even if some deletions + * already succeeded. + */ +export function pruneTraces(now = Date.now()): number { + try { + const store = openStore(); + const settings = getTraceSettings(); + let deleted = deleteExpiredRows(store, now); + const capBytes = settings.maxDbMb * 1024 * 1024; + const path = traceDbPath(); + let guard = 0; + while (dbBytes(path) > capBytes && guard < 200) { + guard += 1; + const changes = store + .query( + "DELETE FROM traces WHERE trace_id IN (SELECT trace_id FROM traces ORDER BY created_at ASC LIMIT 50)", + ) + .run().changes; + if (changes === 0) break; + deleted += changes; + compactAfterDelete(store); + } + if (deleted > 0) { + compactAfterDelete(store); + } + lastPruneAt = now; + writesSincePrune = 0; + return deleted; + } catch { + return 0; + } +} + +/** + * Insert or replace a trace by ID, gzip-compressing supplied bodies as-is. + * Retention starts at write time, independently of `row.createdAt` (Unix + * milliseconds). May create the store and prune expired/oldest rows. The caller + * controls sampling, redaction, and body limits. Return false on write failure; + * pruning failures are swallowed and do not change a successful result. + */ +export function writeTrace(row: TraceRow): boolean { + try { + const store = openStore(); + const now = Date.now(); + const settings = getTraceSettings(); + store + .query( + ` + INSERT OR REPLACE INTO traces + (trace_id, created_at, expires_at, mode, conversation_id, provider, model, status, meta, inbound, outbound, response, truncated) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + `, + ) + .run( + row.traceId, + row.createdAt, + now + settings.ttlHours * 3_600_000, + row.mode, + row.conversationId ?? null, + row.provider ?? null, + row.model ?? null, + row.status ?? null, + JSON.stringify(row.meta), + pack(row.inbound), + pack(row.outbound), + pack(row.response), + row.truncated ? 1 : 0, + ); + writesSincePrune += 1; + if ( + writesSincePrune >= PRUNE_EVERY_WRITES || + now - lastPruneAt >= PRUNE_EVERY_MS + ) + pruneTraces(now); + return true; + } catch { + return false; + } +} + +function summaryFromRow(r: Record): TraceRowSummary { + return { + traceId: String(r.trace_id), + createdAt: Number(r.created_at), + expiresAt: Number(r.expires_at), + mode: r.mode as TraceMode, + ...(typeof r.conversation_id === "string" + ? { conversationId: r.conversation_id } + : {}), + ...(typeof r.provider === "string" ? { provider: r.provider } : {}), + ...(typeof r.model === "string" ? { model: r.model } : {}), + ...(typeof r.status === "number" ? { status: r.status } : {}), + meta: JSON.parse(String(r.meta)) as UsageTraceMeta, + truncated: Number(r.truncated) === 1, + }; +} + +/** + * Read an unexpired trace with decompressed bodies, deleting expired rows first. + * `now` is Unix milliseconds; rows expiring at that instant are expired. Return + * null for a missing database/trace or any read, cleanup, or decoding failure. + */ +export function readTrace( + traceId: string, + now = Date.now(), +): + | (TraceRowSummary & Pick) + | null { + try { + if (!existsSync(traceDbPath())) return null; + const store = openStore(); + if (deleteExpiredRows(store, now) > 0) compactAfterDelete(store); + const r = store + .query("SELECT * FROM traces WHERE trace_id = ? AND expires_at > ?") + .get(traceId, now) as Record | null; + if (!r) return null; + const inbound = unpack(r.inbound); + const outbound = unpack(r.outbound); + const response = unpack(r.response); + return { + ...summaryFromRow(r), + ...(inbound !== undefined ? { inbound } : {}), + ...(outbound !== undefined ? { outbound } : {}), + ...(response !== undefined ? { response } : {}), + }; + } catch { + return null; + } +} + +/** + * List unexpired summaries newest first, optionally matching a nonempty + * `conversationId`. Default to 50 results and clamp `limit` to 1–500. + * `now` defaults to the current Unix time in milliseconds; expired rows are deleted + * first. Return [] for a missing database or any cleanup, query, or decoding error. + */ +export function listTraces( + options: { conversationId?: string; limit?: number; now?: number } = {}, +): TraceRowSummary[] { + try { + if (!existsSync(traceDbPath())) return []; + const limit = Math.min(Math.max(options.limit ?? 50, 1), 500); + const now = options.now ?? Date.now(); + const store = openStore(); + if (deleteExpiredRows(store, now) > 0) compactAfterDelete(store); + const rows = ( + options.conversationId + ? store + .query( + "SELECT * FROM traces WHERE expires_at > ? AND conversation_id = ? ORDER BY created_at DESC LIMIT ?", + ) + .all(now, options.conversationId, limit) + : store + .query( + "SELECT * FROM traces WHERE expires_at > ? ORDER BY created_at DESC LIMIT ?", + ) + .all(now, limit) + ) as Record[]; + return rows.map(summaryFromRow); + } catch { + return []; + } +} diff --git a/src/trace/types.ts b/src/trace/types.ts new file mode 100644 index 00000000..8c747285 --- /dev/null +++ b/src/trace/types.ts @@ -0,0 +1,81 @@ +/** Trace metadata types shared by the usage log and the trace store. No runtime imports. */ + +export type TraceMode = "off" | "metadata" | "redacted" | "full"; + +/** Compact per-request trace summary persisted on the usage.jsonl row (no payload content). */ +export interface UsageTraceMeta { + mode: TraceMode; + /** True when request/response bodies were written to trace.sqlite under `traceId`. */ + stored: boolean; + requestBytes?: number; + outboundBytes?: number; + responseBytes?: number; + messageCount?: number; + toolCallCount?: number; + toolDefCount?: number; + attachmentCount?: number; + /** Number of times a provider wire body was sent (retries/continuations). */ + outboundCount?: number; + requestHash?: string; + outboundHash?: string; + responseHash?: string; + /** Section hashes of the final wire body: which part changed between two turns explains a cache miss. */ + systemHash?: string; + toolsHash?: string; + /** Hash of the wire body minus its last message/input item: the cacheable prefix. */ + prefixHash?: string; +} + +const HASH_RE = /^[0-9a-f]{8,64}$/; +const MODES = new Set(["off", "metadata", "redacted", "full"]); + +function nonNegInt(value: unknown): number | undefined { + return typeof value === "number" && Number.isFinite(value) && value >= 0 + ? Math.floor(value) + : undefined; +} + +function hash(value: unknown): string | undefined { + return typeof value === "string" && HASH_RE.test(value) ? value : undefined; +} + +/** Whitelist-normalize persisted trace metadata; unknown or malformed fields are dropped. */ +export function normalizeUsageTraceMeta( + raw: unknown, +): UsageTraceMeta | undefined { + if (!raw || typeof raw !== "object" || Array.isArray(raw)) return undefined; + const r = raw as Record; + if (typeof r.mode !== "string" || !MODES.has(r.mode as TraceMode)) + return undefined; + const out: UsageTraceMeta = { + mode: r.mode as TraceMode, + stored: r.stored === true, + }; + const ints = [ + "requestBytes", + "outboundBytes", + "responseBytes", + "messageCount", + "toolCallCount", + "toolDefCount", + "attachmentCount", + "outboundCount", + ] as const; + for (const key of ints) { + const v = nonNegInt(r[key]); + if (v !== undefined) out[key] = v; + } + const hashes = [ + "requestHash", + "outboundHash", + "responseHash", + "systemHash", + "toolsHash", + "prefixHash", + ] as const; + for (const key of hashes) { + const v = hash(r[key]); + if (v !== undefined) out[key] = v; + } + return out; +} diff --git a/src/usage/log.ts b/src/usage/log.ts index 57e720be..b80a9f0a 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -8,6 +8,7 @@ import { normalizePromptCacheRequestObservation, type PromptCacheRequestObservation, } from "../prompt-cache/observability"; +import { normalizeUsageTraceMeta, type UsageTraceMeta } from "../trace/types"; export type UsageStatus = "reported" | "unreported" | "unsupported" | "estimated"; @@ -104,6 +105,10 @@ export interface PersistedUsageEntry { closeReason?: "terminal" | "client_cancel" | "non_stream" | "body_stall" | "body_overflow"; /** Already redacted + capped at capture (request-log.ts redactSecretString().slice(0,500)). */ upstreamError?: string; + /** Key into trace.sqlite; present only when this request's bodies were stored there. */ + traceId?: string; + /** Payload-free trace summary (sizes, counts, hashes); see src/trace. */ + trace?: UsageTraceMeta; } const KNOWN_USAGE_SURFACES = new Set>([ @@ -373,6 +378,10 @@ function normalizeUsageEntry(entry: PersistedUsageEntry): PersistedUsageEntry { ...(entry.terminalStatus ? { terminalStatus: entry.terminalStatus } : {}), ...(entry.closeReason ? { closeReason: entry.closeReason } : {}), ...(entry.upstreamError ? { upstreamError: entry.upstreamError } : {}), + ...(typeof entry.traceId === "string" && entry.traceId + ? { traceId: capMetadataString(entry.traceId) } + : {}), + ...(normalizeUsageTraceMeta(entry.trace) ? { trace: normalizeUsageTraceMeta(entry.trace) } : {}), }; } diff --git a/tests/request-log.test.ts b/tests/request-log.test.ts index 04d562e7..77369d5d 100644 --- a/tests/request-log.test.ts +++ b/tests/request-log.test.ts @@ -1441,6 +1441,13 @@ describe("request log restart hydrate", () => { terminalStatus: "failed", closeReason: "terminal", upstreamError: "socket connection was closed unexpectedly", + traceId: "ocx-revive", + trace: { + mode: "redacted", + stored: true, + requestBytes: 42, + requestHash: "a".repeat(32), + }, }; expect(requestLogEntryFromPersistedUsage(persisted)).toEqual({ requestId: "ocx-revive", @@ -1464,6 +1471,13 @@ describe("request log restart hydrate", () => { terminalStatus: "failed", closeReason: "terminal", upstreamError: "socket connection was closed unexpectedly", + traceId: "ocx-revive", + trace: { + mode: "redacted", + stored: true, + requestBytes: 42, + requestHash: "a".repeat(32), + }, }); }); diff --git a/tests/trace.test.ts b/tests/trace.test.ts new file mode 100644 index 00000000..dcf4ff69 --- /dev/null +++ b/tests/trace.test.ts @@ -0,0 +1,381 @@ +import { afterEach, beforeEach, describe, expect, test } from "bun:test"; +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { + appendTraceResponse, + beginTrace, + createTraceCapture, + finalizeTrace, + noteOutboundRequestBody, + runWithTrace, + sectionHashes, + shapeOf, + type TraceCapture, +} from "../src/trace/capture"; +import { + getTraceSettings, + resetTraceSettingsForTests, + setTraceSettings, +} from "../src/trace/settings"; +import { + closeTraceStore, + listTraces, + pruneTraces, + readTrace, + writeTrace, +} from "../src/trace/store"; +import { normalizeUsageTraceMeta } from "../src/trace/types"; +import { + appendUsageEntry, + readUsageEntries, + resetUsageReadCacheForTests, +} from "../src/usage/log"; + +let testDir = ""; +let previousHome: string | undefined; +const traceEnvNames = [ + "OCX_TRACE", + "OCX_TRACE_TTL_HOURS", + "OCX_TRACE_MAX_BODY_BYTES", + "OCX_TRACE_MAX_DB_MB", + "OCX_TRACE_SAMPLE", +] as const; +let previousTraceEnv: Partial> = {}; + +beforeEach(() => { + previousHome = process.env.OPENCODEX_HOME; + previousTraceEnv = {}; + for (const name of traceEnvNames) { + const value = process.env[name]; + if (value !== undefined) previousTraceEnv[name] = value; + delete process.env[name]; + } + testDir = mkdtempSync(join(tmpdir(), "ocx-trace-")); + process.env.OPENCODEX_HOME = testDir; + resetUsageReadCacheForTests(); + resetTraceSettingsForTests(); +}); + +afterEach(() => { + closeTraceStore(); + resetTraceSettingsForTests(); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + for (const name of traceEnvNames) { + const value = previousTraceEnv[name]; + if (value === undefined) delete process.env[name]; + else process.env[name] = value; + } + if (testDir) rmSync(testDir, { recursive: true, force: true }); +}); + +const body = (extra: Record = {}, last = "hello") => + JSON.stringify({ + model: "gpt-x", + instructions: "be brief", + tools: [ + { type: "function", name: "a" }, + { type: "function", name: "b" }, + ], + input: [ + { + role: "user", + content: [ + { type: "input_text", text: "first" }, + { type: "input_image", image_url: "data:x" }, + ], + }, + { type: "function_call", name: "a", arguments: "{}" }, + { role: "user", content: last }, + ], + ...extra, + }); + +describe("trace settings", () => { + test("defaults to off and clamps env values", () => { + expect(getTraceSettings().mode).toBe("off"); + process.env.OCX_TRACE = "REDACTED"; + process.env.OCX_TRACE_TTL_HOURS = "99999"; + process.env.OCX_TRACE_SAMPLE = "7"; + const s = getTraceSettings(); + expect(s.mode).toBe("redacted"); + expect(s.ttlHours).toBe(720); + expect(s.sample).toBe(1); + }); + + test("invalid mode falls back to off", () => { + process.env.OCX_TRACE = "everything"; + expect(getTraceSettings().mode).toBe("off"); + }); + + test("off creates no capture", () => { + expect(createTraceCapture()).toBeUndefined(); + }); +}); + +describe("body shape and section hashes", () => { + test("counts messages, tool calls, tool defs and attachments", () => { + const shape = shapeOf(JSON.parse(body())); + expect(shape.messageCount).toBe(3); + expect(shape.toolCallCount).toBe(1); + expect(shape.toolDefCount).toBe(2); + expect(shape.attachmentCount).toBe(1); + }); + + test("prefixHash survives a new last turn; systemHash/toolsHash localize the change", () => { + const a = sectionHashes(JSON.parse(body({}, "one"))); + const b = sectionHashes(JSON.parse(body({}, "two"))); + expect(a.prefixHash).toBe(b.prefixHash); + const c = sectionHashes( + JSON.parse(body({ instructions: "be verbose" }, "one")), + ); + expect(c.systemHash).not.toBe(a.systemHash); + expect(c.toolsHash).toBe(a.toolsHash); + expect(c.prefixHash).toBe(a.prefixHash); + const d = sectionHashes( + JSON.parse(body({ tools: [{ type: "function", name: "z" }] }, "one")), + ); + expect(d.toolsHash).not.toBe(a.toolsHash); + }); +}); + +describe("trace store", () => { + test("round-trips gzip bodies and lists by conversation", () => { + const meta = { mode: "full" as const, stored: true }; + expect( + writeTrace({ + traceId: "r1", + createdAt: 1, + mode: "full", + conversationId: "c1", + meta, + inbound: "in", + outbound: "out", + response: "res", + truncated: false, + }), + ).toBe(true); + const row = readTrace("r1"); + expect(row?.inbound).toBe("in"); + expect(row?.outbound).toBe("out"); + expect(row?.response).toBe("res"); + expect(listTraces({ conversationId: "c1" }).map((r) => r.traceId)).toEqual([ + "r1", + ]); + expect(listTraces({ conversationId: "other" })).toEqual([]); + }); + + test("expired rows are invisible and physically pruned on read activity", () => { + setTraceSettings({ ttlHours: 1 }); + writeTrace({ + traceId: "old", + createdAt: 1, + mode: "full", + meta: { mode: "full", stored: true }, + inbound: "x", + truncated: false, + }); + const later = Date.now() + 2 * 3_600_000; + expect(readTrace("old", later)).toBeNull(); + expect(pruneTraces(later)).toBe(0); + }); +}); + +describe("finalizeTrace", () => { + function capture(mode: "metadata" | "redacted" | "full"): TraceCapture { + setTraceSettings({ mode }); + const trace = createTraceCapture(); + if (!trace) throw new Error("capture expected"); + return trace; + } + + test("metadata mode stores no bodies but fills hashes and counts", () => { + const trace = capture("metadata"); + trace.inbound = body(); + trace.inboundBytes = Buffer.byteLength(trace.inbound); + const out = finalizeTrace("m1", { trace }, { timestamp: 1 }); + expect(out.traceId).toBeUndefined(); + expect(out.trace?.stored).toBe(false); + expect(out.trace?.messageCount).toBe(3); + expect(out.trace?.requestHash).toMatch(/^[0-9a-f]{32}$/); + expect(readTrace("m1")).toBeNull(); + }); + + test("redacted mode strips secrets from stored bodies but hashes the raw body", () => { + const trace = capture("redacted"); + const secret = "sk-abcdefghijklmnop"; + trace.inbound = body({ + api_key: "k-123456789", + note: `use ${secret} please`, + }); + trace.inboundBytes = Buffer.byteLength(trace.inbound); + const raw = trace.inbound; + appendTraceResponse(trace, `data: {"token":"tok_live_abcdef123456"}`); + const out = finalizeTrace( + "r1", + { trace }, + { + timestamp: 5, + conversationId: "c9", + provider: "p", + model: "m", + status: 200, + }, + ); + expect(out.traceId).toBe("r1"); + expect(out.trace?.stored).toBe(true); + const row = readTrace("r1"); + expect(row?.inbound).not.toContain(secret); + expect(row?.inbound).not.toContain("k-123456789"); + expect(row?.inbound).toContain("[REDACTED]"); + expect(row?.response).not.toContain("tok_live_abcdef123456"); + expect(row?.conversationId).toBe("c9"); + expect(out.trace?.requestHash).toBeDefined(); + expect(out.trace?.requestHash).toBe(finalizeHashOf(raw)); + }); + + test("redacted mode preserves structured redaction across response chunks", () => { + const trace = capture("redacted"); + appendTraceResponse(trace, JSON.stringify({ password: "plain-private-value" })); + appendTraceResponse(trace, JSON.stringify({ ok: true })); + const out = finalizeTrace("r2", { trace }, { timestamp: 1 }); + expect(out.trace?.stored).toBe(true); + const row = readTrace("r2"); + expect(row?.response).not.toContain("plain-private-value"); + expect(row?.response).toContain("[REDACTED]"); + }); + + test("full mode keeps bodies verbatim and caps stored UTF-8 bytes", () => { + setTraceSettings({ mode: "full", maxBodyBytes: 4096 }); + const trace = createTraceCapture()!; + trace.inbound = "😀".repeat(3000); + trace.inboundBytes = Buffer.byteLength(trace.inbound); + const out = finalizeTrace("f1", { trace }, { timestamp: 1 }); + expect(out.trace?.requestBytes).toBe(12_000); + const row = readTrace("f1"); + expect(Buffer.byteLength(row?.inbound ?? "", "utf8")).toBeLessThanOrEqual(4096); + expect(row?.inbound).not.toContain("�"); + expect(row?.truncated).toBe(true); + }); + + test("response capture counts UTF-8 bytes and separators against the cap", () => { + setTraceSettings({ mode: "full", maxBodyBytes: 4096 }); + const trace = createTraceCapture()!; + appendTraceResponse(trace, "😀".repeat(1500)); + appendTraceResponse(trace, "界".repeat(1500)); + const out = finalizeTrace("f2", { trace }, { timestamp: 1 }); + expect(out.trace?.responseBytes).toBe(10_500); + const row = readTrace("f2"); + expect(Buffer.byteLength(row?.response ?? "", "utf8")).toBeLessThanOrEqual(4096); + expect(row?.response).not.toContain("�"); + expect(row?.truncated).toBe(true); + }); + + test("is idempotent", () => { + const trace = capture("full"); + trace.inbound = "{}"; + const ctx = { trace }; + expect(finalizeTrace("i1", ctx, { timestamp: 1 }).trace).toBeDefined(); + expect(finalizeTrace("i1", ctx, { timestamp: 1 })).toEqual({}); + }); +}); + +function finalizeHashOf(text: string): string { + return new Bun.CryptoHasher("sha256").update(text).digest("hex").slice(0, 32); +} + +describe("outbound capture", () => { + test("records the last wire body inside runWithTrace and ignores calls outside it", () => { + setTraceSettings({ mode: "metadata" }); + const logCtx: { trace?: TraceCapture } = { trace: createTraceCapture() }; + noteOutboundRequestBody("outside"); + runWithTrace(logCtx, () => { + noteOutboundRequestBody("first"); + noteOutboundRequestBody(body()); + }); + expect(logCtx.trace?.outboundCount).toBe(2); + const out = finalizeTrace("o1", logCtx, { timestamp: 1 }); + expect(out.trace?.outboundCount).toBe(2); + expect(out.trace?.toolsHash).toBeDefined(); + expect(out.trace?.outboundHash).toBe(finalizeHashOf(body())); + }); +}); + +describe("beginTrace", () => { + test("reads the inbound body from a clone and leaves the request readable", async () => { + setTraceSettings({ mode: "metadata" }); + const req = new Request("http://localhost/v1/responses", { + method: "POST", + body: body(), + }); + const logCtx: { trace?: TraceCapture } = {}; + await beginTrace(logCtx, req); + expect(logCtx.trace?.inboundBytes).toBe(Buffer.byteLength(body())); + expect(await req.text()).toBe(body()); + }); + + test("stops reading a chunked clone once the inbound byte limit is exceeded", async () => { + setTraceSettings({ mode: "metadata" }); + const chunk = new Uint8Array(1024 * 1024); + const totalChunks = 20; + let pulls = 0; + const stream = new ReadableStream({ + pull(controller) { + pulls += 1; + controller.enqueue(chunk); + if (pulls >= totalChunks) controller.close(); + }, + }); + const req = new Request("http://localhost/v1/responses", { + method: "POST", + body: stream, + }); + const logCtx: { trace?: TraceCapture } = {}; + await beginTrace(logCtx, req); + expect(logCtx.trace).toBeDefined(); + expect(logCtx.trace?.inbound).toBeUndefined(); + expect(pulls).toBeLessThan(totalChunks); + }); + + test("does nothing when off", async () => { + const logCtx: { trace?: TraceCapture } = {}; + await beginTrace( + logCtx, + new Request("http://localhost/x", { method: "POST", body: "{}" }), + ); + expect(logCtx.trace).toBeUndefined(); + }); +}); + +describe("usage.jsonl integration", () => { + test("persists traceId and normalized trace meta, drops junk", () => { + appendUsageEntry({ + requestId: "u1", + timestamp: 1, + provider: "p", + model: "m", + status: 200, + durationMs: 1, + usageStatus: "reported", + traceId: "u1", + trace: { + mode: "redacted", + stored: true, + requestBytes: 10, + requestHash: "a".repeat(32), + prefixHash: "nothex!", + } as never, + }); + const [entry] = readUsageEntries(); + expect(entry?.traceId).toBe("u1"); + expect(entry?.trace?.mode).toBe("redacted"); + expect(entry?.trace?.requestBytes).toBe(10); + expect(entry?.trace?.prefixHash).toBeUndefined(); + }); + + test("normalizeUsageTraceMeta rejects unknown modes", () => { + expect(normalizeUsageTraceMeta({ mode: "bogus" })).toBeUndefined(); + expect(normalizeUsageTraceMeta(null)).toBeUndefined(); + }); +});