From dbc37d9d2a88a34874da9544ba4e5fff15a809aa Mon Sep 17 00:00:00 2001 From: chilung Date: Sat, 22 Aug 2026 08:48:06 +0000 Subject: [PATCH 1/9] feat(usage): add durable stream timeline and failure attribution to request history (closes #1217) --- src/usage/log.ts | 92 +++++++++++++++++++++++++++++++++++++++++ tests/usage-log.test.ts | 43 +++++++++++++++++++ 2 files changed, 135 insertions(+) diff --git a/src/usage/log.ts b/src/usage/log.ts index 7e056f97b0..7cb6ca7bea 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -21,6 +21,30 @@ export type UsageStatus = "reported" | "unreported" | "unsupported" | "estimated export type UsageAccountLogLabel = "main" | `p${string}` | `o${string}`; export type CodexUsageAccountLogLabel = UsageAccountLogLabel; +/** + * Bounded stream timing breakdown in elapsed ms from request/attempt start (issue #1217). + * Best-effort correlation metrics for streaming observability. + */ +export interface StreamTimeline { + upstreamDispatchMs?: number; + upstreamHeadersMs?: number; + upstreamFirstByteMs?: number; + upstreamFirstSemanticOutputMs?: number; + downstreamFirstWriteMs?: number; + upstreamEndMs?: number; + downstreamEndMs?: number; +} + +export type FailureSide = "upstream" | "relay" | "downstream" | "client" | "local"; +export type FailureStage = + | "pre_dispatch" + | "upstream_wait_headers" + | "upstream_read" + | "relay_transform" + | "downstream_write" + | "client_cancel" + | "terminal_delivery"; + /** * Accepts EITHER label family. This is the predicate the persistence writers use, so widening * it here is what stops six separate call sites from silently dropping an `o`-label -- including @@ -94,6 +118,10 @@ export interface PersistedUsageAttempt { reasoningWireValue?: string | number | boolean; /** Adapter-produced tier fact for this physical attempt; absent on pre-B0 rows. */ tierOutcome?: AttemptTierOutcome; + /** Bounded streaming timeline for this attempt (issue #1217). */ + streamTimeline?: StreamTimeline; + failureSide?: FailureSide; + failureStage?: FailureStage; } export interface PersistedUsageEntry { @@ -148,6 +176,14 @@ 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; + /** + * Bounded streaming timeline and causal failure attribution (issue #1217). + */ + streamTimeline?: StreamTimeline; + failureSide?: FailureSide; + failureStage?: FailureStage; + transportPhase?: string; + terminalSource?: string; /** * Bounded route-decision trace (RI-01): why this provider/model/account was * selected. Additive field; old rows without it parse unchanged. Never @@ -281,6 +317,38 @@ const TIER_CONFIRMATIONS = new Set([ const FAST_DOWNGRADE_REASONS = new Set>([ "route-unsupported", "wire-unavailable", "response-declined", ]); +const KNOWN_FAILURE_SIDES = new Set([ + "upstream", "relay", "downstream", "client", "local", +]); +const KNOWN_FAILURE_STAGES = new Set([ + "pre_dispatch", + "upstream_wait_headers", + "upstream_read", + "relay_transform", + "downstream_write", + "client_cancel", + "terminal_delivery", +]); + +function normalizeStreamTimeline(raw: unknown): StreamTimeline | null { + if (!raw || typeof raw !== "object" || Array.isArray(raw)) return null; + const t = raw as Record; + const out: StreamTimeline = {}; + for (const key of [ + "upstreamDispatchMs", + "upstreamHeadersMs", + "upstreamFirstByteMs", + "upstreamFirstSemanticOutputMs", + "downstreamFirstWriteMs", + "upstreamEndMs", + "downstreamEndMs", + ] as const) { + if (key in t && isNonNegativeFiniteNumber(t[key])) { + out[key] = t[key] as number; + } + } + return Object.keys(out).length > 0 ? out : null; +} export function isLabRouteSubjectId(value: unknown): value is string { return typeof value === "string" && LAB_ROUTE_SUBJECT_ID_RE.test(value); @@ -439,6 +507,15 @@ function normalizeUsageAttempt(raw: unknown): PersistedUsageAttempt | null { : { reasoningWireValue: attempt.reasoningWireValue } : {}), ...(tierOutcome ? { tierOutcome } : {}), + ...(normalizeStreamTimeline(attempt.streamTimeline) + ? { streamTimeline: normalizeStreamTimeline(attempt.streamTimeline)! } + : {}), + ...(typeof attempt.failureSide === "string" && KNOWN_FAILURE_SIDES.has(attempt.failureSide as FailureSide) + ? { failureSide: attempt.failureSide as FailureSide } + : {}), + ...(typeof attempt.failureStage === "string" && KNOWN_FAILURE_STAGES.has(attempt.failureStage as FailureStage) + ? { failureStage: attempt.failureStage as FailureStage } + : {}), }; } @@ -551,6 +628,21 @@ function normalizeUsageEntry(entry: PersistedUsageEntry): PersistedUsageEntry { ...(entry.terminalStatus ? { terminalStatus: entry.terminalStatus } : {}), ...(entry.closeReason ? { closeReason: entry.closeReason } : {}), ...(entry.upstreamError ? { upstreamError: entry.upstreamError } : {}), + ...(normalizeStreamTimeline(entry.streamTimeline) + ? { streamTimeline: normalizeStreamTimeline(entry.streamTimeline)! } + : {}), + ...(typeof entry.failureSide === "string" && KNOWN_FAILURE_SIDES.has(entry.failureSide as FailureSide) + ? { failureSide: entry.failureSide as FailureSide } + : {}), + ...(typeof entry.failureStage === "string" && KNOWN_FAILURE_STAGES.has(entry.failureStage as FailureStage) + ? { failureStage: entry.failureStage as FailureStage } + : {}), + ...(typeof entry.transportPhase === "string" && entry.transportPhase + ? { transportPhase: capMetadataString(entry.transportPhase) } + : {}), + ...(typeof entry.terminalSource === "string" && entry.terminalSource + ? { terminalSource: capMetadataString(entry.terminalSource) } + : {}), ...(routeDecision ? { routeDecision } : {}), }; } diff --git a/tests/usage-log.test.ts b/tests/usage-log.test.ts index cc96db4921..58d9273002 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -913,4 +913,47 @@ describe("usage log", () => { expect(readRecentUsageEntries(1)).toEqual([]); }, STORE_BUDGET_MS); + + test("normalizes and preserves streamTimeline and failure attribution (#1217)", () => { + appendUsageEntry({ + requestId: "ocx-stream-timeline-test", + timestamp: Date.now(), + provider: "anthropic", + model: "claude-sonnet-5", + status: 502, + durationMs: 61342, + firstOutputMs: 9107, + usageStatus: "unreported", + streamTimeline: { + upstreamDispatchMs: 12, + upstreamHeadersMs: 4410, + upstreamFirstByteMs: 4421, + upstreamFirstSemanticOutputMs: 9107, + downstreamFirstWriteMs: 4423, + upstreamEndMs: 61340, + downstreamEndMs: 61342, + }, + failureSide: "upstream", + failureStage: "upstream_read", + transportPhase: "mid_stream", + terminalSource: "synthetic", + }); + + const entries = readRecentUsageEntries(10); + const row = entries.find(e => e.requestId === "ocx-stream-timeline-test"); + expect(row).toBeDefined(); + expect(row?.streamTimeline).toEqual({ + upstreamDispatchMs: 12, + upstreamHeadersMs: 4410, + upstreamFirstByteMs: 4421, + upstreamFirstSemanticOutputMs: 9107, + downstreamFirstWriteMs: 4423, + upstreamEndMs: 61340, + downstreamEndMs: 61342, + }); + expect(row?.failureSide).toBe("upstream"); + expect(row?.failureStage).toBe("upstream_read"); + expect(row?.transportPhase).toBe("mid_stream"); + expect(row?.terminalSource).toBe("synthetic"); + }); }); From 0058054d0dbc47c9613c3839c3d3ae4eb85187ca Mon Sep 17 00:00:00 2001 From: chilung Date: Sat, 22 Aug 2026 10:39:48 +0000 Subject: [PATCH 2/9] fix(usage): clean up stream timeline normalization (#1217) --- src/usage/log.ts | 8 ++------ 1 file changed, 2 insertions(+), 6 deletions(-) diff --git a/src/usage/log.ts b/src/usage/log.ts index 7cb6ca7bea..e5f74f3b9c 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -507,9 +507,7 @@ function normalizeUsageAttempt(raw: unknown): PersistedUsageAttempt | null { : { reasoningWireValue: attempt.reasoningWireValue } : {}), ...(tierOutcome ? { tierOutcome } : {}), - ...(normalizeStreamTimeline(attempt.streamTimeline) - ? { streamTimeline: normalizeStreamTimeline(attempt.streamTimeline)! } - : {}), + ...(normalizeStreamTimeline(attempt.streamTimeline) ? { streamTimeline: normalizeStreamTimeline(attempt.streamTimeline) as StreamTimeline } : {}), ...(typeof attempt.failureSide === "string" && KNOWN_FAILURE_SIDES.has(attempt.failureSide as FailureSide) ? { failureSide: attempt.failureSide as FailureSide } : {}), @@ -628,9 +626,7 @@ function normalizeUsageEntry(entry: PersistedUsageEntry): PersistedUsageEntry { ...(entry.terminalStatus ? { terminalStatus: entry.terminalStatus } : {}), ...(entry.closeReason ? { closeReason: entry.closeReason } : {}), ...(entry.upstreamError ? { upstreamError: entry.upstreamError } : {}), - ...(normalizeStreamTimeline(entry.streamTimeline) - ? { streamTimeline: normalizeStreamTimeline(entry.streamTimeline)! } - : {}), + ...(normalizeStreamTimeline(entry.streamTimeline) ? { streamTimeline: normalizeStreamTimeline(entry.streamTimeline) as StreamTimeline } : {}), ...(typeof entry.failureSide === "string" && KNOWN_FAILURE_SIDES.has(entry.failureSide as FailureSide) ? { failureSide: entry.failureSide as FailureSide } : {}), From 1459a5cec2f9eff110420b7592a392409a601a3a Mon Sep 17 00:00:00 2001 From: chilung Date: Sat, 22 Aug 2026 11:44:08 +0000 Subject: [PATCH 3/9] fix(usage): enforce strict validation on transportPhase and terminalSource (#1217) --- src/usage/log.ts | 31 ++++++++++++++++++++++------ tests/usage-log.test.ts | 45 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 70 insertions(+), 6 deletions(-) diff --git a/src/usage/log.ts b/src/usage/log.ts index e5f74f3b9c..9bfa3c7ad4 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -44,6 +44,8 @@ export type FailureStage = | "downstream_write" | "client_cancel" | "terminal_delivery"; +export type TransportPhase = "pre_headers" | "mid_stream" | "terminal_sse"; +export type TerminalSource = "upstream" | "synthetic"; /** * Accepts EITHER label family. This is the predicate the persistence writers use, so widening @@ -122,6 +124,8 @@ export interface PersistedUsageAttempt { streamTimeline?: StreamTimeline; failureSide?: FailureSide; failureStage?: FailureStage; + transportPhase?: TransportPhase; + terminalSource?: TerminalSource; } export interface PersistedUsageEntry { @@ -182,8 +186,8 @@ export interface PersistedUsageEntry { streamTimeline?: StreamTimeline; failureSide?: FailureSide; failureStage?: FailureStage; - transportPhase?: string; - terminalSource?: string; + transportPhase?: TransportPhase; + terminalSource?: TerminalSource; /** * Bounded route-decision trace (RI-01): why this provider/model/account was * selected. Additive field; old rows without it parse unchanged. Never @@ -329,6 +333,15 @@ const KNOWN_FAILURE_STAGES = new Set([ "client_cancel", "terminal_delivery", ]); +const KNOWN_TRANSPORT_PHASES = new Set([ + "pre_headers", + "mid_stream", + "terminal_sse", +]); +const KNOWN_TERMINAL_SOURCES = new Set([ + "upstream", + "synthetic", +]); function normalizeStreamTimeline(raw: unknown): StreamTimeline | null { if (!raw || typeof raw !== "object" || Array.isArray(raw)) return null; @@ -514,6 +527,12 @@ function normalizeUsageAttempt(raw: unknown): PersistedUsageAttempt | null { ...(typeof attempt.failureStage === "string" && KNOWN_FAILURE_STAGES.has(attempt.failureStage as FailureStage) ? { failureStage: attempt.failureStage as FailureStage } : {}), + ...(typeof attempt.transportPhase === "string" && KNOWN_TRANSPORT_PHASES.has(attempt.transportPhase as TransportPhase) + ? { transportPhase: attempt.transportPhase as TransportPhase } + : {}), + ...(typeof attempt.terminalSource === "string" && KNOWN_TERMINAL_SOURCES.has(attempt.terminalSource as TerminalSource) + ? { terminalSource: attempt.terminalSource as TerminalSource } + : {}), }; } @@ -633,11 +652,11 @@ function normalizeUsageEntry(entry: PersistedUsageEntry): PersistedUsageEntry { ...(typeof entry.failureStage === "string" && KNOWN_FAILURE_STAGES.has(entry.failureStage as FailureStage) ? { failureStage: entry.failureStage as FailureStage } : {}), - ...(typeof entry.transportPhase === "string" && entry.transportPhase - ? { transportPhase: capMetadataString(entry.transportPhase) } + ...(typeof entry.transportPhase === "string" && KNOWN_TRANSPORT_PHASES.has(entry.transportPhase as TransportPhase) + ? { transportPhase: entry.transportPhase as TransportPhase } : {}), - ...(typeof entry.terminalSource === "string" && entry.terminalSource - ? { terminalSource: capMetadataString(entry.terminalSource) } + ...(typeof entry.terminalSource === "string" && KNOWN_TERMINAL_SOURCES.has(entry.terminalSource as TerminalSource) + ? { terminalSource: entry.terminalSource as TerminalSource } : {}), ...(routeDecision ? { routeDecision } : {}), }; diff --git a/tests/usage-log.test.ts b/tests/usage-log.test.ts index 58d9273002..4e3df5d4c2 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -956,4 +956,49 @@ describe("usage log", () => { expect(row?.transportPhase).toBe("mid_stream"); expect(row?.terminalSource).toBe("synthetic"); }); + + test("drops unknown or invalid transportPhase and terminalSource (#1217)", () => { + appendUsageEntry({ + requestId: "ocx-stream-invalid-attribution-test", + timestamp: Date.now(), + provider: "anthropic", + model: "claude-sonnet-5", + status: 502, + durationMs: 1000, + usageStatus: "unreported", + failureSide: "invalid_side" as unknown as any, + failureStage: "invalid_stage" as unknown as any, + transportPhase: "invalid_phase" as unknown as any, + terminalSource: "invalid_source" as unknown as any, + attempts: [ + { + ordinal: 1, + adapter: "anthropic", + sendCount: 1, + usageStatus: "unreported", + timestamp: Date.now(), + provider: "anthropic", + model: "claude-sonnet-5", + status: 502, + durationMs: 1000, + failureSide: "bogus_side" as unknown as any, + failureStage: "bogus_stage" as unknown as any, + transportPhase: "bogus_phase" as unknown as any, + terminalSource: "bogus_source" as unknown as any, + }, + ], + }); + + const entries = readRecentUsageEntries(10); + const row = entries.find(e => e.requestId === "ocx-stream-invalid-attribution-test"); + expect(row).toBeDefined(); + expect(row?.failureSide).toBeUndefined(); + expect(row?.failureStage).toBeUndefined(); + expect(row?.transportPhase).toBeUndefined(); + expect(row?.terminalSource).toBeUndefined(); + expect(row?.attempts?.[0].failureSide).toBeUndefined(); + expect(row?.attempts?.[0].failureStage).toBeUndefined(); + expect(row?.attempts?.[0].transportPhase).toBeUndefined(); + expect(row?.attempts?.[0].terminalSource).toBeUndefined(); + }); }); From 1bf3b08be212c3d81e5adf806e6273886aae6aa2 Mon Sep 17 00:00:00 2001 From: chilung Date: Sun, 23 Aug 2026 10:40:53 +0000 Subject: [PATCH 4/9] fix(usage): populate stream timeline and attribution from streaming events (#1217) --- src/server/relay.ts | 22 ++++++++++- src/server/request-log.ts | 64 ++++++++++++++++++++++++++++-- tests/usage-log.test.ts | 82 +++++++++++++++++++++++++++++++++++++++ 3 files changed, 162 insertions(+), 6 deletions(-) diff --git a/src/server/relay.ts b/src/server/relay.ts index a523ffd7e5..783848d636 100644 --- a/src/server/relay.ts +++ b/src/server/relay.ts @@ -15,6 +15,7 @@ import { httpStatusForRequestLogTerminal, inspectResponseLogJson, inspectResponseLogSsePayloadParsed, + noteStreamTimelineEvent, recordFirstOutput, type RequestLogContext, type RequestLogEntry, @@ -1021,6 +1022,9 @@ export function createSseInspector(handlers: SseInspectorHandlers): SseInspector try { handlers.onParsedPayload(parsed); } catch { /* inspection must never throw into the pump */ } } reportFirstOutput.parsed(parsed); + if (handlers.logCtx && handlers.logCtx.firstOutputMs !== undefined) { + noteStreamTimelineEvent(handlers.logCtx, "upstreamFirstSemanticOutputMs"); + } const status = terminalStatusFromParsed(parsed); const policyTerminal = status === "failed" && isPolicyRewriteType(parsed) @@ -1032,6 +1036,10 @@ export function createSseInspector(handlers: SseInspectorHandlers): SseInspector if (handlers.logCtx) { handlers.logCtx.transportPhase = "terminal_sse"; handlers.logCtx.terminalSource = "upstream"; + if (status === "failed") { + handlers.logCtx.failureSide = "upstream"; + handlers.logCtx.failureStage = "terminal_delivery"; + } } handlers.onTerminal(status, policyTerminal ? 400 : undefined); } finally { @@ -1158,7 +1166,11 @@ export function createSseInspector(handlers: SseInspectorHandlers): SseInspector return { feed(chunk) { - if (!disposed) scanChunk(chunk); + if (disposed) return; + if (chunk.byteLength > 0 && handlers.logCtx) { + noteStreamTimelineEvent(handlers.logCtx, "upstreamFirstByteMs"); + } + scanChunk(chunk); }, finish() { if (disposed) return; @@ -1367,7 +1379,11 @@ export function consumeForInspection( onCancel, onCleanEof: () => { if (!inspector.reported()) { - if (logCtx) logCtx.terminalSource = "synthetic"; + if (logCtx) { + logCtx.terminalSource = "synthetic"; + logCtx.failureSide = "upstream"; + logCtx.failureStage = "upstream_read"; + } onTerminal("incomplete"); } }, @@ -1379,6 +1395,8 @@ export function consumeForInspection( if (logCtx) { logCtx.transportPhase = "mid_stream"; logCtx.terminalSource = "synthetic"; + logCtx.failureSide = "upstream"; + logCtx.failureStage = "upstream_read"; // A truncated 200 body must not meter as a success the client never // received; the router's equivalent turn carries 502 + streamAborted // (codex-router #139). diff --git a/src/server/request-log.ts b/src/server/request-log.ts index 2c8d3e179c..7473a475a5 100644 --- a/src/server/request-log.ts +++ b/src/server/request-log.ts @@ -29,8 +29,13 @@ import { usageStatusForFinalLog, usageTotalTokens, type AttemptRecoveryKind, + type FailureSide, + type FailureStage, type PersistedUsageAttempt, type PersistedUsageEntry, + type StreamTimeline, + type TerminalSource, + type TransportPhase, type UsageStatus, } from "../usage/log"; import { @@ -50,6 +55,8 @@ import { modelRecordValue } from "../reasoning-effort"; export interface RequestLogContext { model: string; provider: string; + /** Internal request start timestamp in wall-clock ms. */ + requestStartedAt?: number; /** TTFT: ms from request start to the first non-empty model output delta (WP4, devlog 040). */ firstOutputMs?: number; /** Best-effort chat/session correlation for Logs grouping (#330). Opaque; omit when unknown. */ @@ -134,8 +141,11 @@ export interface RequestLogContext { /** Structured reason from `response.incomplete`; internal-only input to log classification. */ terminalIncompleteReason?: string; affinity?: "reused" | "new_bind" | "rebound" | "cleared"; - transportPhase?: "pre_headers" | "mid_stream" | "terminal_sse"; - terminalSource?: "upstream" | "synthetic"; + transportPhase?: TransportPhase; + terminalSource?: TerminalSource; + streamTimeline?: StreamTimeline; + failureSide?: FailureSide; + failureStage?: FailureStage; /** Bounded route-decision trace (RI-01); never contains secrets. */ routeDecision?: RouteDecisionTraceV1; } @@ -198,9 +208,12 @@ export interface RequestLogEntry { /** Codex pool affinity decision for this request (diagnostics for #186). */ affinity?: "reused" | "new_bind" | "rebound" | "cleared"; /** Where the upstream terminal/failure was observed. */ - transportPhase?: "pre_headers" | "mid_stream" | "terminal_sse"; + transportPhase?: TransportPhase; /** Whether the terminal came from a real upstream SSE event or a proxy synthetic tail. */ - terminalSource?: "upstream" | "synthetic"; + terminalSource?: TerminalSource; + streamTimeline?: StreamTimeline; + failureSide?: FailureSide; + failureStage?: FailureStage; /** Bounded route-decision trace (RI-01); never contains secrets. */ routeDecision?: RouteDecisionTraceV1; } @@ -312,6 +325,11 @@ export function requestLogEntryFromPersistedUsage(entry: PersistedUsageEntry): R ...(entry.usage ? { usage: entry.usage } : {}), ...(entry.totalTokens !== undefined ? { totalTokens: entry.totalTokens } : {}), ...(entry.attempts !== undefined ? { attempts: entry.attempts } : {}), + ...(entry.streamTimeline ? { streamTimeline: entry.streamTimeline } : {}), + ...(entry.failureSide ? { failureSide: entry.failureSide } : {}), + ...(entry.failureStage ? { failureStage: entry.failureStage } : {}), + ...(entry.transportPhase ? { transportPhase: entry.transportPhase } : {}), + ...(entry.terminalSource ? { terminalSource: entry.terminalSource } : {}), ...(routeDecision ? { routeDecision } : {}), }; } @@ -430,6 +448,11 @@ export function addRequestLog(entry: RequestLogEntry) { ...(entry.totalTokens !== undefined ? { totalTokens: entry.totalTokens } : {}), ...(entry.attempts !== undefined ? { attempts: entry.attempts } : {}), ...failureDiagnostics, + ...(entry.streamTimeline ? { streamTimeline: entry.streamTimeline } : {}), + ...(entry.failureSide ? { failureSide: entry.failureSide } : {}), + ...(entry.failureStage ? { failureStage: entry.failureStage } : {}), + ...(entry.transportPhase ? { transportPhase: entry.transportPhase } : {}), + ...(entry.terminalSource ? { terminalSource: entry.terminalSource } : {}), ...(entry.routeDecision ? { routeDecision: entry.routeDecision } : {}), }); } catch { @@ -461,6 +484,30 @@ export function recordFirstOutput( } } +export function noteStreamTimelineEvent( + logCtx: RequestLogContext | undefined, + event: keyof StreamTimeline, + requestStartedAt?: number, + now = Date.now(), +): void { + if (!logCtx) return; + if (requestStartedAt && !logCtx.requestStartedAt) { + logCtx.requestStartedAt = requestStartedAt; + } + if (!logCtx.streamTimeline) logCtx.streamTimeline = {}; + const origin = requestStartedAt ?? logCtx.requestStartedAt ?? logCtx.activeAttemptStartedAt; + if (origin !== undefined && logCtx.streamTimeline[event] === undefined) { + logCtx.streamTimeline[event] = Math.max(0, now - origin); + } + if (logCtx.activeAttempt) { + if (!logCtx.activeAttempt.streamTimeline) logCtx.activeAttempt.streamTimeline = {}; + const attemptOrigin = logCtx.activeAttemptStartedAt ?? requestStartedAt ?? logCtx.requestStartedAt; + if (attemptOrigin !== undefined && logCtx.activeAttempt.streamTimeline[event] === undefined) { + logCtx.activeAttempt.streamTimeline[event] = Math.max(0, now - attemptOrigin); + } + } +} + /** Snapshot target-specific requested effort even for runTurn adapters with no AdapterRequest. */ export function recordAttemptRequestedEffort(logCtx: RequestLogContext): void { const attempt = logCtx.activeAttempt; @@ -926,6 +973,7 @@ export function addFinalRequestLog( meta?: Pick, addLog: (entry: RequestLogEntry) => void = addRequestLog, ): void { + if (!logCtx.requestStartedAt) logCtx.requestStartedAt = start; // 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) @@ -954,6 +1002,11 @@ export function addFinalRequestLog( // semantic code on both so detailed attempt telemetry cannot regress to a generic status code. if (errorCode) logCtx.activeAttempt.errorCode = errorCode; else delete logCtx.activeAttempt.errorCode; + if (logCtx.streamTimeline) logCtx.activeAttempt.streamTimeline = { ...logCtx.streamTimeline }; + if (logCtx.failureSide) logCtx.activeAttempt.failureSide = logCtx.failureSide; + if (logCtx.failureStage) logCtx.activeAttempt.failureStage = logCtx.failureStage; + if (logCtx.transportPhase) logCtx.activeAttempt.transportPhase = logCtx.transportPhase; + if (logCtx.terminalSource) logCtx.activeAttempt.terminalSource = logCtx.terminalSource; } const existing = finalizedUsage( logCtx.providerAdapter ?? logCtx.provider, @@ -1027,6 +1080,9 @@ export function addFinalRequestLog( ...(logCtx.affinity ? { affinity: logCtx.affinity } : {}), ...(logCtx.transportPhase ? { transportPhase: logCtx.transportPhase } : {}), ...(logCtx.terminalSource ? { terminalSource: logCtx.terminalSource } : {}), + ...(logCtx.streamTimeline ? { streamTimeline: logCtx.streamTimeline } : {}), + ...(logCtx.failureSide ? { failureSide: logCtx.failureSide } : {}), + ...(logCtx.failureStage ? { failureStage: logCtx.failureStage } : {}), ...(logCtx.routeDecision ? { routeDecision: logCtx.routeDecision } : {}), }); if (isUsageDebugEnabled()) { diff --git a/tests/usage-log.test.ts b/tests/usage-log.test.ts index 4e3df5d4c2..421582c8c0 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -1001,4 +1001,86 @@ describe("usage log", () => { expect(row?.attempts?.[0].transportPhase).toBeUndefined(); expect(row?.attempts?.[0].terminalSource).toBeUndefined(); }); + + test("drops credential-shaped values from diagnostic metadata (#1217)", () => { + appendUsageEntry({ + requestId: "ocx-stream-credential-drop-test", + timestamp: Date.now(), + provider: "anthropic", + model: "claude-sonnet-5", + status: 502, + durationMs: 1000, + usageStatus: "unreported", + failureSide: "Basic dXNlcjpwYXNz" as unknown as any, + failureStage: "token=eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9" as unknown as any, + transportPhase: "password=super-secret-pw" as unknown as any, + terminalSource: "api-key=supersecret12345" as unknown as any, + }); + + const entries = readRecentUsageEntries(10); + const row = entries.find(e => e.requestId === "ocx-stream-credential-drop-test"); + expect(row).toBeDefined(); + expect(row?.failureSide).toBeUndefined(); + expect(row?.failureStage).toBeUndefined(); + expect(row?.transportPhase).toBeUndefined(); + expect(row?.terminalSource).toBeUndefined(); + }); + + test("shared request logging flow populates streamTimeline and failure attribution from real streaming failure (#1217)", async () => { + const { addFinalRequestLog } = await import("../src/server/request-log"); + const { consumeForInspection } = await import("../src/server/relay"); + type ReqLogCtx = import("../src/server/request-log").RequestLogContext; + + const requestId = "ocx-real-stream-fail-test"; + const start = Date.now() - 50; + const logCtx: ReqLogCtx = { + model: "claude-sonnet-5", + provider: "anthropic", + providerAdapter: "anthropic", + requestStartedAt: start, + }; + + // Simulate an SSE stream that delivers data then encounters an upstream read error (socket reset) + let sentFirst = false; + const stream = new ReadableStream({ + async pull(controller) { + if (!sentFirst) { + sentFirst = true; + controller.enqueue(new TextEncoder().encode('data: {"type":"response.output_item.added"}\n\n')); + return; + } + controller.error(new Error("socket reset mid-stream")); + }, + }); + + await new Promise(resolve => { + consumeForInspection( + stream, + (terminalStatus, httpStatusOverride) => { + addFinalRequestLog( + requestId, + start, + logCtx, + httpStatusOverride ?? (terminalStatus === "completed" ? 200 : 502), + { terminalStatus }, + ); + resolve(); + }, + undefined, + () => resolve(), + logCtx, + ); + }); + + const entries = readRecentUsageEntries(10); + const row = entries.find(e => e.requestId === requestId); + expect(row).toBeDefined(); + expect(row?.status).toBe(502); + expect(row?.terminalStatus).toBe("failed"); + expect(row?.transportPhase).toBe("mid_stream"); + expect(row?.terminalSource).toBe("synthetic"); + expect(row?.failureSide).toBe("upstream"); + expect(row?.failureStage).toBe("upstream_read"); + expect(row?.streamTimeline?.upstreamFirstByteMs).toBeGreaterThanOrEqual(0); + }); }); From ac57ff22b1089217a52169bb7e7978cd50a7279e Mon Sep 17 00:00:00 2001 From: chilung Date: Mon, 24 Aug 2026 04:33:01 +0000 Subject: [PATCH 5/9] fix(usage): preserve attempt-relative timeline on retry and set transportPhase on synthetic EOF (#1217) --- src/server/relay.ts | 1 + src/server/request-log.ts | 4 +- tests/usage-log.test.ts | 98 +++++++++++++++++++++++++++++++++++++++ 3 files changed, 102 insertions(+), 1 deletion(-) diff --git a/src/server/relay.ts b/src/server/relay.ts index 783848d636..d43ed6d5ea 100644 --- a/src/server/relay.ts +++ b/src/server/relay.ts @@ -1380,6 +1380,7 @@ export function consumeForInspection( onCleanEof: () => { if (!inspector.reported()) { if (logCtx) { + logCtx.transportPhase = "mid_stream"; logCtx.terminalSource = "synthetic"; logCtx.failureSide = "upstream"; logCtx.failureStage = "upstream_read"; diff --git a/src/server/request-log.ts b/src/server/request-log.ts index 7473a475a5..bed655f1c1 100644 --- a/src/server/request-log.ts +++ b/src/server/request-log.ts @@ -1002,7 +1002,9 @@ export function addFinalRequestLog( // semantic code on both so detailed attempt telemetry cannot regress to a generic status code. if (errorCode) logCtx.activeAttempt.errorCode = errorCode; else delete logCtx.activeAttempt.errorCode; - if (logCtx.streamTimeline) logCtx.activeAttempt.streamTimeline = { ...logCtx.streamTimeline }; + if (logCtx.streamTimeline && !logCtx.activeAttempt.streamTimeline) { + logCtx.activeAttempt.streamTimeline = { ...logCtx.streamTimeline }; + } if (logCtx.failureSide) logCtx.activeAttempt.failureSide = logCtx.failureSide; if (logCtx.failureStage) logCtx.activeAttempt.failureStage = logCtx.failureStage; if (logCtx.transportPhase) logCtx.activeAttempt.transportPhase = logCtx.transportPhase; diff --git a/tests/usage-log.test.ts b/tests/usage-log.test.ts index 421582c8c0..25572d0cec 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -1083,4 +1083,102 @@ describe("usage log", () => { expect(row?.failureStage).toBe("upstream_read"); expect(row?.streamTimeline?.upstreamFirstByteMs).toBeGreaterThanOrEqual(0); }); + + test("synthetic clean EOF sets transportPhase to mid_stream", async () => { + const { addFinalRequestLog } = await import("../src/server/request-log"); + const { consumeForInspection } = await import("../src/server/relay"); + const requestId = "ocx-stream-clean-eof-test"; + const start = Date.now() - 50; + const logCtx: RequestLogContext = { + requestId, + provider: "anthropic", + model: "claude-sonnet-5", + requestStartedAt: start, + activeAttemptStartedAt: start, + activeAttempt: { + ordinal: 1, + provider: "anthropic", + model: "claude-sonnet-5", + adapter: "anthropic", + status: 200, + durationMs: 50, + sendCount: 1, + recoveryKinds: [], + usageStatus: "unreported", + }, + }; + + // An empty SSE stream that terminates with clean EOF without a completed terminal event + const stream = new ReadableStream({ + start(controller) { + controller.close(); + }, + }); + + await new Promise(resolve => { + consumeForInspection( + stream, + (terminalStatus, httpStatusOverride) => { + addFinalRequestLog( + requestId, + start, + logCtx, + httpStatusOverride ?? 502, + { terminalStatus }, + ); + resolve(); + }, + undefined, + () => resolve(), + logCtx, + ); + }); + + const entries = readRecentUsageEntries(10); + const row = entries.find(e => e.requestId === requestId); + expect(row).toBeDefined(); + expect(row?.transportPhase).toBe("mid_stream"); + expect(row?.terminalSource).toBe("synthetic"); + expect(row?.failureSide).toBe("upstream"); + expect(row?.failureStage).toBe("upstream_read"); + }); + + test("preserves distinct attempt-relative and request-relative streamTimeline on retry", async () => { + const { addFinalRequestLog, noteStreamTimelineEvent } = await import("../src/server/request-log"); + const requestId = "ocx-stream-retry-timeline-test"; + const requestStart = 10000; + const attemptStart = 12000; // attempt starts 2000ms after request + const now = 13500; + + const attempt: PersistedUsageAttempt = { + ordinal: 2, + provider: "anthropic", + model: "claude-sonnet-5", + adapter: "anthropic", + status: 200, + durationMs: 1500, + sendCount: 2, + recoveryKinds: ["retry"], + usageStatus: "reported", + }; + const logCtx: RequestLogContext = { + requestId, + provider: "anthropic", + model: "claude-sonnet-5", + requestStartedAt: requestStart, + activeAttemptStartedAt: attemptStart, + activeAttempt: attempt, + attempts: [attempt], + }; + + noteStreamTimelineEvent(logCtx, "upstreamFirstByteMs", requestStart, now); + expect(logCtx.streamTimeline?.upstreamFirstByteMs).toBe(3500); // 13500 - 10000 + expect(logCtx.activeAttempt?.streamTimeline?.upstreamFirstByteMs).toBe(1500); // 13500 - 12000 + + addFinalRequestLog(requestId, requestStart, logCtx, 200, { terminalStatus: "completed" }); + const entries = readRecentUsageEntries(10); + const row = entries.find(e => e.requestId === requestId); + expect(row?.streamTimeline?.upstreamFirstByteMs).toBe(3500); + expect(row?.attempts?.[0]?.streamTimeline?.upstreamFirstByteMs).toBe(1500); + }); }); From c37ea24e0f41305fb5c7449655d6ce3fd79f377d Mon Sep 17 00:00:00 2001 From: chilung Date: Mon, 24 Aug 2026 15:04:13 +0000 Subject: [PATCH 6/9] fix(request-log): seed stream origins and normalize diagnostics --- src/server/index.ts | 9 +++ src/server/relay.ts | 8 +++ src/server/request-log.ts | 28 ++++++-- src/server/responses/core.ts | 2 + src/usage/log.ts | 59 +++++++++------- tests/request-log.test.ts | 128 ++++++++++++++++++++++++++++++++++- 6 files changed, 203 insertions(+), 31 deletions(-) diff --git a/src/server/index.ts b/src/server/index.ts index eab0ba1d7f..c2e29d8503 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -1624,6 +1624,7 @@ export function startServer(port?: number, deps: StartServerDeps = {}): Server { @@ -1737,6 +1740,7 @@ export function startServer(port?: number, deps: StartServerDeps = {}): Server { @@ -1909,6 +1916,7 @@ export function startServer(port?: number, deps: StartServerDeps = {}): Server void = addRequestLog, ): Response { + if (logCtx.requestStartedAt === undefined) logCtx.requestStartedAt = start; const contentType = response.headers.get("content-type")?.toLowerCase() ?? ""; if (isUsageDebugEnabled() && !logCtx.usageDebugContentType && contentType) { logCtx.usageDebugContentType = contentType; @@ -1196,6 +1197,7 @@ export function createSseInspector(handlers: SseInspectorHandlers): SseInspector export type InspectionDrainBounds = { ms: number; bytes: number }; export type InspectionConsumerOptions = { + requestStartedAt?: number; clientGoneSignal?: AbortSignal; drainBounds?: Partial; upstream?: AbortController; @@ -1361,6 +1363,9 @@ export function consumeForInspection( onFirstOutput?: () => void, options?: InspectionConsumerOptions, ): void { + if (logCtx && options?.requestStartedAt !== undefined && logCtx.requestStartedAt === undefined) { + logCtx.requestStartedAt = options.requestStartedAt; + } const reader = body.getReader(); const inspector = (options?.inspectorFactory ?? createSseInspector)({ onTerminal, @@ -1418,6 +1423,9 @@ export function consumeForResponseLogMetadata( onFirstOutput?: () => void, options?: InspectionConsumerOptions, ): void { + if (options?.requestStartedAt !== undefined && logCtx.requestStartedAt === undefined) { + logCtx.requestStartedAt = options.requestStartedAt; + } const reader = body.getReader(); // No onTerminal → the inspector's `reported` gate stays permanently false, // reproducing this consumer's unconditional logCtx inspection. diff --git a/src/server/request-log.ts b/src/server/request-log.ts index bed655f1c1..c799793569 100644 --- a/src/server/request-log.ts +++ b/src/server/request-log.ts @@ -24,6 +24,7 @@ import { isKnownUsageSurface, isCodexUsageAccountLogLabel, isValidReasoningWireValue, + normalizeStreamDiagnostics, readRecentUsageEntries, usageForFinalLog, usageStatusForFinalLog, @@ -386,10 +387,29 @@ export function addRequestLog(entry: RequestLogEntry) { // line-oriented viewer — while `usage.jsonl` looked clean, which is the worst shape for a // sanitization bug because the safe surface is the one you check. const shadowCallRewrittenFrom = sanitizeLogMetadataString(entry.shadowCallRewrittenFrom); - const retained: RequestLogEntry = shadowCallRewrittenFrom === entry.shadowCallRewrittenFrom - ? entry - : { ...entry, ...(shadowCallRewrittenFrom ? { shadowCallRewrittenFrom } : {}) }; - if (!shadowCallRewrittenFrom && retained !== entry) delete retained.shadowCallRewrittenFrom; + const diagnostics = normalizeStreamDiagnostics(entry); + const attempts = entry.attempts?.map(attempt => { + const normalized = { ...attempt }; + const attemptDiagnostics = normalizeStreamDiagnostics(attempt); + if (!attemptDiagnostics.streamTimeline) delete normalized.streamTimeline; + if (!attemptDiagnostics.failureSide) delete normalized.failureSide; + if (!attemptDiagnostics.failureStage) delete normalized.failureStage; + if (!attemptDiagnostics.transportPhase) delete normalized.transportPhase; + if (!attemptDiagnostics.terminalSource) delete normalized.terminalSource; + return { ...normalized, ...attemptDiagnostics }; + }); + const retained: RequestLogEntry = { + ...entry, + ...(shadowCallRewrittenFrom ? { shadowCallRewrittenFrom } : {}), + ...(attempts ? { attempts } : {}), + ...diagnostics, + }; + if (!shadowCallRewrittenFrom) delete retained.shadowCallRewrittenFrom; + if (!diagnostics.streamTimeline) delete retained.streamTimeline; + if (!diagnostics.failureSide) delete retained.failureSide; + if (!diagnostics.failureStage) delete retained.failureStage; + if (!diagnostics.transportPhase) delete retained.transportPhase; + if (!diagnostics.terminalSource) delete retained.terminalSource; entry = retained; retainRequestLogEntry(entry); try { diff --git a/src/server/responses/core.ts b/src/server/responses/core.ts index e2e23f241d..4998171970 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -2346,6 +2346,7 @@ export async function handleComboResponses( const childLog: RequestLogContext = { model: pick.target.model, provider: pick.target.provider, + ...(logCtx.requestStartedAt !== undefined ? { requestStartedAt: logCtx.requestStartedAt } : {}), ...(logCtx.conversationId ? { conversationId: logCtx.conversationId } : {}), ...(logCtx.surface ? { surface: logCtx.surface } : {}), }; @@ -4801,6 +4802,7 @@ async function handleResponsesInner( linkAbortSignal(upstream, turnAc.signal); registerTurn(turnAc, options.turnAdmissionLease); const inspectionConsumerOptions = { + requestStartedAt: logCtx.requestStartedAt, clientGoneSignal: clientGone.signal, drainBounds: { ms: 15_000, bytes: 32 * 1024 * 1024 }, upstream, diff --git a/src/usage/log.ts b/src/usage/log.ts index 9bfa3c7ad4..293195972d 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -363,6 +363,37 @@ function normalizeStreamTimeline(raw: unknown): StreamTimeline | null { return Object.keys(out).length > 0 ? out : null; } +export function normalizeStreamDiagnostics(raw: { + streamTimeline?: unknown; + failureSide?: unknown; + failureStage?: unknown; + transportPhase?: unknown; + terminalSource?: unknown; +}): { + streamTimeline?: StreamTimeline; + failureSide?: FailureSide; + failureStage?: FailureStage; + transportPhase?: TransportPhase; + terminalSource?: TerminalSource; +} { + const streamTimeline = normalizeStreamTimeline(raw.streamTimeline); + return { + ...(streamTimeline ? { streamTimeline } : {}), + ...(typeof raw.failureSide === "string" && KNOWN_FAILURE_SIDES.has(raw.failureSide as FailureSide) + ? { failureSide: raw.failureSide as FailureSide } + : {}), + ...(typeof raw.failureStage === "string" && KNOWN_FAILURE_STAGES.has(raw.failureStage as FailureStage) + ? { failureStage: raw.failureStage as FailureStage } + : {}), + ...(typeof raw.transportPhase === "string" && KNOWN_TRANSPORT_PHASES.has(raw.transportPhase as TransportPhase) + ? { transportPhase: raw.transportPhase as TransportPhase } + : {}), + ...(typeof raw.terminalSource === "string" && KNOWN_TERMINAL_SOURCES.has(raw.terminalSource as TerminalSource) + ? { terminalSource: raw.terminalSource as TerminalSource } + : {}), + }; +} + export function isLabRouteSubjectId(value: unknown): value is string { return typeof value === "string" && LAB_ROUTE_SUBJECT_ID_RE.test(value); } @@ -520,19 +551,7 @@ function normalizeUsageAttempt(raw: unknown): PersistedUsageAttempt | null { : { reasoningWireValue: attempt.reasoningWireValue } : {}), ...(tierOutcome ? { tierOutcome } : {}), - ...(normalizeStreamTimeline(attempt.streamTimeline) ? { streamTimeline: normalizeStreamTimeline(attempt.streamTimeline) as StreamTimeline } : {}), - ...(typeof attempt.failureSide === "string" && KNOWN_FAILURE_SIDES.has(attempt.failureSide as FailureSide) - ? { failureSide: attempt.failureSide as FailureSide } - : {}), - ...(typeof attempt.failureStage === "string" && KNOWN_FAILURE_STAGES.has(attempt.failureStage as FailureStage) - ? { failureStage: attempt.failureStage as FailureStage } - : {}), - ...(typeof attempt.transportPhase === "string" && KNOWN_TRANSPORT_PHASES.has(attempt.transportPhase as TransportPhase) - ? { transportPhase: attempt.transportPhase as TransportPhase } - : {}), - ...(typeof attempt.terminalSource === "string" && KNOWN_TERMINAL_SOURCES.has(attempt.terminalSource as TerminalSource) - ? { terminalSource: attempt.terminalSource as TerminalSource } - : {}), + ...normalizeStreamDiagnostics(attempt), }; } @@ -645,19 +664,7 @@ function normalizeUsageEntry(entry: PersistedUsageEntry): PersistedUsageEntry { ...(entry.terminalStatus ? { terminalStatus: entry.terminalStatus } : {}), ...(entry.closeReason ? { closeReason: entry.closeReason } : {}), ...(entry.upstreamError ? { upstreamError: entry.upstreamError } : {}), - ...(normalizeStreamTimeline(entry.streamTimeline) ? { streamTimeline: normalizeStreamTimeline(entry.streamTimeline) as StreamTimeline } : {}), - ...(typeof entry.failureSide === "string" && KNOWN_FAILURE_SIDES.has(entry.failureSide as FailureSide) - ? { failureSide: entry.failureSide as FailureSide } - : {}), - ...(typeof entry.failureStage === "string" && KNOWN_FAILURE_STAGES.has(entry.failureStage as FailureStage) - ? { failureStage: entry.failureStage as FailureStage } - : {}), - ...(typeof entry.transportPhase === "string" && KNOWN_TRANSPORT_PHASES.has(entry.transportPhase as TransportPhase) - ? { transportPhase: entry.transportPhase as TransportPhase } - : {}), - ...(typeof entry.terminalSource === "string" && KNOWN_TERMINAL_SOURCES.has(entry.terminalSource as TerminalSource) - ? { terminalSource: entry.terminalSource as TerminalSource } - : {}), + ...normalizeStreamDiagnostics(entry), ...(routeDecision ? { routeDecision } : {}), }; } diff --git a/tests/request-log.test.ts b/tests/request-log.test.ts index bbf5b1071f..57938f4382 100644 --- a/tests/request-log.test.ts +++ b/tests/request-log.test.ts @@ -31,9 +31,10 @@ import { appendUsageEntry, readUsageEntries, resetUsageReadCacheForTests, + usageLogPath, type PersistedUsageEntry, } from "../src/usage/log"; -import { mkdtempSync} from "node:fs"; +import { mkdtempSync, readFileSync, rmSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { removeTreeWithRetry } from "./helpers/remove-tree"; @@ -356,6 +357,131 @@ describe("request log metadata", () => { } }); + test("the direct addRequestLog ingress normalizes nested attempt diagnostics", () => { + const home = mkdtempSync(join(tmpdir(), "ocx-diagnostic-ingress-")); + const previousHome = process.env.OPENCODEX_HOME; + process.env.OPENCODEX_HOME = home; + try { + clearRequestLogsForTests(); + resetUsageReadCacheForTests(); + addRequestLog({ + requestId: "ocx-diagnostic-direct", + timestamp: Date.now(), + provider: "anthropic", + model: "claude-sonnet-5", + status: 502, + durationMs: 10, + usageStatus: "unreported", + transportPhase: "password=super-secret-pw" as never, + terminalSource: "api-key=supersecret12345" as never, + attempts: [{ + ordinal: 1, + provider: "anthropic", + model: "claude-sonnet-5", + adapter: "anthropic", + status: 502, + durationMs: 10, + sendCount: 1, + recoveryKinds: [], + usageStatus: "unreported", + transportPhase: "password=super-secret-pw" as never, + terminalSource: "api-key=supersecret12345" as never, + }], + }); + + const inMemory = getRequestLogEntries()[0]?.attempts?.[0]; + const inMemoryEntry = getRequestLogEntries()[0]; + expect(inMemoryEntry?.transportPhase).toBeUndefined(); + expect(inMemoryEntry?.terminalSource).toBeUndefined(); + expect(inMemory?.transportPhase).toBeUndefined(); + expect(inMemory?.terminalSource).toBeUndefined(); + const inMemoryJson = JSON.stringify(inMemoryEntry); + expect(inMemoryJson).not.toContain("super-secret-pw"); + expect(inMemoryJson).not.toContain("supersecret12345"); + const persistedRaw = readFileSync(usageLogPath(), "utf8"); + expect(persistedRaw).not.toContain("super-secret-pw"); + expect(persistedRaw).not.toContain("supersecret12345"); + const persisted = JSON.parse(persistedRaw.trim()) as PersistedUsageEntry; + expect(persisted.transportPhase).toBeUndefined(); + expect(persisted.terminalSource).toBeUndefined(); + expect(persisted.attempts?.[0]?.transportPhase).toBeUndefined(); + expect(persisted.attempts?.[0]?.terminalSource).toBeUndefined(); + } finally { + clearRequestLogsForTests(); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + resetUsageReadCacheForTests(); + rmSync(home, { recursive: true, force: true }); + } + }); + + test("inspection seeds request-relative timeline origin before a retry attempt", async () => { + const { consumeForInspection } = await import("../src/server/relay"); + const requestStart = Date.now() - 50; + const attemptStart = requestStart + 10; + const attempt = { + ordinal: 2, + provider: "anthropic", + model: "claude-sonnet-5", + adapter: "anthropic", + status: 200, + durationMs: 1, + sendCount: 1, + recoveryKinds: ["transient-5xx" as const], + usageStatus: "unreported" as const, + }; + const logCtx: RequestLogContext = { + model: "claude-sonnet-5", + provider: "anthropic", + activeAttempt: attempt, + activeAttemptStartedAt: attemptStart, + }; + const entries: RequestLogEntry[] = []; + const payload = JSON.stringify({ + type: "response.completed", + response: { status: "completed", output: [] }, + }); + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(`data: ${payload}\n\n`)); + controller.close(); + }, + }); + + await new Promise(resolve => { + consumeForInspection( + body, + (terminalStatus, httpStatusOverride) => { + addFinalRequestLog( + "ocx-retry-origin", + requestStart, + logCtx, + httpStatusOverride ?? 200, + { terminalStatus }, + entry => { + entries.push(entry); + resolve(); + }, + ); + }, + undefined, + undefined, + logCtx, + undefined, + undefined, + undefined, + { requestStartedAt: requestStart }, + ); + }); + + expect(logCtx.requestStartedAt).toBe(requestStart); + expect(logCtx.streamTimeline?.upstreamFirstByteMs).toBeGreaterThanOrEqual(0); + expect(logCtx.activeAttempt?.streamTimeline?.upstreamFirstByteMs).toBeGreaterThanOrEqual(0); + expect(logCtx.streamTimeline?.upstreamFirstByteMs) + .toBeGreaterThanOrEqual(logCtx.activeAttempt?.streamTimeline?.upstreamFirstByteMs ?? 0); + expect(entries).toHaveLength(1); + }); + test("records ordered attempts with sealed identity, fresh estimates, and deduplicated recoveries", () => { const a = beginRequestAttempt(1, "provisional-a", "model-a", "openai-chat"); noteAttemptSend(a, 100); From 83933caca7f56d8a99629a866f02c003b71a2a72 Mon Sep 17 00:00:00 2001 From: chilung Date: Mon, 24 Aug 2026 17:41:38 +0000 Subject: [PATCH 7/9] fix(usage): preserve safe stream diagnostics --- src/usage/log.ts | 42 ++++++++++++++++++++++++----------------- tests/usage-log.test.ts | 10 +++++----- 2 files changed, 30 insertions(+), 22 deletions(-) diff --git a/src/usage/log.ts b/src/usage/log.ts index 293195972d..1fb6fc2b9e 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -44,8 +44,10 @@ export type FailureStage = | "downstream_write" | "client_cancel" | "terminal_delivery"; -export type TransportPhase = "pre_headers" | "mid_stream" | "terminal_sse"; -export type TerminalSource = "upstream" | "synthetic"; +/** Bounded, redacted diagnostic string describing the transport phase. */ +export type TransportPhase = string; +/** Bounded, redacted diagnostic string describing the terminal source. */ +export type TerminalSource = string; /** * Accepts EITHER label family. This is the predicate the persistence writers use, so widening @@ -333,15 +335,23 @@ const KNOWN_FAILURE_STAGES = new Set([ "client_cancel", "terminal_delivery", ]); -const KNOWN_TRANSPORT_PHASES = new Set([ - "pre_headers", - "mid_stream", - "terminal_sse", -]); -const KNOWN_TERMINAL_SOURCES = new Set([ - "upstream", - "synthetic", -]); +const MAX_STREAM_DIAGNOSTIC_LENGTH = 64; +const STREAM_DIAGNOSTIC_CONTROL_CHARS = /[\u0000-\u001f\u007f-\u009f\u2028\u2029]/g; + +/** + * Keep bounded diagnostic values when they are safe, but drop the complete value + * when redaction had to rewrite it. Keeping a redacted fragment would make the + * field look trustworthy while still retaining attacker-controlled context next + * to a credential-shaped marker. + */ +function normalizeBoundedStreamDiagnostic(value: unknown): string | undefined { + if (typeof value !== "string") return undefined; + const cleaned = value.trim().replace(STREAM_DIAGNOSTIC_CONTROL_CHARS, ""); + if (!cleaned) return undefined; + const unbounded = sanitizeLogMetadataString(cleaned, cleaned.length); + if (!unbounded || unbounded !== cleaned) return undefined; + return cleaned.slice(0, MAX_STREAM_DIAGNOSTIC_LENGTH); +} function normalizeStreamTimeline(raw: unknown): StreamTimeline | null { if (!raw || typeof raw !== "object" || Array.isArray(raw)) return null; @@ -377,6 +387,8 @@ export function normalizeStreamDiagnostics(raw: { terminalSource?: TerminalSource; } { const streamTimeline = normalizeStreamTimeline(raw.streamTimeline); + const transportPhase = normalizeBoundedStreamDiagnostic(raw.transportPhase); + const terminalSource = normalizeBoundedStreamDiagnostic(raw.terminalSource); return { ...(streamTimeline ? { streamTimeline } : {}), ...(typeof raw.failureSide === "string" && KNOWN_FAILURE_SIDES.has(raw.failureSide as FailureSide) @@ -385,12 +397,8 @@ export function normalizeStreamDiagnostics(raw: { ...(typeof raw.failureStage === "string" && KNOWN_FAILURE_STAGES.has(raw.failureStage as FailureStage) ? { failureStage: raw.failureStage as FailureStage } : {}), - ...(typeof raw.transportPhase === "string" && KNOWN_TRANSPORT_PHASES.has(raw.transportPhase as TransportPhase) - ? { transportPhase: raw.transportPhase as TransportPhase } - : {}), - ...(typeof raw.terminalSource === "string" && KNOWN_TERMINAL_SOURCES.has(raw.terminalSource as TerminalSource) - ? { terminalSource: raw.terminalSource as TerminalSource } - : {}), + ...(transportPhase ? { transportPhase: transportPhase as TransportPhase } : {}), + ...(terminalSource ? { terminalSource: terminalSource as TerminalSource } : {}), }; } diff --git a/tests/usage-log.test.ts b/tests/usage-log.test.ts index 25572d0cec..db72861741 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -957,7 +957,7 @@ describe("usage log", () => { expect(row?.terminalSource).toBe("synthetic"); }); - test("drops unknown or invalid transportPhase and terminalSource (#1217)", () => { + test("preserves bounded diagnostics but drops invalid closed attribution (#1217)", () => { appendUsageEntry({ requestId: "ocx-stream-invalid-attribution-test", timestamp: Date.now(), @@ -994,12 +994,12 @@ describe("usage log", () => { expect(row).toBeDefined(); expect(row?.failureSide).toBeUndefined(); expect(row?.failureStage).toBeUndefined(); - expect(row?.transportPhase).toBeUndefined(); - expect(row?.terminalSource).toBeUndefined(); + expect(row?.transportPhase).toBe("invalid_phase"); + expect(row?.terminalSource).toBe("invalid_source"); expect(row?.attempts?.[0].failureSide).toBeUndefined(); expect(row?.attempts?.[0].failureStage).toBeUndefined(); - expect(row?.attempts?.[0].transportPhase).toBeUndefined(); - expect(row?.attempts?.[0].terminalSource).toBeUndefined(); + expect(row?.attempts?.[0].transportPhase).toBe("bogus_phase"); + expect(row?.attempts?.[0].terminalSource).toBe("bogus_source"); }); test("drops credential-shaped values from diagnostic metadata (#1217)", () => { From 20c27a589e5d34824f1e6128f825f8c3aec0608f Mon Sep 17 00:00:00 2001 From: chilung Date: Sun, 30 Aug 2026 06:04:10 +0000 Subject: [PATCH 8/9] fix(relay): isolate inspected source attribution and gate semantic timeline on payload (#1217) --- src/server/relay.ts | 20 ++++++++--- tests/request-log.test.ts | 75 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 91 insertions(+), 4 deletions(-) diff --git a/src/server/relay.ts b/src/server/relay.ts index 595ec8a1bc..bc560a5755 100644 --- a/src/server/relay.ts +++ b/src/server/relay.ts @@ -586,6 +586,7 @@ export function trackSseForRequestLog( onCancel: () => void, logCtx?: RequestLogContext, onFirstOutput?: () => void, + inspectedSource?: "upstream" | "relay", ): ReadableStream { const reader = body.getReader(); let terminalReported = false; @@ -602,6 +603,7 @@ export function trackSseForRequestLog( onTerminal: reportTerminal, logCtx, onFirstOutput, + inspectedSource, }); return new ReadableStream({ @@ -646,6 +648,7 @@ export function responseWithDeferredRequestLog( start: number, logCtx: RequestLogContext, addLog: (entry: RequestLogEntry) => void = addRequestLog, + inspectedSource?: "upstream" | "relay", ): Response { if (logCtx.requestStartedAt === undefined) logCtx.requestStartedAt = start; const contentType = response.headers.get("content-type")?.toLowerCase() ?? ""; @@ -709,6 +712,7 @@ export function responseWithDeferredRequestLog( }, logCtx, () => recordFirstOutput(logCtx, start), + inspectedSource, ); return new Response(body, { status: response.status, @@ -841,6 +845,11 @@ export type SseInspectorHandlers = { * between `response.created` and `response.completed`. */ pinCompletedResponseIdToFirstSeen?: boolean; + /** + * Provenance of the stream being inspected. + * Raw upstream SSE defaults to "upstream"; adapter or bridge produced SSE uses "relay". + */ + inspectedSource?: "upstream" | "relay"; }; type CompletedOutputItem = { item: unknown; sourceBytes: number }; @@ -1023,7 +1032,7 @@ export function createSseInspector(handlers: SseInspectorHandlers): SseInspector try { handlers.onParsedPayload(parsed); } catch { /* inspection must never throw into the pump */ } } reportFirstOutput.parsed(parsed); - if (handlers.logCtx && handlers.logCtx.firstOutputMs !== undefined) { + if (handlers.logCtx && firstOutputFromParsed(parsed)) { noteStreamTimelineEvent(handlers.logCtx, "upstreamFirstSemanticOutputMs"); } const status = terminalStatusFromParsed(parsed); @@ -1035,11 +1044,14 @@ export function createSseInspector(handlers: SseInspectorHandlers): SseInspector try { reported = true; if (handlers.logCtx) { + const source = handlers.inspectedSource ?? "upstream"; handlers.logCtx.transportPhase = "terminal_sse"; - handlers.logCtx.terminalSource = "upstream"; + handlers.logCtx.terminalSource = source; if (status === "failed") { - handlers.logCtx.failureSide = "upstream"; - handlers.logCtx.failureStage = "terminal_delivery"; + handlers.logCtx.failureSide = handlers.logCtx.failureSide + ?? (source === "relay" ? "relay" : "upstream"); + handlers.logCtx.failureStage = handlers.logCtx.failureStage + ?? (source === "relay" ? "relay_transform" : "terminal_delivery"); } } handlers.onTerminal(status, policyTerminal ? 400 : undefined); diff --git a/tests/request-log.test.ts b/tests/request-log.test.ts index 57938f4382..7a546f5989 100644 --- a/tests/request-log.test.ts +++ b/tests/request-log.test.ts @@ -26,6 +26,8 @@ import { } from "../src/server/request-log"; import { handleResponses } from "../src/server/responses"; import { bridgeToResponsesSSE } from "../src/bridge"; +import { createSseInspector } from "../src/server/relay"; +import { createTranslatorBudget } from "../src/lib/translator-budget"; import type { AdapterEvent, OcxConfig, OcxUsage } from "../src/types"; import { appendUsageEntry, @@ -1722,6 +1724,79 @@ describe("request log metadata", () => { expect(entries[0].upstreamError).toContain("adapter_eof"); expect(entries[0].upstreamError).toContain("ended unexpectedly"); }); + + test("bridge-to-request-log translator overflow persists relay failure attribution", async () => { + const entries: RequestLogEntry[] = []; + const budget = createTranslatorBudget({ maxTurnBytes: 1_024 }); + async function* events(): AsyncGenerator { + yield { type: "tool_call_start", id: "call_1", name: "exec_command" }; + yield { type: "tool_call_delta", arguments: "x".repeat(2_048) }; + yield { type: "tool_call_end" }; + yield { type: "done" }; + } + const sse = bridgeToResponsesSSE(events(), "test-model", undefined, undefined, undefined, undefined, 2_000, { + translatorBudget: budget, + }); + const logCtx: RequestLogContext = { + model: "anthropic/claude-sonnet-4", + provider: "anthropic", + }; + const response = responseWithDeferredRequestLog( + new Response(sse, { status: 200, headers: { "content-type": "text/event-stream" } }), + "ocx-test-bridge-overflow", + Date.now(), + logCtx, + entry => entries.push(entry), + "relay", + ); + + await response.text(); + expect(entries).toHaveLength(1); + expect(entries[0]).toMatchObject({ + terminalStatus: "failed", + status: 502, + failureSide: "relay", + failureStage: "relay_transform", + terminalSource: "relay", + transportPhase: "terminal_sse", + }); + }); + + test("upstreamFirstSemanticOutputMs is gated by semantic payload, not pre-populated context", () => { + const logCtx: RequestLogContext = { + model: "openai/gpt-5.6-sol", + provider: "openai", + requestStartedAt: 1_000, + firstOutputMs: 50, + }; + const inspector = createSseInspector({ logCtx }); + const nonSemanticEvent = new TextEncoder().encode( + `data: ${JSON.stringify({ type: "response.created", response: { id: "resp_1" } })}\n\n` + ); + inspector.feed(nonSemanticEvent); + expect(logCtx.streamTimeline?.upstreamFirstByteMs).toBeDefined(); + expect(logCtx.streamTimeline?.upstreamFirstSemanticOutputMs).toBeUndefined(); + + const semanticEvent = new TextEncoder().encode( + `data: ${JSON.stringify({ type: "response.output_text.delta", delta: "hello" })}\n\n` + ); + inspector.feed(semanticEvent); + expect(logCtx.streamTimeline?.upstreamFirstSemanticOutputMs).toBeDefined(); + }); + + test("upstreamFirstSemanticOutputMs records when onFirstOutput is absent but semantic delta arrives", () => { + const logCtx: RequestLogContext = { + model: "openai/gpt-5.6-sol", + provider: "openai", + requestStartedAt: 1_000, + }; + const inspector = createSseInspector({ logCtx }); + const semanticEvent = new TextEncoder().encode( + `data: ${JSON.stringify({ type: "response.output_text.delta", delta: "hello" })}\n\n` + ); + inspector.feed(semanticEvent); + expect(logCtx.streamTimeline?.upstreamFirstSemanticOutputMs).toBeDefined(); + }); }); describe("request log restart hydrate", () => { From 309aa29ef9a3052b0b4011d1c387305d4a39bfd4 Mon Sep 17 00:00:00 2001 From: chilung Date: Sun, 30 Aug 2026 06:16:25 +0000 Subject: [PATCH 9/9] docs(structure): document stream timeline and failure attribution persistence contract (#1217) --- structure/05_gui-and-management-api.md | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/structure/05_gui-and-management-api.md b/structure/05_gui-and-management-api.md index ff75b15e88..748835ca88 100644 --- a/structure/05_gui-and-management-api.md +++ b/structure/05_gui-and-management-api.md @@ -337,6 +337,16 @@ keeps the saved state and renders fixed `ocx sync` guidance without server/accou An opt-in shadow-call rewrite persists the bounded, redacted original helper model as `shadowCallRewrittenFrom`, so helper traffic remains identifiable after restart without storing request content or inferring a helper subtype from timing. +Streaming request execution records bounded, attempt-isolated `streamTimeline` milestones +(`upstreamDispatchMs`, `upstreamHeadersMs`, `upstreamFirstByteMs`, `upstreamFirstSemanticOutputMs`, +`downstreamFirstWriteMs`, `upstreamEndMs`, `downstreamEndMs`) and closed-enum failure attribution +(`failureSide`: `upstream | relay | downstream | client | local`, `failureStage`: `pre_dispatch | upstream_wait_headers | upstream_read | relay_transform | downstream_write | client_cancel | terminal_delivery`). +Timeline recording is strictly content-gated: `upstreamFirstSemanticOutputMs` requires actual +non-empty text or reasoning output deltas rather than pre-populated context state. Inspected stream +provenance (`inspectedSource`: `upstream` vs `relay`) explicitly isolates bridge/adapter translation +failures (such as buffer or schema limits) into `relay` / `relay_transform` attribution rather than +misclassifying them as upstream failures. Diagnostic `transportPhase` and `terminalSource` fields are +strictly validated, sanitized of credential shapes, and preserved across process restarts. `src/usage/summary.ts` turns that file into the `/api/usage` shape — totals, daily zero-filled grid, model and provider breakdowns, and `measured / reported / unreported / unsupported / estimated` counts. A Codex-surface response also includes an `accounts` breakdown keyed by the stable non-PII