From 84f105ba56240f31ca57f8e348bcc7290c6cc7bd Mon Sep 17 00:00:00 2001 From: chefgroep Date: Sun, 4 Oct 2026 23:30:57 +0200 Subject: [PATCH] feat(trace): cover compact and cache hits --- docs-site/src/content/docs/reference/cli.md | 7 +++- src/cache/kv-cache.ts | 12 +++++- src/cache/response-cache-middleware.ts | 24 +++++++++-- src/server/index.ts | 45 ++++++++++++++------- src/trace/capture.ts | 33 +++++++++++++++ src/trace/types.ts | 11 +++++ tests/kv-cache.test.ts | 12 +++++- tests/response-cache-e2e.test.ts | 8 ++++ tests/response-cache-middleware.test.ts | 11 ++++- tests/trace.test.ts | 39 ++++++++++++++++++ 10 files changed, 178 insertions(+), 24 deletions(-) diff --git a/docs-site/src/content/docs/reference/cli.md b/docs-site/src/content/docs/reference/cli.md index 4abfab58..8c05648a 100644 --- a/docs-site/src/content/docs/reference/cli.md +++ b/docs-site/src/content/docs/reference/cli.md @@ -534,8 +534,11 @@ Tracing is **off by default** and is configured by environment variable at proxy 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`. +Outbound capture covers adapters that send through the shared upstream fetch helper; inbound and response +capture covers `/v1/responses`, `/v1/messages`, `/v1/chat/completions`, and `/v1/responses/compact`. +Proxy response-cache hits are traced without duplicating the cached response body: the usage trace records +`cacheHit`, a response hash/byte count, and the request id that originally populated the cache when known. +WebSocket/live/realtime traffic is not yet covered. ## Updating diff --git a/src/cache/kv-cache.ts b/src/cache/kv-cache.ts index 1841a297..fe131e0f 100644 --- a/src/cache/kv-cache.ts +++ b/src/cache/kv-cache.ts @@ -34,6 +34,8 @@ export interface CacheEntry { storedAt: number; /** Byte size of `body`, for capacity accounting. */ size: number; + /** Request id that produced this cached response, for trace correlation only. */ + sourceTraceId?: string; } export interface ResponseCacheOptions { @@ -180,7 +182,7 @@ export class ResponseCache { normalizedBody: string, endpoint = "responses", now = Date.now(), - ): { body: string; contentType: string } | null { + ): { body: string; contentType: string; sourceTraceId?: string } | null { if (!this.opts.enabled) return null; const key = cacheKeyFor(endpoint, provider, model, normalizedBody); const entry = this.map.get(key); @@ -198,7 +200,11 @@ export class ResponseCache { this.map.delete(key); this.map.set(key, entry); this.stats.hits += 1; - return { body: entry.body, contentType: entry.contentType }; + return { + body: entry.body, + contentType: entry.contentType, + ...(entry.sourceTraceId ? { sourceTraceId: entry.sourceTraceId } : {}), + }; } /** Store a completed response. No-op when disabled or over capacity after eviction. */ @@ -210,6 +216,7 @@ export class ResponseCache { contentType: string, endpoint = "responses", now = Date.now(), + sourceTraceId?: string, ): void { if (!this.opts.enabled) return; // Measure real UTF-8 bytes, not UTF-16 code units: `body.length` under-counts multi-byte @@ -229,6 +236,7 @@ export class ResponseCache { expiresAt: now + this.opts.ttlMs, storedAt: now, size: byteSize, + ...(sourceTraceId ? { sourceTraceId } : {}), }; this.map.delete(key); this.map.set(key, entry); diff --git a/src/cache/response-cache-middleware.ts b/src/cache/response-cache-middleware.ts index 057df834..2b0c7360 100644 --- a/src/cache/response-cache-middleware.ts +++ b/src/cache/response-cache-middleware.ts @@ -51,6 +51,12 @@ export function getInstalledCache(): ResponseCache | null { export interface CacheHit { hit: Response; + /** Rebuilt request whose body was consumed while computing the cache key. */ + request: Request; + /** Cached response body for payload-free trace hashing/counting; never exposed to clients separately. */ + responseBody: string; + /** Request id that originally populated this entry, when available. */ + sourceTraceId?: string; /** Route this hit replayed, so the server can log the served request without re-parsing the body. */ provider: string; model: string; @@ -65,7 +71,7 @@ export interface CacheMiss { normalizedBody: string; endpoint: "responses" | "messages" | "chat-completions"; /** Store a 2xx non-streaming response into the cache (no-op if not cacheable). */ - store: (response: Response) => void; + store: (response: Response, sourceTraceId?: string) => void; } export type CacheProbe = CacheHit | CacheMiss | null; @@ -173,6 +179,9 @@ export async function probeResponseCache( ); return { hit: new Response(hit.body, { status: 200, headers }), + request: rebuild(req, raw), + responseBody: hit.body, + ...(hit.sourceTraceId ? { sourceTraceId: hit.sourceTraceId } : {}), provider: route.providerName, model: route.modelId, }; @@ -180,7 +189,7 @@ export async function probeResponseCache( const rebuilt = rebuild(req, raw); - const store = (response: Response) => { + const store = (response: Response, sourceTraceId?: string) => { if (response.status < 200 || response.status >= 300) return; const ct = response.headers.get("content-type") ?? "application/json"; if (ct.includes("text/event-stream")) return; @@ -190,7 +199,16 @@ export async function probeResponseCache( .text() .then((body) => { if (!body) return; - cache.set(route.providerName, route.modelId, scopedNormalized, body, ct, endpoint); + cache.set( + route.providerName, + route.modelId, + scopedNormalized, + body, + ct, + endpoint, + Date.now(), + sourceTraceId, + ); }) .catch(() => { /* clone read failure is non-fatal */ diff --git a/src/server/index.ts b/src/server/index.ts index 0156c151..6000f5bb 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -1,5 +1,10 @@ import { markActivity } from "../lib/sidecar-tracker"; -import { beginTrace, runWithTrace } from "../trace/capture"; +import { + appendTraceResponse, + beginTrace, + noteTraceCacheHit, + runWithTrace, +} from "../trace/capture"; import { initServerSentry } from "../telemetry/sentry-server"; import { buildWarmupCompletionFrames, @@ -221,14 +226,16 @@ const LIVE_SIDEBAND_PENDING_MAX = 32; * and addFinalRequestLog omits the usage/totalTokens fields when it is absent. The response itself * already carries `x-cache: HIT`; this is purely the request-count / observability entry. */ -function logCacheHitRequest(hit: CacheHit): void { +async function logCacheHitRequest(hit: CacheHit): Promise { const start = Date.now(); - addFinalRequestLog( - nextRequestLogId(start), - start, - { model: hit.model || "unknown", provider: hit.provider || "unknown" }, - 200, - ); + const requestId = nextRequestLogId(start); + const logCtx: RequestLogContext = { + model: hit.model || "unknown", + provider: hit.provider || "unknown", + }; + await beginTrace(logCtx, hit.request); + noteTraceCacheHit(logCtx.trace, hit.responseBody, hit.sourceTraceId); + addFinalRequestLog(requestId, start, logCtx, 200); } function closeLiveSideband( @@ -969,9 +976,13 @@ export function startServer(port?: number) { model: "unknown", provider: "unknown", }; + await beginTrace(logCtx, req); let response: Response; try { - response = await handleResponsesCompact(req, config, logCtx); + response = await runWithTrace( + logCtx, + () => handleResponsesCompact(req, config, logCtx), + ); } catch { response = formatErrorResponse( 500, @@ -979,6 +990,10 @@ export function startServer(port?: number) { "Unexpected compact request failure", ); } + if (logCtx.trace) { + const tracedBody = await response.clone().text().catch(() => ""); + if (tracedBody) appendTraceResponse(logCtx.trace, tracedBody); + } addFinalRequestLog( requestId, start, @@ -1172,7 +1187,7 @@ export function startServer(port?: number) { "responses", ); if (responsesCacheProbe && "hit" in responsesCacheProbe) { - logCacheHitRequest(responsesCacheProbe); + await logCacheHitRequest(responsesCacheProbe); return withCors(responsesCacheProbe.hit, req, config); } const responsesWorkReq = responsesCacheProbe?.request ?? req; @@ -1215,7 +1230,7 @@ export function startServer(port?: number) { }, }, )); - responsesCacheProbe?.store(response); + responsesCacheProbe?.store(response, requestId); return withCors( responseWithDeferredRequestLog(response, requestId, start, logCtx), req, @@ -1312,7 +1327,7 @@ export function startServer(port?: number) { "messages", ); if (messagesCacheProbe && "hit" in messagesCacheProbe) { - logCacheHitRequest(messagesCacheProbe); + await logCacheHitRequest(messagesCacheProbe); return withCors(messagesCacheProbe.hit, req, config); } const messagesWorkReq = messagesCacheProbe?.request ?? req; @@ -1332,7 +1347,7 @@ export function startServer(port?: number) { logCtx, { requestId, start }, )); - messagesCacheProbe?.store(response); + messagesCacheProbe?.store(response, requestId); return withCors(response, req, config); } @@ -1368,7 +1383,7 @@ export function startServer(port?: number) { "chat-completions", ); if (chatCacheProbe && "hit" in chatCacheProbe) { - logCacheHitRequest(chatCacheProbe); + await logCacheHitRequest(chatCacheProbe); return withCors(chatCacheProbe.hit, req, config); } const chatWorkReq = chatCacheProbe?.request ?? req; @@ -1385,7 +1400,7 @@ export function startServer(port?: number) { logCtx, { requestId, start }, )); - chatCacheProbe?.store(response); + chatCacheProbe?.store(response, requestId); return withCors(response, req, config); } diff --git a/src/trace/capture.ts b/src/trace/capture.ts index c7f4bd25..4d0776db 100644 --- a/src/trace/capture.ts +++ b/src/trace/capture.ts @@ -34,6 +34,8 @@ export interface TraceCapture { responseStoredBytes: number; responseTruncated: boolean; responseHasher: Hash; + cacheHit: boolean; + cacheSourceTraceId?: string; done: boolean; } @@ -85,6 +87,7 @@ export function createTraceCapture(): TraceCapture | undefined { responseStoredBytes: 0, responseTruncated: false, responseHasher: createHash("sha256"), + cacheHit: false, done: false, }; } @@ -162,6 +165,32 @@ export function noteOutboundRequestBody(body: unknown): void { } } +/** + * Record a proxy response-cache hit without copying the cached response body into trace.sqlite. + * The full cached body is hashed/count-measured for correlation, while persisted payload capture + * remains limited to the hit request itself. + */ +export function noteTraceCacheHit( + trace: TraceCapture | undefined, + responseBody: string, + sourceTraceId?: string, +): void { + if (!trace || trace.done) return; + try { + trace.cacheHit = true; + if ( + typeof sourceTraceId === "string" + && /^[A-Za-z0-9._:-]{1,128}$/u.test(sourceTraceId) + ) { + trace.cacheSourceTraceId = sourceTraceId; + } + trace.responseHasher.update(responseBody); + trace.responseBytes += byteLength(responseBody); + } catch { + /* ignore */ + } +} + /** Append one response payload (SSE data block or JSON body) to the bounded response copy. */ export function appendTraceResponse( trace: TraceCapture | undefined, @@ -370,6 +399,10 @@ export function finalizeTrace( ...(trace.outboundCount > 0 ? { outboundCount: trace.outboundCount } : {}), + ...(trace.cacheHit ? { cacheHit: true } : {}), + ...(trace.cacheSourceTraceId + ? { cacheSourceTraceId: trace.cacheSourceTraceId } + : {}), ...(trace.inbound !== undefined ? { requestHash: sha(trace.inbound) } : {}), diff --git a/src/trace/types.ts b/src/trace/types.ts index 8c747285..7ec9c2b7 100644 --- a/src/trace/types.ts +++ b/src/trace/types.ts @@ -16,6 +16,10 @@ export interface UsageTraceMeta { attachmentCount?: number; /** Number of times a provider wire body was sent (retries/continuations). */ outboundCount?: number; + /** True when the proxy response cache served this request without an upstream call. */ + cacheHit?: boolean; + /** Request id that originally populated the response-cache entry, when known. */ + cacheSourceTraceId?: string; requestHash?: string; outboundHash?: string; responseHash?: string; @@ -65,6 +69,13 @@ export function normalizeUsageTraceMeta( const v = nonNegInt(r[key]); if (v !== undefined) out[key] = v; } + if (r.cacheHit === true) out.cacheHit = true; + if ( + typeof r.cacheSourceTraceId === "string" + && /^[A-Za-z0-9._:-]{1,128}$/u.test(r.cacheSourceTraceId) + ) { + out.cacheSourceTraceId = r.cacheSourceTraceId; + } const hashes = [ "requestHash", "outboundHash", diff --git a/tests/kv-cache.test.ts b/tests/kv-cache.test.ts index fde2abc6..735be559 100644 --- a/tests/kv-cache.test.ts +++ b/tests/kv-cache.test.ts @@ -184,7 +184,16 @@ describe("ResponseCache observability + guards (Fase D quality round)", () => { { enabled: true, ttlMs: 60_000, maxEntries: 8, persist: true }, dir, ); - first.set("p", "m", "req", '{"warm":true}', "application/json"); + first.set( + "p", + "m", + "req", + '{"warm":true}', + "application/json", + "responses", + Date.now(), + "ocx-source-123", + ); const second = new ResponseCache( { enabled: true, ttlMs: 60_000, maxEntries: 8, persist: true }, @@ -193,6 +202,7 @@ describe("ResponseCache observability + guards (Fase D quality round)", () => { const hit = second.get("p", "m", "req"); expect(hit).not.toBeNull(); expect(hit!.body).toBe('{"warm":true}'); + expect(hit!.sourceTraceId).toBe("ocx-source-123"); } finally { rmSync(dir, { recursive: true, force: true }); } diff --git a/tests/response-cache-e2e.test.ts b/tests/response-cache-e2e.test.ts index a670ba35..8cff2c78 100644 --- a/tests/response-cache-e2e.test.ts +++ b/tests/response-cache-e2e.test.ts @@ -13,6 +13,7 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import { saveConfig } from "../src/config"; import { startServer } from "../src/server"; +import { clearRequestLogsForTests, getRequestLogEntries } from "../src/server/request-log"; import type { OcxConfig } from "../src/types"; import { installIsolatedCodexHome, type IsolatedCodexHome } from "./helpers/isolated-codex-home"; import { managementFetch } from "./helpers/management-auth"; @@ -24,6 +25,7 @@ const originalFetch = globalThis.fetch; beforeEach(() => { previousHome = process.env.OPENCODEX_HOME; + clearRequestLogsForTests(); isolatedCodexHome = installIsolatedCodexHome("ocx-resp-cache-e2e-"); testDir = mkdtempSync(join(tmpdir(), "ocx-resp-cache-e2e-")); process.env.OPENCODEX_HOME = testDir; @@ -31,6 +33,7 @@ beforeEach(() => { }); afterEach(() => { + delete process.env.OCX_TRACE; if (previousHome === undefined) delete process.env.OPENCODEX_HOME; else process.env.OPENCODEX_HOME = previousHome; isolatedCodexHome?.restore(); @@ -78,6 +81,7 @@ function mockConfig(baseUrl: string): OcxConfig { } test("identical non-streaming request hits the cache on the second call", async () => { + process.env.OCX_TRACE = "metadata"; const upstream = mockJsonUpstream(); saveConfig(mockConfig(`${upstream.url.toString().replace(/\/$/, "")}/v1`)); const server = startServer(0); @@ -104,6 +108,10 @@ test("identical non-streaming request hits the cache on the second call", async expect(second.status).toBe(200); expect(second.headers.get("x-cache")).toBe("HIT"); expect(upstreamHits).toBe(1); // upstream was NOT called again + const hitLog = getRequestLogEntries().at(-1); + expect(hitLog?.trace?.cacheHit).toBe(true); + expect(hitLog?.trace?.cacheSourceTraceId).toMatch(/^ocx-[A-Za-z0-9-]+$/); + expect(hitLog?.trace?.cacheSourceTraceId).not.toBe(hitLog?.requestId); const replayed = await second.json() as { object: string; choices: Array<{ message: { content: string } }> }; expect(replayed.object).toBe("chat.completion"); expect(replayed.choices[0]?.message.content).toBe("pong"); diff --git a/tests/response-cache-middleware.test.ts b/tests/response-cache-middleware.test.ts index c90d8015..447beea8 100644 --- a/tests/response-cache-middleware.test.ts +++ b/tests/response-cache-middleware.test.ts @@ -183,13 +183,22 @@ describe("probeResponseCache", () => { probeA.store(new Response(JSON.stringify({ ok: true }), { status: 200, headers: { "content-type": "application/json" }, - })); + }), "ocx-source-123"); await new Promise(r => setTimeout(r, 0)); // store() resolves asynchronously } const probeB = await probeResponseCache(b, baseConfig(), "messages"); expect(probeB).not.toBeNull(); expect("hit" in probeB!).toBe(true); // stable order → same key → HIT + if (probeB && "hit" in probeB) { + expect(probeB.sourceTraceId).toBe("ocx-source-123"); + expect(probeB.responseBody).toBe(JSON.stringify({ ok: true })); + expect(await probeB.request.json()).toMatchObject({ + model: "claude-opus-4-8", + a: 2, + z: 1, + }); + } }); test("miss exposes a rebuilt request + store that captures a 2xx body", async () => { diff --git a/tests/trace.test.ts b/tests/trace.test.ts index dcf4ff69..bb77600f 100644 --- a/tests/trace.test.ts +++ b/tests/trace.test.ts @@ -8,6 +8,7 @@ import { createTraceCapture, finalizeTrace, noteOutboundRequestBody, + noteTraceCacheHit, runWithTrace, sectionHashes, shapeOf, @@ -302,6 +303,44 @@ describe("outbound capture", () => { }); }); +describe("cache-hit trace metadata", () => { + test("hashes cached response and links source trace without duplicating its body", () => { + setTraceSettings({ mode: "full", maxBodyBytes: 4096 }); + const trace = createTraceCapture(); + if (!trace) throw new Error("trace expected"); + trace.inbound = body(); + trace.inboundBytes = Buffer.byteLength(trace.inbound); + noteTraceCacheHit(trace, '{"cached":"response"}', "ocx-source-123"); + + const out = finalizeTrace("ocx-hit-456", { trace }, { timestamp: Date.now() }); + expect(out.trace).toMatchObject({ + cacheHit: true, + cacheSourceTraceId: "ocx-source-123", + responseBytes: Buffer.byteLength('{"cached":"response"}'), + }); + expect(out.trace?.responseHash).toBe( + finalizeHashOf('{"cached":"response"}'), + ); + + const row = readTrace("ocx-hit-456"); + expect(row?.inbound).toBe(body()); + expect(row?.response).toBeUndefined(); + }); + + test("drops an unsafe persisted cache source trace id", () => { + expect(normalizeUsageTraceMeta({ + mode: "metadata", + stored: false, + cacheHit: true, + cacheSourceTraceId: "bad\nterminal-id", + })).toEqual({ + mode: "metadata", + stored: false, + cacheHit: true, + }); + }); +}); + describe("beginTrace", () => { test("reads the inbound body from a clone and leaves the request readable", async () => { setTraceSettings({ mode: "metadata" });