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
26 changes: 26 additions & 0 deletions docs-site/src/content/docs/reference/cli.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand Down
24 changes: 17 additions & 7 deletions src/server/index.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { markActivity } from "../lib/sidecar-tracker";
import { beginTrace, runWithTrace } from "../trace/capture";
import { initServerSentry } from "../telemetry/sentry-server";
import {
buildWarmupCompletionFrames,
Expand Down Expand Up @@ -350,6 +351,12 @@ function attachLiveSidebandUpstream(ws: ServerWebSocket<WsData>): 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(
Expand Down Expand Up @@ -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);
Comment thread
ChefGroep marked this conversation as resolved.
let logged = false;
const finalizeNativePassthroughLog = (
status: number,
Expand All @@ -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,
Expand All @@ -1233,7 +1241,7 @@ export function startServer(port?: number) {
});
},
},
);
));
responsesCacheProbe?.store(response);
return withCors(
responseWithDeferredRequestLog(response, requestId, start, logCtx),
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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);
}
Expand Down
37 changes: 37 additions & 0 deletions src/server/request-log.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@
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 {
Expand Down Expand Up @@ -109,6 +111,8 @@
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 {
Expand Down Expand Up @@ -160,6 +164,10 @@
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[] = [];
Expand Down Expand Up @@ -316,6 +324,8 @@
...(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 } : {}),
};
}

Expand Down Expand Up @@ -432,6 +442,8 @@
...(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 */
Expand Down Expand Up @@ -672,19 +684,31 @@
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));
} catch {
/* 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();
Expand All @@ -695,6 +719,7 @@
/* 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";
Expand Down Expand Up @@ -802,122 +827,132 @@
export function httpStatusForTerminalStatus(status: ResponsesTerminalStatus): number {
return status === "completed" ? 200 : 502;
}

export function httpStatusForRequestLogTerminal(
status: ResponsesTerminalStatus,
logCtx?: RequestLogContext,
): number {
/**
* [Decision Log]
* - 목적과 의도: Keep request logs aligned with the successful HTTP/SSE contract.
* - 기존 구현 및 제약 조건: All incomplete terminals were recorded as 502 even when the
* client-requested output limit was reached normally.
* - 검토한 주요 대안: Treat every incomplete as success, or infer the reason from display text.
* - 선택한 방식: Only structured max_output_tokens incompletes map to 200.
* - 다른 대안 대신 이 방식을 선택한 이유: Stall, EOF, and unknown incompletes must remain
* visible failures, and display text is not a stable classification contract.
* - 장점, 단점 및 영향: Logs stop reporting false upstream errors while retaining the
* incomplete terminal detail; native callers without a structured reason keep old behavior.
*/
if (status === "incomplete" && logCtx?.terminalIncompleteReason === "max_output_tokens") {
return 200;
}
if (status === "failed" && logCtx?.terminalHttpStatus !== undefined) {
return logCtx.terminalHttpStatus;
}
return httpStatusForTerminalStatus(status);
}

/**
* 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
* @param status - The upstream response status
* @param meta - Optional terminal status and close reason
* @param addLog - Function used to store the finalized log entry
*/
export function addFinalRequestLog(
requestId: string,
start: number,
logCtx: RequestLogContext,
status: number,
meta?: Pick<RequestLogEntry, "terminalStatus" | "closeReason">,
addLog: (entry: RequestLogEntry) => void = addRequestLog,
): void {
try {
// Mid-stream web-search aborts used to emit response.failed and land as 502/upstream_server_error.
// Prefer the client-close classification whenever the captured reason says so.
const effectiveStatus = status >= 500 && logCtx.upstreamError && isClientClosedMessage(logCtx.upstreamError)
? 499
: status;
const errorCode = requestLogErrorCode(effectiveStatus, logCtx.upstreamError);
// A response.failed whose classified status is 499 is still a client cancel, not an upstream
// terminal failure — keep /api/logs closeReason aligned with that.
const closeReason = effectiveStatus === 499
? "client_cancel"
: meta?.closeReason;
if (logCtx.activeAttempt) {
finishRequestAttempt(
logCtx.activeAttempt,
effectiveStatus,
Date.now() - (logCtx.activeAttemptStartedAt ?? start),
logCtx.usage,
);
}
const existing = finalizedUsage(
logCtx.providerAdapter ?? logCtx.provider,
logCtx.usage,
logCtx.usageLogInputTokens,
);
const isCombo = logCtx.comboId !== undefined && (logCtx.attempts?.length ?? 0) > 0;
const provider = isCombo ? "combo" : resolveRequestLogProvider(logCtx);
const model = isCombo ? (logCtx.requestedModel ?? resolveRequestLogModel(logCtx)) : resolveRequestLogModel(logCtx);
const account = providerAccountLabel(provider);
const winningProvider = isCombo ? (logCtx.provider ?? provider) : provider;
const providerAccountId = resolveUsageProviderAccountId(logCtx, winningProvider);
if (providerAccountId) {
for (const attempt of logCtx.attempts ?? []) {
if (attempt.providerAccountId) continue;
// Combo hops can be different providers; never inherit the winner's account id cross-provider.
if (isCombo && attempt.provider !== winningProvider) continue;
attempt.providerAccountId = providerAccountId;
}
if (
logCtx.activeAttempt
&& !logCtx.activeAttempt.providerAccountId
&& (!isCombo || logCtx.activeAttempt.provider === winningProvider)
) {
logCtx.activeAttempt.providerAccountId = providerAccountId;
}
}
const attempts = logCtx.attempts?.map(attempt => ({
...attempt,
recoveryKinds: [...attempt.recoveryKinds],
...(attempt.usage ? { usage: { ...attempt.usage } } : {}),
}));
const adapter = isCombo
? attempts?.at(-1)?.adapter ?? logCtx.providerAdapter
: logCtx.providerAdapter;
const aggregate = isCombo ? aggregateAttemptUsage(attempts ?? []) : null;
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,
model,
provider,
...(adapter ? { adapter } : {}),
...(logCtx.promptCache ? { promptCache: { ...logCtx.promptCache } } : {}),
...(account ? { account } : {}),
...(providerAccountId ? { providerAccountId } : {}),
...(logCtx.surface ? { surface: logCtx.surface } : {}),
...(logCtx.conversationId ? { conversationId: logCtx.conversationId } : {}),
...(logCtx.requestedModel ? { requestedModel: logCtx.requestedModel } : {}),
...(logCtx.requestedEffort ? { requestedEffort: logCtx.requestedEffort } : {}),
...(logCtx.effectiveEffort ? { effectiveEffort: logCtx.effectiveEffort } : {}),

Check warning on line 955 in src/server/request-log.ts

View check run for this annotation

codefactor.io / CodeFactor

src/server/request-log.ts#L830-L955

Very Complex Method
...(logCtx.reasoningWireField ? { reasoningWireField: logCtx.reasoningWireField } : {}),
...(logCtx.reasoningWireValue !== undefined ? { reasoningWireValue: logCtx.reasoningWireValue } : {}),
...(logCtx.requestedServiceTier ? { requestedServiceTier: logCtx.requestedServiceTier } : {}),
Expand All @@ -942,6 +977,8 @@
...(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({
Expand Down
10 changes: 10 additions & 0 deletions src/server/responses/fetch-helpers.ts
Original file line number Diff line number Diff line change
@@ -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<Server<WsData>, "timeout"> | undefined): boolean {
Expand Down Expand Up @@ -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<RequestInit, "signal">,
Expand All @@ -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")) {
Expand Down
Loading
Loading