From 17eb2c61ec202a5453463ef9be84c8fa8ffc7dee Mon Sep 17 00:00:00 2001 From: chilung Date: Sat, 22 Aug 2026 08:48:06 +0000 Subject: [PATCH 1/7] 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 7bca650964..a40342c3b4 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -13,6 +13,30 @@ import { CODEX_ACCOUNT_LOG_LABEL_RE } from "../codex/account-label"; export type UsageStatus = "reported" | "unreported" | "unsupported" | "estimated"; export type CodexUsageAccountLogLabel = "main" | `p${string}`; +/** + * 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"; + export function isCodexUsageAccountLogLabel(value: unknown): value is CodexUsageAccountLogLabel { return value === "main" || (typeof value === "string" && CODEX_ACCOUNT_LOG_LABEL_RE.test(value)); } @@ -65,6 +89,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 { @@ -118,6 +146,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 @@ -238,6 +274,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); @@ -392,6 +460,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 } + : {}), }; } @@ -504,6 +581,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 e22c7e84b3..7b155f6b36 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -894,4 +894,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 6e3bb5b36bd0cdcbf10ad9151fcfa5206ad66e27 Mon Sep 17 00:00:00 2001 From: chilung Date: Sat, 22 Aug 2026 10:39:48 +0000 Subject: [PATCH 2/7] 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 a40342c3b4..3c447b53ca 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -460,9 +460,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 } : {}), @@ -581,9 +579,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 fcbeb2bbfa5040bd34f3d38932b1af2903be8550 Mon Sep 17 00:00:00 2001 From: chilung Date: Sat, 22 Aug 2026 11:44:08 +0000 Subject: [PATCH 3/7] 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 3c447b53ca..e5897fe378 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -36,6 +36,8 @@ export type FailureStage = | "downstream_write" | "client_cancel" | "terminal_delivery"; +export type TransportPhase = "pre_headers" | "mid_stream" | "terminal_sse"; +export type TerminalSource = "upstream" | "synthetic"; export function isCodexUsageAccountLogLabel(value: unknown): value is CodexUsageAccountLogLabel { return value === "main" || (typeof value === "string" && CODEX_ACCOUNT_LOG_LABEL_RE.test(value)); @@ -93,6 +95,8 @@ export interface PersistedUsageAttempt { streamTimeline?: StreamTimeline; failureSide?: FailureSide; failureStage?: FailureStage; + transportPhase?: TransportPhase; + terminalSource?: TerminalSource; } export interface PersistedUsageEntry { @@ -152,8 +156,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 @@ -286,6 +290,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; @@ -467,6 +480,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 } + : {}), }; } @@ -586,11 +605,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 7b155f6b36..5d143add96 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -937,4 +937,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 0ca137ca59f1b61a2e5aafffcd7a9756306b141d Mon Sep 17 00:00:00 2001 From: chilung Date: Sun, 23 Aug 2026 10:40:53 +0000 Subject: [PATCH 4/7] 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 bc56319c98..6aafa20394 100644 --- a/src/server/request-log.ts +++ b/src/server/request-log.ts @@ -27,8 +27,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 { @@ -48,6 +53,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. */ @@ -117,8 +124,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; } @@ -174,9 +184,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; } @@ -288,6 +301,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 } : {}), }; } @@ -405,6 +423,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 { @@ -437,6 +460,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; @@ -905,6 +952,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) @@ -931,6 +979,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, @@ -999,6 +1052,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 5d143add96..1936660357 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -982,4 +982,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 ee77d7220c9146f452cfbfa19e7bf46cdb2d4118 Mon Sep 17 00:00:00 2001 From: chilung Date: Mon, 24 Aug 2026 04:33:01 +0000 Subject: [PATCH 5/7] 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 6aafa20394..e7f0cc1556 100644 --- a/src/server/request-log.ts +++ b/src/server/request-log.ts @@ -979,7 +979,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 1936660357..a3547779d7 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -1064,4 +1064,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 c530fdbfa366e523aa9cac2ece4a27a451732aa4 Mon Sep 17 00:00:00 2001 From: chilung Date: Mon, 24 Aug 2026 15:04:13 +0000 Subject: [PATCH 6/7] 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 f59b4d0ae3..4b677a1bb6 100644 --- a/src/server/index.ts +++ b/src/server/index.ts @@ -1217,6 +1217,7 @@ export function startServer(port?: number, deps: StartServerDeps = {}): Server { @@ -1330,6 +1333,7 @@ export function startServer(port?: number, deps: StartServerDeps = {}): Server { @@ -1499,6 +1506,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 e7f0cc1556..e15ff4aa4d 100644 --- a/src/server/request-log.ts +++ b/src/server/request-log.ts @@ -22,6 +22,7 @@ import { isKnownUsageSurface, isCodexUsageAccountLogLabel, isValidReasoningWireValue, + normalizeStreamDiagnostics, readRecentUsageEntries, usageForFinalLog, usageStatusForFinalLog, @@ -362,10 +363,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 e66e85bd6f..a5d9ddfd99 100644 --- a/src/server/responses/core.ts +++ b/src/server/responses/core.ts @@ -1894,6 +1894,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 } : {}), }; @@ -3919,6 +3920,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 e5897fe378..0a511f9ff1 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -320,6 +320,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); } @@ -473,19 +504,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), }; } @@ -598,19 +617,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 4e57325112..6767073ed1 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, rmSync } from "node:fs"; +import { mkdtempSync, readFileSync, rmSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; @@ -355,6 +356,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 2cb06dd659bdc9323633f1a35e99d04f80ad9311 Mon Sep 17 00:00:00 2001 From: chilung Date: Mon, 24 Aug 2026 17:41:38 +0000 Subject: [PATCH 7/7] 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 0a511f9ff1..3bebba1a46 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -36,8 +36,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; export function isCodexUsageAccountLogLabel(value: unknown): value is CodexUsageAccountLogLabel { return value === "main" || (typeof value === "string" && CODEX_ACCOUNT_LOG_LABEL_RE.test(value)); @@ -290,15 +292,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; @@ -334,6 +344,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) @@ -342,12 +354,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 a3547779d7..667d0cdb28 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -938,7 +938,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(), @@ -975,12 +975,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)", () => {