From 2c193636c82ea15e1c305b415d01b9ea8d0e6f0a Mon Sep 17 00:00:00 2001 From: Randi <55005611+rdself@users.noreply.github.com> Date: Tue, 16 Jun 2026 08:12:31 -0400 Subject: [PATCH] Capture actual upstream provider requests (#3941) Integrated into release/v3.8.27 --- config/quality/file-size-baseline.json | 9 +- open-sse/executors/antigravity.ts | 6 + open-sse/executors/base.ts | 5 +- open-sse/executors/bedrock.ts | 10 +- open-sse/executors/codex.ts | 2 + open-sse/handlers/chatCore.ts | 104 +++++---- open-sse/utils/providerRequestLogging.ts | 243 ++++++++++++++++++++ src/lib/usage/pendingRequestScope.ts | 26 +++ src/lib/usage/usageHistory.ts | 9 +- tests/unit/executor-antigravity.test.ts | 37 ++- tests/unit/executor-bedrock.test.ts | 31 ++- tests/unit/executor-codex.test.ts | 57 +++++ tests/unit/executor-default-base.test.ts | 68 ++++++ tests/unit/provider-request-logging.test.ts | 140 +++++++++++ tests/unit/request-logger-endpoints.test.ts | 46 ++++ 15 files changed, 723 insertions(+), 70 deletions(-) create mode 100644 open-sse/utils/providerRequestLogging.ts create mode 100644 src/lib/usage/pendingRequestScope.ts create mode 100644 tests/unit/provider-request-logging.test.ts diff --git a/config/quality/file-size-baseline.json b/config/quality/file-size-baseline.json index 558cd45551..5054b664f1 100644 --- a/config/quality/file-size-baseline.json +++ b/config/quality/file-size-baseline.json @@ -19,11 +19,11 @@ "cap": 800, "frozen": { "open-sse/config/providerRegistry.ts": 4731, - "open-sse/executors/antigravity.ts": 1649, + "open-sse/executors/antigravity.ts": 1664, "open-sse/executors/base.ts": 1218, "open-sse/executors/chatgpt-web.ts": 2870, "open-sse/executors/claude-web.ts": 1057, - "open-sse/executors/codex.ts": 1439, + "open-sse/executors/codex.ts": 1447, "open-sse/executors/cursor.ts": 1391, "open-sse/executors/deepseek-web.ts": 1117, "open-sse/executors/duckduckgo-web.ts": 917, @@ -31,7 +31,7 @@ "open-sse/executors/muse-spark-web.ts": 1284, "open-sse/executors/perplexity-web.ts": 939, "open-sse/handlers/audioSpeech.ts": 965, - "open-sse/handlers/chatCore.ts": 5823, + "open-sse/handlers/chatCore.ts": 5830, "open-sse/handlers/imageGeneration.ts": 3777, "open-sse/handlers/responseSanitizer.ts": 1103, "open-sse/handlers/search.ts": 1442, @@ -147,5 +147,6 @@ "_rebaseline_2026_06_14_r3_3836_kiro_discovery": "PR #3836 own growth: models/route.ts 2426→2487 (+61 = kiro live per-account discovery branch wiring fetchKiroAvailableModels into the existing cache/auto-fetch/fallback discovery flow). Structural shrink of this route tracked in #3789.", "_rebaseline_2026_06_14_2997_disable_cooling": "Re-baseline #2997 (per-connection disable-cooling): EditConnectionModal.tsx 1171→1174 (+toggle UI) + auth.ts 2207→2216 (honor de disableCooling no markAccountUnavailable, pós-prettier). Lógica coesa; não-extraível. (combo.ts/RequestLoggerV2 drift já documentado em _r3_3835.)", "_rebaseline_2026_06_14_r3_3848_compression": "PR #3848 own growth: chatCore.ts 5808→5811 (+3 = compression engine pipeline hooks). Also carries inherited release/v3.8.25 drift not touched by this PR: models/route.ts 2487→2489 (+2, post-#3836/prettier). Updating the frozen values restores Fast Quality Gates on the current base.", - "_rebaseline_2026_06_14_3861_gitlab_duo": "PR #3861 own growth: oauth/[provider]/[action]/route.ts 903→916 (+13 = gitlab-duo authorize guard mirroring the existing qoder guard — returns a clear 'register an OAuth app + set GITLAB_DUO_OAUTH_CLIENT_ID' message instead of letting buildAuthUrl's throw become an opaque 500). Cohesive with the qoder branch right above it; not separately extractable." + "_rebaseline_2026_06_14_3861_gitlab_duo": "PR #3861 own growth: oauth/[provider]/[action]/route.ts 903→916 (+13 = gitlab-duo authorize guard mirroring the existing qoder guard — returns a clear 'register an OAuth app + set GITLAB_DUO_OAUTH_CLIENT_ID' message instead of letting buildAuthUrl's throw become an opaque 500). Cohesive with the qoder branch right above it; not separately extractable.", + "_rebaseline_2026_06_15_3941_provider_request_capture": "PR #3941 own growth: chatCore.ts 5823->5830 (+7 by check-file-size counting = run executor attempts inside the unified provider request capture scope), antigravity.ts 1649->1664 (+15) and codex.ts 1439->1447 (+8) = bridge hand-written upstream transports that bypass normal fetch/BaseExecutor capture. Cohesive logging-fidelity refactor; not extractable without hiding the actual dispatch boundary." } diff --git a/open-sse/executors/antigravity.ts b/open-sse/executors/antigravity.ts index eb2cdc4db5..b44fbe7fd0 100644 --- a/open-sse/executors/antigravity.ts +++ b/open-sse/executors/antigravity.ts @@ -54,6 +54,7 @@ import { getAntigravityEnvelopeUserAgent, getAntigravitySessionId, } from "../services/antigravityIdentity.ts"; +import * as prl from "../utils/providerRequestLogging.ts"; const MAX_RETRY_AFTER_MS = 60_000; const LONG_RETRY_THRESHOLD_MS = 60_000; @@ -1184,6 +1185,8 @@ export class AntigravityExecutor extends BaseExecutor { transformedBody ); let finalHeaders = serializedRequest.headers; + const capture = (h: Record, s: string) => + prl.captureCurrentProviderBody(url, h, s, log); const clientProfile = applyAntigravityClientProfileHeaders( finalHeaders, credentials, @@ -1220,6 +1223,7 @@ export class AntigravityExecutor extends BaseExecutor { ); } + await capture(finalHeaders, serializedRequest.bodyString); let response = await fetchWithReadinessTimeout(url, { method: "POST", headers: finalHeaders, @@ -1232,6 +1236,7 @@ export class AntigravityExecutor extends BaseExecutor { const retryHeaders = { ...finalHeaders }; removeHeaderCaseInsensitive(retryHeaders, "x-goog-user-project"); log?.debug?.("RETRY", "403 with x-goog-user-project, retrying once without it"); + await capture(retryHeaders, serializedRequest.bodyString); response = await fetchWithReadinessTimeout(url, { method: "POST", headers: retryHeaders, @@ -1314,6 +1319,7 @@ export class AntigravityExecutor extends BaseExecutor { ); const finalCreditsHeaders = serializedCreditsRequest.headers; try { + await capture(finalCreditsHeaders, serializedCreditsRequest.bodyString); const creditsResp = await fetchWithReadinessTimeout(url, { method: "POST", headers: finalCreditsHeaders, diff --git a/open-sse/executors/base.ts b/open-sse/executors/base.ts index a82af80579..77690d30ed 100644 --- a/open-sse/executors/base.ts +++ b/open-sse/executors/base.ts @@ -34,6 +34,7 @@ import { obfuscateInBody } from "../services/claudeCodeObfuscation.ts"; import { sanitizeClaudeToolSchemas } from "../translator/helpers/schemaCoercion.ts"; import { sanitizeResponsesInputItems } from "../services/responsesInputSanitizer.ts"; import { applySystemTransformPipeline, PROVIDER_CLAUDE } from "../services/systemTransforms.ts"; +import * as prl from "../utils/providerRequestLogging.ts"; import { fixToolPairs, fixToolAdjacency, @@ -1144,7 +1145,7 @@ export class BaseExecutor { } mergeUpstreamExtraHeaders(finalHeaders, upstreamExtraHeaders); - + const serializedBody = prl.parseBody(bodyString); const fetchOptions: RequestInit = { method: "POST", headers: finalHeaders, @@ -1190,7 +1191,7 @@ export class BaseExecutor { continue; } - return { response, url, headers: finalHeaders, transformedBody }; + return { response, url, headers: finalHeaders, transformedBody: serializedBody }; } catch (error) { // Distinguish timeout errors from other abort errors const err = error instanceof Error ? error : new Error(String(error)); diff --git a/open-sse/executors/bedrock.ts b/open-sse/executors/bedrock.ts index 43a38cf567..740e3e8bd5 100644 --- a/open-sse/executors/bedrock.ts +++ b/open-sse/executors/bedrock.ts @@ -9,6 +9,7 @@ import { randomUUID } from "node:crypto"; import { BaseExecutor } from "./base.ts"; import { PROVIDERS } from "../config/constants.ts"; import { buildBedrockNativeConverseUrl, resolveBedrockRegion } from "../config/bedrock.ts"; +import * as prl from "../utils/providerRequestLogging.ts"; const encoder = new TextEncoder(); @@ -631,7 +632,7 @@ export class BedrockExecutor extends BaseExecutor { }); } - async execute({ model, body, stream, credentials, signal }) { + async execute({ model, body, stream, credentials, signal, log }) { const url = this.buildUrl(model, stream, 0, credentials); const headers = this.buildHeaders(credentials); @@ -661,6 +662,13 @@ export class BedrockExecutor extends BaseExecutor { try { const client = this.createClient(credentials); + await prl.captureCurrentProviderRequest( + url, + headers, + transformedBody, + JSON.stringify(transformedBody), + log + ); if (stream) { const output = await client.send(new ConverseStreamCommand(transformedBody), { abortSignal: signal || undefined, diff --git a/open-sse/executors/codex.ts b/open-sse/executors/codex.ts index 93fc326da8..4581e866e5 100644 --- a/open-sse/executors/codex.ts +++ b/open-sse/executors/codex.ts @@ -28,6 +28,7 @@ import { getAccessToken } from "../services/tokenRefresh.ts"; import { sanitizeResponsesInputItems } from "../services/responsesInputSanitizer.ts"; import { getThinkingBudgetConfig, ThinkingMode } from "../services/thinkingBudget.ts"; import { CORS_HEADERS } from "../utils/cors.ts"; +import * as prl from "../utils/providerRequestLogging.ts"; import { createRequire } from "module"; // ─── wreq-js lazy loader ─────────────────────────────────────────────────── @@ -1000,6 +1001,7 @@ export class CodexExecutor extends BaseExecutor { finishStream({ reason: "upstream_closed", closeSocket: false }); }; if (!closed) { + await prl.captureCurrentProviderBody(url, headers, bodyString, nextInput.log); ws.send(bodyString); } } catch (error) { diff --git a/open-sse/handlers/chatCore.ts b/open-sse/handlers/chatCore.ts index 25488919fa..ed6dd48bc4 100644 --- a/open-sse/handlers/chatCore.ts +++ b/open-sse/handlers/chatCore.ts @@ -29,6 +29,7 @@ import { runWithOnPersist, } from "../services/tokenRefresh.ts"; import { createRequestLogger } from "../utils/requestLogger.ts"; +import { createPreparedRequestLogger, runWithCapture } from "../utils/providerRequestLogging.ts"; import { applyResponsesPreviousResponseIdPolicy } from "../utils/responsesStatePolicy.ts"; import { getModelTargetFormat, @@ -98,11 +99,10 @@ import { handleBypassRequest } from "../utils/bypassHandler.ts"; import { saveRequestUsage, trackPendingRequest, - updatePendingRequest, - finalizePendingRequest, appendRequestLog, saveCallLog, } from "@/lib/usageDb"; +import { finalizePendingScope, updatePendingScope } from "@/lib/usage/pendingRequestScope"; import { formatUsageLog, getLoggedInputTokens, @@ -2311,7 +2311,8 @@ export async function handleChatCore({ provider: provider || undefined, connectionId: connectionId || credentials?.connectionId || undefined, }); - + const pendingScope = { id: pendingRequestId, model, provider, connectionId: pendingConnId }; + const providerRequestCapture = createPreparedRequestLogger(reqLogger, pendingScope); // 0. Log client raw request (before format conversion) if (clientRawRequest) { reqLogger.logClientRawRequest( @@ -3795,7 +3796,7 @@ export async function handleChatCore({ } } - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { providerRequest: bodyToSend, stage: "payload_prepared", }); @@ -3805,7 +3806,7 @@ export async function handleChatCore({ max: accountSemaphoreMaxConcurrency, }); if (accountSemaphoreKey && accountSemaphoreMaxConcurrency != null) { - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { stage: "waiting_account_slot", }); } @@ -3817,7 +3818,7 @@ export async function handleChatCore({ }) : () => {}; trace("post_semaphore"); - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { stage: "waiting_rate_limit", }); @@ -3829,7 +3830,7 @@ export async function handleChatCore({ modelToCall, async () => { trace("inside_rate_limit"); - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { stage: "rate_limit_slot_acquired", }); let attempts = 0; @@ -3853,7 +3854,7 @@ export async function handleChatCore({ while (attempts < maxAttempts) { trace("pre_executor", { attempt: attempts }); - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { stage: "sending_to_provider", }); const execCreds = getExecutionCredentials(); @@ -3864,19 +3865,24 @@ export async function handleChatCore({ signal: streamController.signal, log, execute: (signal) => - executor.execute({ - model: modelToCall, - body: bodyToSend, - stream: upstreamStream, - credentials: execCreds, - signal, - log, - extendedContext, - upstreamExtraHeaders: buildUpstreamHeadersForExecute(modelToCall), - clientHeaders: buildExecutorClientHeaders(clientRawRequest?.headers, userAgent), - onCredentialsRefreshed, - skipUpstreamRetry, - }), + runWithCapture(providerRequestCapture, () => + executor.execute({ + model: modelToCall, + body: bodyToSend, + stream: upstreamStream, + credentials: execCreds, + signal, + log, + extendedContext, + upstreamExtraHeaders: buildUpstreamHeadersForExecute(modelToCall), + clientHeaders: buildExecutorClientHeaders( + clientRawRequest?.headers, + userAgent + ), + onCredentialsRefreshed, + skipUpstreamRetry, + }) + ), }); const res = normalizeExecutorResult(rawExecutorResult); trace("post_executor", { status: res?.response?.status }); @@ -3886,7 +3892,7 @@ export async function handleChatCore({ incrementRequestCount(modelToCall); } - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { stage: "provider_response_started", }); @@ -4117,10 +4123,9 @@ export async function handleChatCore({ } : translatedBody; - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { providerRequest: registeredProviderRequest, }); - // T5: track which models we've tried for intra-family fallback const triedModels = new Set([effectiveModel]); let currentModel = effectiveModel; @@ -4183,7 +4188,7 @@ export async function handleChatCore({ providerResponse = result.response; providerUrl = result.url; providerHeaders = result.headers; - finalBody = result.transformedBody; + finalBody = providerRequestCapture.body(result.transformedBody); effectiveServiceTier = resolveEffectiveServiceTier(finalBody); claudePromptCacheLogMeta = buildClaudePromptCacheLogMeta( targetFormat, @@ -4194,12 +4199,11 @@ export async function handleChatCore({ // Log target request (final request to provider) reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { providerRequest: finalBody, providerUrl, stage: "provider_response_started", }); - // Update rate limiter from response headers (learn limits dynamically) updateFromHeaders( provider, @@ -4409,27 +4413,29 @@ export async function handleChatCore({ // stay aligned if this block ever runs after a path that mutates body.model (e.g. fallback). try { const retryModelId = String(translatedBody.model || effectiveModel); - const retryResult = await executor.execute({ - model: retryModelId, - body: translatedBody, - stream: upstreamStream, - credentials: getExecutionCredentials(), - signal: streamController.signal, - log, - extendedContext, - upstreamExtraHeaders: buildUpstreamHeadersForExecute(retryModelId), - clientHeaders: buildExecutorClientHeaders(clientRawRequest?.headers, userAgent), - onCredentialsRefreshed, - skipUpstreamRetry: isCombo, - }); + const retryResult = await runWithCapture(providerRequestCapture, () => + executor.execute({ + model: retryModelId, + body: translatedBody, + stream: upstreamStream, + credentials: getExecutionCredentials(), + signal: streamController.signal, + log, + extendedContext, + upstreamExtraHeaders: buildUpstreamHeadersForExecute(retryModelId), + clientHeaders: buildExecutorClientHeaders(clientRawRequest?.headers, userAgent), + onCredentialsRefreshed, + skipUpstreamRetry: isCombo, + }) + ); if (retryResult.response.ok) { providerResponse = retryResult.response; providerUrl = retryResult.url; providerHeaders = new Headers(retryResult.headers || {}); - finalBody = retryResult.transformedBody; + finalBody = providerRequestCapture.body(retryResult.transformedBody); reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { providerRequest: finalBody, providerUrl, stage: "provider_response_started", @@ -4689,9 +4695,9 @@ export async function handleChatCore({ providerResponse = fallbackResult.response; providerUrl = fallbackResult.url; providerHeaders = fallbackResult.headers; - finalBody = fallbackResult.transformedBody; + finalBody = providerRequestCapture.body(fallbackResult.transformedBody); reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { providerRequest: finalBody, providerUrl, stage: "provider_response_started", @@ -4776,9 +4782,9 @@ export async function handleChatCore({ providerResponse = fallbackResult.response; providerUrl = fallbackResult.url; providerHeaders = fallbackResult.headers; - finalBody = fallbackResult.transformedBody; + finalBody = providerRequestCapture.body(fallbackResult.transformedBody); reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); - updatePendingRequest(model, provider, connectionId, { + updatePendingScope(pendingScope, { providerRequest: finalBody, providerUrl, stage: "provider_response_started", @@ -4999,7 +5005,7 @@ export async function handleChatCore({ responseBody = fallbackRaw ? JSON.parse(fallbackRaw) : {}; providerUrl = fallbackResult.url; providerHeaders = fallbackResult.headers; - finalBody = fallbackResult.transformedBody; + finalBody = providerRequestCapture.body(fallbackResult.transformedBody); reqLogger.logTargetRequest(providerUrl, providerHeaders, finalBody); log?.info?.( "EMPTY_CONTENT_FALLBACK", @@ -5264,7 +5270,7 @@ export async function handleChatCore({ "GUARDRAIL", `Response blocked by ${postCallGuardrails.guardrail || "guardrail"}: ${guardrailMessage}` ); - finalizePendingRequest(model, provider, connectionId, { + finalizePendingScope(pendingScope, { providerResponse: responseBody, clientResponse: translatedResponse, }); @@ -5349,7 +5355,7 @@ export async function handleChatCore({ } } - finalizePendingRequest(model, provider, connectionId, { + finalizePendingScope(pendingScope, { providerResponse: responseBody, clientResponse: translatedResponse, }); diff --git a/open-sse/utils/providerRequestLogging.ts b/open-sse/utils/providerRequestLogging.ts new file mode 100644 index 0000000000..e4e4998672 --- /dev/null +++ b/open-sse/utils/providerRequestLogging.ts @@ -0,0 +1,243 @@ +import { AsyncLocalStorage } from "node:async_hooks"; + +import { updatePendingScope, type PendingRequestScope } from "@/lib/usage/pendingRequestScope"; + +export type ProviderRequestPrepared = { + url: string; + headers: Record; + body: unknown; + bodyString: string; +}; + +export type Capture = { + capture: (request: ProviderRequestPrepared) => Promise | void; + body: (fallback: unknown) => unknown; + latest?: () => ProviderRequestPrepared | null; +}; + +type RequestLoggerLike = { + logTargetRequest: (url: unknown, headers: Record, body: unknown) => void; +}; + +type WarnLog = { + warn?: (tag: string, message: string) => void; +}; + +type FetchInput = Parameters[0]; +type FetchInit = Parameters[1]; + +type CaptureState = { + context: AsyncLocalStorage; + wrappedFetch: typeof fetch | null; + wrappedInnerFetch: typeof fetch | null; +}; + +const CAPTURE_STATE_KEY = Symbol.for("omniroute.providerRequestCapture.state"); + +function getCaptureState(): CaptureState { + const scopedGlobal = globalThis as typeof globalThis & { + [CAPTURE_STATE_KEY]?: CaptureState; + }; + + if (!scopedGlobal[CAPTURE_STATE_KEY]) { + scopedGlobal[CAPTURE_STATE_KEY] = { + context: new AsyncLocalStorage(), + wrappedFetch: null, + wrappedInnerFetch: null, + }; + } + return scopedGlobal[CAPTURE_STATE_KEY]; +} + +const captureState = getCaptureState(); +const BODY_METHODS = new Set(["POST", "PUT", "PATCH"]); +const AUTH_BODY_KEYS = new Set([ + "access_token", + "client_secret", + "grant_type", + "id_token", + "refresh_token", +]); +const REQUEST_BODY_KEYS = new Set([ + "conversationId", + "conversation_id", + "contents", + "input", + "messages", + "model", + "prompt", + "request", + "tools", + "userSelectedModel", +]); + +export function parseBody(bodyString: string): unknown { + try { + return JSON.parse(bodyString); + } catch { + return bodyString; + } +} + +async function capturePreparedRequest( + requestCapture: Capture | null | undefined, + url: string, + headers: Record, + body: unknown, + bodyString: string, + log?: WarnLog | null +) { + if (!requestCapture) return; + const latest = requestCapture.latest?.(); + if (latest?.url === url && latest.bodyString === bodyString) return; + + try { + await requestCapture.capture({ url, headers, body, bodyString }); + } catch (error) { + log?.warn?.( + "REQUEST_LOG", + `Provider request logging hook failed: ${error instanceof Error ? error.message : String(error)}` + ); + } +} + +export function captureCurrentProviderRequest( + url: string, + headers: Record, + body: unknown, + bodyString: string, + log?: WarnLog | null +) { + return capturePreparedRequest( + captureState.context.getStore(), + url, + headers, + body, + bodyString, + log + ); +} + +export function captureCurrentProviderBody( + url: string, + headers: Record, + bodyString: string, + log?: WarnLog | null +) { + return captureCurrentProviderRequest(url, headers, parseBody(bodyString), bodyString, log); +} + +export function runWithCapture(requestCapture: Capture, fn: () => Promise): Promise { + installFetchCapture(); + return captureState.context.run(requestCapture, fn); +} + +function installFetchCapture() { + if (globalThis.fetch === captureState.wrappedFetch) return; + + captureState.wrappedInnerFetch = globalThis.fetch.bind(globalThis); + captureState.wrappedFetch = (async (input: FetchInput, init?: FetchInit) => { + const activeCapture = captureState.context.getStore(); + if (activeCapture) { + await captureFetchRequest(activeCapture, input, init); + } + return captureState.wrappedInnerFetch!(input, init); + }) as typeof fetch; + globalThis.fetch = captureState.wrappedFetch; +} + +async function captureFetchRequest(requestCapture: Capture, input: FetchInput, init?: FetchInit) { + const method = getFetchMethod(input, init); + if (!BODY_METHODS.has(method)) return; + + const bodyString = bodyToString(init?.body); + if (!bodyString) return; + + const body = parseBody(bodyString); + if (!looksLikeProviderRequestBody(body)) return; + + await capturePreparedRequest( + requestCapture, + getFetchUrl(input), + getFetchHeaders(input, init), + body, + bodyString + ); +} + +function getFetchMethod(input: FetchInput, init?: FetchInit) { + const method = init?.method || (isRequest(input) ? input.method : "GET"); + return String(method || "GET").toUpperCase(); +} + +function getFetchUrl(input: FetchInput) { + if (typeof input === "string") return input; + if (input instanceof URL) return input.toString(); + if (isRequest(input)) return input.url; + return String(input); +} + +function getFetchHeaders(input: FetchInput, init?: FetchInit) { + const headers = new Headers(isRequest(input) ? input.headers : undefined); + if (init?.headers) { + new Headers(init.headers).forEach((value, key) => headers.set(key, value)); + } + + const result: Record = {}; + headers.forEach((value, key) => { + result[key] = value; + }); + return result; +} + +function bodyToString(body: BodyInit | null | undefined): string | null { + if (typeof body === "string") return body; + if (body instanceof URLSearchParams) return body.toString(); + if (body instanceof ArrayBuffer) return new TextDecoder().decode(body); + if (ArrayBuffer.isView(body)) { + return new TextDecoder().decode( + body.buffer.slice(body.byteOffset, body.byteOffset + body.byteLength) + ); + } + return null; +} + +function isRequest(input: FetchInput): input is Request { + return typeof Request !== "undefined" && input instanceof Request; +} + +function looksLikeProviderRequestBody(body: unknown) { + if (!body || typeof body !== "object" || Array.isArray(body)) return false; + const record = body as Record; + + if (Object.keys(record).some((key) => AUTH_BODY_KEYS.has(key))) return false; + if (Object.keys(record).some((key) => REQUEST_BODY_KEYS.has(key))) return true; + + return ( + typeof record.query === "string" && !!record.variables && typeof record.variables === "object" + ); +} + +export function createPreparedRequestLogger( + reqLogger: RequestLoggerLike, + scope: PendingRequestScope +): Capture { + let latest: ProviderRequestPrepared | null = null; + return { + capture(request) { + latest = request; + reqLogger.logTargetRequest(request.url, request.headers, request.body); + updatePendingScope(scope, { + providerRequest: request.body, + providerUrl: request.url, + stage: "sending_to_provider", + }); + }, + body(fallback) { + return latest?.body ?? fallback; + }, + latest() { + return latest; + }, + }; +} diff --git a/src/lib/usage/pendingRequestScope.ts b/src/lib/usage/pendingRequestScope.ts new file mode 100644 index 0000000000..60a4b77c5d --- /dev/null +++ b/src/lib/usage/pendingRequestScope.ts @@ -0,0 +1,26 @@ +import { + finalizePendingRequest, + finalizePendingRequestById, + updatePendingRequest, + updatePendingRequestById, + type PendingRequestMetadata, +} from "./usageHistory"; + +export type PendingRequestScope = { + id: string | null | undefined; + model: string; + provider: string; + connectionId: string | null; +}; + +export function updatePendingScope(scope: PendingRequestScope, metadata: PendingRequestMetadata) { + if (!updatePendingRequestById(scope.id || null, metadata)) { + updatePendingRequest(scope.model, scope.provider, scope.connectionId, metadata); + } +} + +export function finalizePendingScope(scope: PendingRequestScope, metadata: PendingRequestMetadata) { + if (!finalizePendingRequestById(scope.id, metadata)) { + finalizePendingRequest(scope.model, scope.provider, scope.connectionId, metadata); + } +} diff --git a/src/lib/usage/usageHistory.ts b/src/lib/usage/usageHistory.ts index d194fe03fe..5935c83ed0 100644 --- a/src/lib/usage/usageHistory.ts +++ b/src/lib/usage/usageHistory.ts @@ -26,7 +26,7 @@ import { } from "./tokenAccounting"; type JsonRecord = Record; -type PendingRequestMetadata = { +export type PendingRequestMetadata = { clientEndpoint?: string | null; clientRequest?: unknown; providerRequest?: unknown; @@ -308,6 +308,13 @@ export function updatePendingRequest( Object.assign(details[lastIdx], normalizePendingMetadata(metadata)); } +export function updatePendingRequestById(id: string | null, metadata: PendingRequestMetadata) { + const detail = id ? pendingById.get(id) : null; + if (!detail) return false; + Object.assign(detail, normalizePendingMetadata(metadata)); + return true; +} + /** * Update the first (oldest) pending request detail and then remove it. * Unlike updatePendingRequest which targets the last entry, this is designed diff --git a/tests/unit/executor-antigravity.test.ts b/tests/unit/executor-antigravity.test.ts index 442a44432c..036e7eee78 100644 --- a/tests/unit/executor-antigravity.test.ts +++ b/tests/unit/executor-antigravity.test.ts @@ -10,6 +10,7 @@ import { seedAntigravityVersionCache, } from "../../open-sse/services/antigravityVersion.ts"; import { clearAntigravityProjectCache } from "../../open-sse/services/antigravityProjectBootstrap.ts"; +import { runWithCapture } from "../../open-sse/utils/providerRequestLogging.ts"; type AntigravityTransformResult = Exclude< Awaited>, @@ -818,12 +819,18 @@ test("AntigravityExecutor.execute tags pre-response stalls with a fallbackable t test("AntigravityExecutor.execute applies CLI fingerprint when enabled", async () => { const executor = new AntigravityExecutor(); const originalFetch = globalThis.fetch; + let fetchStarted = false; + let fetchBody: Record | null = null; + let prepared: unknown = null; + let preparedBeforeFetch = false; seedAntigravityVersionCache("2026.04.17-test"); setCliCompatProviders(["antigravity"]); globalThis.fetch = async (_url, init) => { + fetchStarted = true; const headers = init?.headers as Record; const parsedBody = JSON.parse(String(init?.body)); + fetchBody = parsedBody; assert.equal(headers["User-Agent"], antigravityUserAgent("2026.04.17-test")); assert.equal(headers["x-client-name"], "antigravity"); @@ -849,17 +856,33 @@ test("AntigravityExecutor.execute applies CLI fingerprint when enabled", async ( }; try { + const requestCapture = { + capture(request) { + preparedBeforeFetch = !fetchStarted; + prepared = request.body; + }, + body(fallback) { + return prepared ?? fallback; + }, + latest() { + return null; + }, + }; const result = await withEnv("ANTIGRAVITY_CREDITS", "always", () => - executor.execute({ - model: "antigravity/gemini-2.5-flash", - body: { request: { contents: [] } }, - stream: false, - credentials: { accessToken: "token", projectId: "project-1" }, - log: { debug() {}, warn() {}, info() {} }, - }) + runWithCapture(requestCapture, () => + executor.execute({ + model: "antigravity/gemini-2.5-flash", + body: { request: { contents: [] } }, + stream: false, + credentials: { accessToken: "token", projectId: "project-1" }, + log: { debug() {}, warn() {}, info() {} }, + }) + ) ); assert.equal(result.response.status, 200); + assert.equal(preparedBeforeFetch, true); + assert.deepEqual(prepared, fetchBody); } finally { setCliCompatProviders([]); globalThis.fetch = originalFetch; diff --git a/tests/unit/executor-bedrock.test.ts b/tests/unit/executor-bedrock.test.ts index 816f823284..14fa822378 100644 --- a/tests/unit/executor-bedrock.test.ts +++ b/tests/unit/executor-bedrock.test.ts @@ -2,6 +2,7 @@ import test from "node:test"; import assert from "node:assert/strict"; import { BedrockExecutor, openAIToBedrockConverse } from "../../open-sse/executors/bedrock.ts"; +import { runWithCapture } from "../../open-sse/utils/providerRequestLogging.ts"; function credentials(region = "eu-west-2") { return { @@ -271,6 +272,8 @@ test("openAIToBedrockConverse removes tool uses whose results are not immediatel test("BedrockExecutor converts non-streaming Converse output to OpenAI chat completion JSON", async () => { const sent = []; + let prepared = null; + let preparedBeforeSend = false; const executor = new BedrockExecutor(() => ({ send: async (command) => { sent.push(command); @@ -282,15 +285,31 @@ test("BedrockExecutor converts non-streaming Converse output to OpenAI chat comp }, })); - const result = await executor.execute({ - model: "anthropic.claude-sonnet-4-6", - body: { messages: [{ role: "user", content: "Hi" }], max_tokens: 8 }, - stream: false, - credentials: credentials(), - }); + const requestCapture = { + capture(request) { + preparedBeforeSend = sent.length === 0; + prepared = request; + }, + body(fallback) { + return prepared?.body ?? fallback; + }, + latest() { + return prepared; + }, + }; + const result = await runWithCapture(requestCapture, () => + executor.execute({ + model: "anthropic.claude-sonnet-4-6", + body: { messages: [{ role: "user", content: "Hi" }], max_tokens: 8 }, + stream: false, + credentials: credentials(), + }) + ); assert.equal(sent[0].constructor.name, "ConverseCommand"); assert.equal(sent[0].input.modelId, "anthropic.claude-sonnet-4-6"); + assert.equal(preparedBeforeSend, true); + assert.deepEqual(prepared.body, sent[0].input); assert.equal(result.response.status, 200); const body = await result.response.json(); assert.equal(body.model, "anthropic.claude-sonnet-4-6"); diff --git a/tests/unit/executor-codex.test.ts b/tests/unit/executor-codex.test.ts index a1697f36f6..d92b134787 100644 --- a/tests/unit/executor-codex.test.ts +++ b/tests/unit/executor-codex.test.ts @@ -23,6 +23,7 @@ import { setThinkingBudgetConfig, ThinkingMode, } from "../../open-sse/services/thinkingBudget.ts"; +import { runWithCapture } from "../../open-sse/utils/providerRequestLogging.ts"; import { CODEX_CHAT_DEFAULT_INSTRUCTIONS } from "../../open-sse/config/codexInstructions.ts"; type MockCodexWebSocket = { @@ -912,6 +913,62 @@ test("CodexExecutor.execute falls back to HTTP when websocket transport is unava } }); +test("CodexExecutor.execute captures the exact websocket request body before send", async () => { + const executor = new CodexExecutor(); + let sent: string | null = null; + let sendStarted = false; + let prepared: unknown = null; + let preparedBeforeSend = false; + const ws: MockCodexWebSocket = { + send(data) { + sendStarted = true; + sent = data; + queueMicrotask(() => { + ws.onmessage?.({ + data: JSON.stringify({ type: "response.completed", response: { status: "completed" } }), + }); + }); + }, + close() {}, + onmessage: null, + onerror: null, + onclose: null, + }; + __setCodexWebSocketTransportForTesting(async () => ws); + + const requestCapture = { + capture(request) { + preparedBeforeSend = !sendStarted; + prepared = request.body; + }, + body(fallback) { + return prepared ?? fallback; + }, + latest() { + return null; + }, + }; + const result = await runWithCapture(requestCapture, () => + executor.execute({ + model: "gpt-5.5-xhigh", + body: { model: "gpt-5.5-xhigh", input: [{ role: "user", content: "hello" }] }, + stream: true, + credentials: { + accessToken: "codex-token", + providerSpecificData: { codexTransport: "websocket" }, + }, + }) + ); + await result.response.text(); + + assert.ok(sent); + const sentBody = JSON.parse(sent); + assert.equal(preparedBeforeSend, true); + assert.deepEqual(prepared, sentBody); + assert.equal(sentBody.type, "response.create"); + assert.equal(sentBody.model, "gpt-5.5"); +}); + test("CodexExecutor.execute adds CLI-like session identity headers without changing response flow", async () => { const executor = new CodexExecutor(); const originalFetch = globalThis.fetch; diff --git a/tests/unit/executor-default-base.test.ts b/tests/unit/executor-default-base.test.ts index 95657cf61e..75d8c5fe84 100644 --- a/tests/unit/executor-default-base.test.ts +++ b/tests/unit/executor-default-base.test.ts @@ -17,6 +17,7 @@ import { CLAUDE_CODE_COMPATIBLE_REDACT_THINKING_BETA, CONTEXT_1M_BETA_HEADER, } from "../../open-sse/services/claudeCodeCompatible.ts"; +import { runWithCapture } from "../../open-sse/utils/providerRequestLogging.ts"; class TestExecutor extends BaseExecutor { constructor(config = {}) { @@ -608,6 +609,73 @@ test("DefaultExecutor.execute uses CC-compatible connection defaults to append 1 assert.equal(calls[2].headers["anthropic-beta"], undefined); }); +test("DefaultExecutor.execute reports the exact serialized provider request before fetch", async () => { + const originalFetch = globalThis.fetch; + let fetchStarted = false; + let fetchBody: any = null; + let prepared: any = null; + let preparedBeforeFetch = false; + + globalThis.fetch = async (_url, init = {}) => { + fetchStarted = true; + fetchBody = JSON.parse(String(init.body || "{}")); + return new Response(JSON.stringify({ ok: true }), { + status: 200, + headers: { "Content-Type": "application/json" }, + }); + }; + + try { + const cc = new DefaultExecutor("anthropic-compatible-cc-test"); + const requestCapture = { + capture(request) { + preparedBeforeFetch = !fetchStarted; + prepared = request; + }, + body(fallback) { + return prepared?.body ?? fallback; + }, + latest() { + return prepared; + }, + }; + const result = await runWithCapture(requestCapture, () => + cc.execute({ + model: "claude-sonnet-4-6", + body: { + model: "claude-sonnet-4-6", + system: [ + { + type: "text", + text: "x-anthropic-billing-header: cc_version=1.0.0; cc_entrypoint=sdk-cli; cch=00000;", + }, + ], + messages: [{ role: "user", content: "hi" }], + max_tokens: 1, + reasoning_effort: "xhigh", + }, + stream: false, + credentials: { + apiKey: "cc-key", + providerSpecificData: { + ccSessionId: "session-1", + }, + }, + }) + ); + + assert.ok(prepared, "prepared request hook should fire before fetch"); + assert.equal(preparedBeforeFetch, true); + assert.deepEqual(prepared.body, fetchBody); + assert.deepEqual(result.transformedBody, fetchBody); + assert.equal(prepared.body.reasoning_effort, "high"); + assert.equal(fetchBody.reasoning_effort, "high"); + assert.match(JSON.stringify(fetchBody), /\bcch=(?!00000)[0-9a-f]{5};/); + } finally { + globalThis.fetch = originalFetch; + } +}); + test("DefaultExecutor.execute only injects adaptive thinking defaults for Claude models that support x-high effort", async () => { const originalFetch = globalThis.fetch; const requestBodies = []; diff --git a/tests/unit/provider-request-logging.test.ts b/tests/unit/provider-request-logging.test.ts new file mode 100644 index 0000000000..f0eeacb95d --- /dev/null +++ b/tests/unit/provider-request-logging.test.ts @@ -0,0 +1,140 @@ +import test from "node:test"; +import assert from "node:assert/strict"; + +import { + runWithCapture, + type Capture, + type ProviderRequestPrepared, +} from "../../open-sse/utils/providerRequestLogging.ts"; + +test("runWithCapture captures the actual JSON provider fetch body", async () => { + const originalFetch = globalThis.fetch; + const prepared: ProviderRequestPrepared[] = []; + const sentBodies: unknown[] = []; + const capture: Capture = { + capture(request) { + prepared.push(request); + }, + body(fallback) { + return prepared.at(-1)?.body ?? fallback; + }, + }; + + globalThis.fetch = async (_url, init = {}) => { + sentBodies.push(JSON.parse(String(init.body))); + return new Response(JSON.stringify({ ok: true }), { + status: 200, + headers: { "Content-Type": "application/json" }, + }); + }; + + try { + await runWithCapture(capture, () => + fetch("https://provider.example/v1/chat/completions", { + method: "POST", + headers: { Authorization: "Bearer provider-key" }, + body: JSON.stringify({ + model: "claude-sonnet-4-6", + messages: [{ role: "user", content: "hi" }], + reasoning_effort: "high", + }), + }) + ); + + assert.equal(prepared.length, 1); + assert.deepEqual(prepared[0].body, sentBodies[0]); + assert.equal(prepared[0].url, "https://provider.example/v1/chat/completions"); + assert.equal(prepared[0].headers.authorization, "Bearer provider-key"); + assert.equal((capture.body(null) as Record).reasoning_effort, "high"); + } finally { + globalThis.fetch = originalFetch; + } +}); + +test("runWithCapture ignores auth fetch bodies in the same executor scope", async () => { + const originalFetch = globalThis.fetch; + const prepared: ProviderRequestPrepared[] = []; + const capture: Capture = { + capture(request) { + prepared.push(request); + }, + body(fallback) { + return prepared.at(-1)?.body ?? fallback; + }, + }; + + globalThis.fetch = async () => + new Response(JSON.stringify({ ok: true }), { + status: 200, + headers: { "Content-Type": "application/json" }, + }); + + try { + await runWithCapture(capture, async () => { + await fetch("https://oauth.example/token", { + method: "POST", + body: JSON.stringify({ + grant_type: "refresh_token", + refresh_token: "refresh", + client_id: "client", + }), + }); + await fetch("https://provider.example/v1/chat/completions", { + method: "POST", + body: JSON.stringify({ + model: "gpt-5", + messages: [{ role: "user", content: "hi" }], + }), + }); + }); + + assert.equal(prepared.length, 1); + assert.equal((prepared[0].body as Record).model, "gpt-5"); + } finally { + globalThis.fetch = originalFetch; + } +}); + +test("runWithCapture does not duplicate an already prepared identical fetch", async () => { + const originalFetch = globalThis.fetch; + const prepared: ProviderRequestPrepared[] = []; + const body = { + model: "gpt-5", + messages: [{ role: "user", content: "hi" }], + }; + const bodyString = JSON.stringify(body); + const url = "https://provider.example/v1/chat/completions"; + let latest: ProviderRequestPrepared | null = null; + const capture: Capture = { + capture(request) { + latest = request; + prepared.push(request); + }, + body(fallback) { + return latest?.body ?? fallback; + }, + latest() { + return latest; + }, + }; + + globalThis.fetch = async () => + new Response(JSON.stringify({ ok: true }), { + status: 200, + headers: { "Content-Type": "application/json" }, + }); + + try { + await runWithCapture(capture, async () => { + await capture.capture({ url, headers: {}, body, bodyString }); + await fetch(url, { + method: "POST", + body: bodyString, + }); + }); + + assert.equal(prepared.length, 1); + } finally { + globalThis.fetch = originalFetch; + } +}); diff --git a/tests/unit/request-logger-endpoints.test.ts b/tests/unit/request-logger-endpoints.test.ts index e504523591..f016e63cd4 100644 --- a/tests/unit/request-logger-endpoints.test.ts +++ b/tests/unit/request-logger-endpoints.test.ts @@ -121,6 +121,52 @@ test("updatePendingRequest keeps pending detail API view in sync", () => { }); }); +test("updatePendingRequestById updates and finalizes the exact overlapping request", () => { + usageHistory.clearPendingRequests(); + const firstId = usageHistory.trackPendingRequest("claude-sonnet-4-6", "cc-test", "conn-1", true, { + providerRequest: { request: "first", reasoning_effort: "xhigh" }, + }); + const secondId = usageHistory.trackPendingRequest( + "claude-sonnet-4-6", + "cc-test", + "conn-1", + true, + { + providerRequest: { request: "second", reasoning_effort: "xhigh" }, + } + ); + assert.ok(firstId); + assert.ok(secondId); + + const updated = usageHistory.updatePendingRequestById(firstId, { + providerRequest: { request: "first", reasoning_effort: "high" }, + stage: "sending_to_provider", + }); + assert.equal(updated, true); + + const modelKey = "claude-sonnet-4-6 (cc-test)"; + const details = usageHistory.getPendingRequests().details["conn-1"]?.[modelKey]; + assert.equal(details?.length, 2); + assert.deepEqual(details?.[0]?.providerRequest, { + request: "first", + reasoning_effort: "high", + }); + assert.deepEqual(details?.[1]?.providerRequest, { + request: "second", + reasoning_effort: "xhigh", + }); + + const completed = usageHistory.finalizePendingRequestById(firstId, { + clientResponse: { ok: true }, + }); + assert.equal(completed, true); + assert.deepEqual(usageHistory.getCompletedDetails().get(firstId)?.providerRequest, { + request: "first", + reasoning_effort: "high", + }); + assert.equal(usageHistory.getPendingById().has(secondId), true); +}); + test("updatePendingRequestStreamChunks stores empty streamChunks object (not null)", () => { usageHistory.clearPendingRequests(); usageHistory.trackPendingRequest("gpt-4", "openai", "conn-1", true);