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, logCtx?: RequestLogContext, onFirstOutput?: () => void, + inspectedSource?: "upstream" | "relay", ): ReadableStream { const reader = body.getReader(); let terminalReported = false; @@ -601,6 +603,7 @@ export function trackSseForRequestLog( onTerminal: reportTerminal, logCtx, onFirstOutput, + inspectedSource, }); return new ReadableStream({ @@ -645,7 +648,9 @@ 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() ?? ""; if (isUsageDebugEnabled() && !logCtx.usageDebugContentType && contentType) { logCtx.usageDebugContentType = contentType; @@ -707,6 +712,7 @@ export function responseWithDeferredRequestLog( }, logCtx, () => recordFirstOutput(logCtx, start), + inspectedSource, ); return new Response(body, { status: response.status, @@ -839,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 }; @@ -1021,6 +1032,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 && firstOutputFromParsed(parsed)) { + noteStreamTimelineEvent(handlers.logCtx, "upstreamFirstSemanticOutputMs"); + } const status = terminalStatusFromParsed(parsed); const policyTerminal = status === "failed" && isPolicyRewriteType(parsed) @@ -1030,8 +1044,15 @@ 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 = 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); } finally { @@ -1158,7 +1179,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; @@ -1184,6 +1209,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; @@ -1349,6 +1375,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, @@ -1367,7 +1396,12 @@ export function consumeForInspection( onCancel, onCleanEof: () => { if (!inspector.reported()) { - if (logCtx) logCtx.terminalSource = "synthetic"; + if (logCtx) { + logCtx.transportPhase = "mid_stream"; + logCtx.terminalSource = "synthetic"; + logCtx.failureSide = "upstream"; + logCtx.failureStage = "upstream_read"; + } onTerminal("incomplete"); } }, @@ -1379,6 +1413,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). @@ -1399,6 +1435,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 2c8d3e179c..c799793569 100644 --- a/src/server/request-log.ts +++ b/src/server/request-log.ts @@ -24,13 +24,19 @@ import { isKnownUsageSurface, isCodexUsageAccountLogLabel, isValidReasoningWireValue, + normalizeStreamDiagnostics, readRecentUsageEntries, usageForFinalLog, 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 +56,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 +142,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 +209,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 +326,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 } : {}), }; } @@ -368,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 { @@ -430,6 +468,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 +504,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 +993,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 +1022,13 @@ 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.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 +1102,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/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 7e056f97b0..1fb6fc2b9e 100644 --- a/src/usage/log.ts +++ b/src/usage/log.ts @@ -21,6 +21,34 @@ 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"; +/** 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 * it here is what stops six separate call sites from silently dropping an `o`-label -- including @@ -94,6 +122,12 @@ 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; + transportPhase?: TransportPhase; + terminalSource?: TerminalSource; } export interface PersistedUsageEntry { @@ -148,6 +182,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?: 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 @@ -281,6 +323,84 @@ 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", +]); +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; + 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 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); + 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) + ? { failureSide: raw.failureSide as FailureSide } + : {}), + ...(typeof raw.failureStage === "string" && KNOWN_FAILURE_STAGES.has(raw.failureStage as FailureStage) + ? { failureStage: raw.failureStage as FailureStage } + : {}), + ...(transportPhase ? { transportPhase: transportPhase as TransportPhase } : {}), + ...(terminalSource ? { terminalSource: terminalSource as TerminalSource } : {}), + }; +} export function isLabRouteSubjectId(value: unknown): value is string { return typeof value === "string" && LAB_ROUTE_SUBJECT_ID_RE.test(value); @@ -439,6 +559,7 @@ function normalizeUsageAttempt(raw: unknown): PersistedUsageAttempt | null { : { reasoningWireValue: attempt.reasoningWireValue } : {}), ...(tierOutcome ? { tierOutcome } : {}), + ...normalizeStreamDiagnostics(attempt), }; } @@ -551,6 +672,7 @@ function normalizeUsageEntry(entry: PersistedUsageEntry): PersistedUsageEntry { ...(entry.terminalStatus ? { terminalStatus: entry.terminalStatus } : {}), ...(entry.closeReason ? { closeReason: entry.closeReason } : {}), ...(entry.upstreamError ? { upstreamError: entry.upstreamError } : {}), + ...normalizeStreamDiagnostics(entry), ...(routeDecision ? { routeDecision } : {}), }; } 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 diff --git a/tests/request-log.test.ts b/tests/request-log.test.ts index bbf5b1071f..7a546f5989 100644 --- a/tests/request-log.test.ts +++ b/tests/request-log.test.ts @@ -26,14 +26,17 @@ 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, 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 +359,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); @@ -1596,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", () => { diff --git a/tests/usage-log.test.ts b/tests/usage-log.test.ts index cc96db4921..db72861741 100644 --- a/tests/usage-log.test.ts +++ b/tests/usage-log.test.ts @@ -913,4 +913,272 @@ 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"); + }); + + test("preserves bounded diagnostics but drops invalid closed attribution (#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).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).toBe("bogus_phase"); + expect(row?.attempts?.[0].terminalSource).toBe("bogus_source"); + }); + + 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); + }); + + 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); + }); });