Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions docs-site/src/content/docs/reference/cli.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Stored traces can be inspected directly from the local SQLite store; the proxy does not need to be running:

Expand Down
12 changes: 10 additions & 2 deletions src/cache/kv-cache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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);
Expand All @@ -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. */
Expand All @@ -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
Expand All @@ -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);
Expand Down
24 changes: 21 additions & 3 deletions src/cache/response-cache-middleware.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -173,14 +179,17 @@ 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,
};
}

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;
Expand All @@ -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 */
Expand Down
45 changes: 30 additions & 15 deletions src/server/index.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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<void> {
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(
Expand Down Expand Up @@ -990,16 +997,24 @@ 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,
"server_error",
"Unexpected compact request failure",
);
}
if (logCtx.trace) {
const tracedBody = await response.clone().text().catch(() => "");
if (tracedBody) appendTraceResponse(logCtx.trace, tracedBody);
}
addFinalRequestLog(
requestId,
start,
Expand Down Expand Up @@ -1193,7 +1208,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;
Expand Down Expand Up @@ -1236,7 +1251,7 @@ export function startServer(port?: number) {
},
},
));
responsesCacheProbe?.store(response);
responsesCacheProbe?.store(response, requestId);
return withCors(
responseWithDeferredRequestLog(response, requestId, start, logCtx),
req,
Expand Down Expand Up @@ -1333,7 +1348,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;
Expand All @@ -1353,7 +1368,7 @@ export function startServer(port?: number) {
logCtx,
{ requestId, start },
));
messagesCacheProbe?.store(response);
messagesCacheProbe?.store(response, requestId);
return withCors(response, req, config);
}

Expand Down Expand Up @@ -1389,7 +1404,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;
Expand All @@ -1406,7 +1421,7 @@ export function startServer(port?: number) {
logCtx,
{ requestId, start },
));
chatCacheProbe?.store(response);
chatCacheProbe?.store(response, requestId);
return withCors(response, req, config);
}

Expand Down
33 changes: 33 additions & 0 deletions src/trace/capture.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,8 @@ export interface TraceCapture {
responseStoredBytes: number;
responseTruncated: boolean;
responseHasher: Hash;
cacheHit: boolean;
cacheSourceTraceId?: string;
done: boolean;
}

Expand Down Expand Up @@ -85,6 +87,7 @@ export function createTraceCapture(): TraceCapture | undefined {
responseStoredBytes: 0,
responseTruncated: false,
responseHasher: createHash("sha256"),
cacheHit: false,
done: false,
};
}
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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) }
: {}),
Expand Down
11 changes: 11 additions & 0 deletions src/trace/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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",
Expand Down
12 changes: 11 additions & 1 deletion tests/kv-cache.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 },
Expand All @@ -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 });
}
Expand Down
8 changes: 8 additions & 0 deletions tests/response-cache-e2e.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -24,13 +25,15 @@ 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;
globalThis.fetch = originalFetch;
});

afterEach(() => {
delete process.env.OCX_TRACE;
if (previousHome === undefined) delete process.env.OPENCODEX_HOME;
else process.env.OPENCODEX_HOME = previousHome;
isolatedCodexHome?.restore();
Expand Down Expand Up @@ -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);
Expand All @@ -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");
Expand Down
Loading
Loading