From ca00505f68932f01b82d48ef0fb9895d11588e85 Mon Sep 17 00:00:00 2001 From: chrarnoldus <12196001+chrarnoldus@users.noreply.github.com> Date: Fri, 24 Jul 2026 15:12:44 +0000 Subject: [PATCH 1/5] refactor(ai-gateway): capture request log body during response rewrite Combine handleRequestLogging and rewriteModelResponse so the event stream is only processed once. handleRequestLogging now returns a capture handle instead of reading a cloned response; rewriteModelResponse accepts the capture and records the upstream body while rewriting it. When request logging is enabled the rewrite pipeline now runs for everyone so the body is captured in the same pass. processUsage is unchanged. Co-authored-by: kiloconnect[bot] <240665456+kiloconnect[bot]@users.noreply.github.com> --- .../src/app/api/openrouter/[...path]/route.ts | 40 +++++-- .../lib/ai-gateway/handleRequestLogging.ts | 72 ++++++++---- apps/web/src/lib/rewriteModelResponse.test.ts | 105 ++++++++++++++++++ apps/web/src/lib/rewriteModelResponse.ts | 82 +++++++++++--- 4 files changed, 252 insertions(+), 47 deletions(-) diff --git a/apps/web/src/app/api/openrouter/[...path]/route.ts b/apps/web/src/app/api/openrouter/[...path]/route.ts index 0b1b913c2b..049c5262b4 100644 --- a/apps/web/src/app/api/openrouter/[...path]/route.ts +++ b/apps/web/src/app/api/openrouter/[...path]/route.ts @@ -957,8 +957,11 @@ export async function POST(request: NextRequest): Promise { + const { status, user, organization_id, session_id, vercel_request_id, provider, model, request } = + params; if (!(await isLoggingEnabledForUser(user, organization_id))) { - return; + return null; } + + let resolveCaptured: (result: CapturedResponseBody) => void = () => {}; + const captured = new Promise(resolve => { + resolveCaptured = resolve; + }); + let isSettled = false; + const settleOnce = (result: CapturedResponseBody) => { + if (!isSettled) { + isSettled = true; + resolveCaptured(result); + } + }; + after(async () => { - // Read the response body in its own try/catch: if it cannot be read (e.g. - // the stream errored mid-flight), still log the request without it. - let response: string | undefined; - let responseReadError: string | undefined; - try { - response = await clonedResponse.text(); - } catch (e) { - responseReadError = String(e).substring(0, 4000); + // Wait until the response pipeline has processed the response body. This + // resolves when the response stream completes (or fails), which happens + // before after() callbacks are awaited. + const result = await captured; + const response = 'text' in result ? result.text : undefined; + const responseReadError = 'readError' in result ? result.readError : undefined; + if (responseReadError !== undefined) { logExceptInTest( - `[handleRequestLogging] failed to read response body (user=${user?.id}, status=${clonedResponse.status}, model=${model}): ${responseReadError}` + `[handleRequestLogging] failed to read response body (user=${user?.id}, status=${status}, model=${model}): ${responseReadError}` ); } try { @@ -68,7 +87,7 @@ export async function handleRequestLogging(params: { organization_id: organization_id, session_id, vercel_request_id, - status_code: clonedResponse.status, + status_code: status, model, provider, request: request.body, @@ -83,8 +102,13 @@ export async function handleRequestLogging(params: { } catch (e) { const cause = e instanceof Error ? e.cause : undefined; logExceptInTest( - `[handleRequestLogging] failed to insert api_request_log (user=${user?.id}, status=${clonedResponse.status}, model=${model}) cause (truncated): ${String(cause).substring(0, 4000)} error (truncated): ${String(e).substring(0, 4000)}` + `[handleRequestLogging] failed to insert api_request_log (user=${user?.id}, status=${status}, model=${model}) cause (truncated): ${String(cause).substring(0, 4000)} error (truncated): ${String(e).substring(0, 4000)}` ); } }); + + return { + setBody: text => settleOnce({ text }), + setReadError: error => settleOnce({ readError: String(error).substring(0, 4000) }), + }; } diff --git a/apps/web/src/lib/rewriteModelResponse.test.ts b/apps/web/src/lib/rewriteModelResponse.test.ts index dc204fd350..64809da6e9 100644 --- a/apps/web/src/lib/rewriteModelResponse.test.ts +++ b/apps/web/src/lib/rewriteModelResponse.test.ts @@ -534,4 +534,109 @@ describe('rewriteModelResponse', () => { usage: {}, }); }); + + test('processes responses it would normally skip when a capture is provided', async () => { + const capture = makeCapture(); + const result = await rewriteModelResponse( + jsonResponse({ model: 'openai/gpt-5' }), + 'openai/gpt-5', + 'openrouter', + 'chat_completions', + '00000000-0000-0000-0000-000000000000', + capture + ); + + expect(result).not.toBeNull(); + expect(capture.setBody).toHaveBeenCalledWith(JSON.stringify({ model: 'openai/gpt-5' })); + }); +}); + +function makeCapture() { + return { setBody: jest.fn(), setReadError: jest.fn() }; +} + +describe('request log capture', () => { + test.each(rewriters)('%s: captures the raw JSON body', async (_name, rewrite) => { + const capture = makeCapture(); + const body = { model: 'upstream-model' }; + + const result = await rewrite(jsonResponse(body), true, capture); + + expect(result.status).toBe(200); + expect(capture.setBody).toHaveBeenCalledTimes(1); + expect(capture.setBody).toHaveBeenCalledWith(JSON.stringify(body)); + expect(capture.setReadError).not.toHaveBeenCalled(); + }); + + test.each(rewriters)('%s: captures the raw event stream', async (_name, rewrite) => { + const capture = makeCapture(); + const sseBody = + 'data: {"id":"gen-1","model":"upstream-model","choices":[]}\n\n' + 'data: [DONE]\n\n'; + + const result = await rewrite(sseResponse(sseBody), true, capture); + await readOutputStream(result); + + expect(capture.setBody).toHaveBeenCalledTimes(1); + expect(capture.setBody).toHaveBeenCalledWith(sseBody); + expect(capture.setReadError).not.toHaveBeenCalled(); + }); + + test.each(rewriters)('%s: captures an empty body when upstream has no body', async (_name, rewrite) => { + const capture = makeCapture(); + + const result = await rewrite( + new Response(null, { headers: { 'content-type': 'text/event-stream' } }), + true, + capture + ); + await readOutputStream(result); + + expect(capture.setBody).toHaveBeenCalledWith(''); + expect(capture.setReadError).not.toHaveBeenCalled(); + }); + + test.each(rewriters)('%s: records a read error when the stream fails', async (_name, rewrite) => { + const capture = makeCapture(); + + const result = await rewrite( + failingResponse( + 'text/event-stream', + 'ResponseAborted', + 'data: {"id":"gen-1","choices":[]}\n\n' + ), + true, + capture + ); + await readOutputStream(result); + + expect(capture.setReadError).toHaveBeenCalledTimes(1); + expect(capture.setBody).not.toHaveBeenCalled(); + }); + + test.each(rewriters)( + '%s: records a read error when a JSON body cannot be read', + async (_name, rewrite) => { + const capture = makeCapture(); + + const result = await rewrite(failingResponse('application/json', 'TimeoutError'), true, capture); + + expect(result.status).toBe(503); + expect(capture.setReadError).toHaveBeenCalledTimes(1); + expect(capture.setBody).not.toHaveBeenCalled(); + } + ); + + test('records a read error when the response stream is cancelled', async () => { + const capture = makeCapture(); + const upstream = new Response(new ReadableStream({ start() {} }), { + headers: { 'content-type': 'text/event-stream' }, + }); + + const result = await rewriteModelResponse_ChatCompletions(upstream, true, capture); + const reader = result.body?.getReader(); + await reader?.cancel(); + + expect(capture.setReadError).toHaveBeenCalled(); + expect(capture.setBody).not.toHaveBeenCalled(); + }); }); diff --git a/apps/web/src/lib/rewriteModelResponse.ts b/apps/web/src/lib/rewriteModelResponse.ts index c36c451bd2..88c468bd41 100644 --- a/apps/web/src/lib/rewriteModelResponse.ts +++ b/apps/web/src/lib/rewriteModelResponse.ts @@ -1,4 +1,5 @@ import { isKiloExclusiveFreeModel } from '@/lib/ai-gateway/models'; +import type { RequestLogCapture } from '@/lib/ai-gateway/handleRequestLogging'; import type { GatewayRequest } from '@/lib/ai-gateway/providers/openrouter/types'; import type { ProviderId } from '@/lib/ai-gateway/providers/types'; import { getOutputHeaders } from '@/lib/ai-gateway/llm-proxy-helpers'; @@ -61,7 +62,7 @@ function getResponseReadError(error: unknown): ResponseReadError | null { async function readResponseText( response: Response, headers: Headers -): Promise<{ text: string } | { errorResponse: NextResponse }> { +): Promise<{ text: string } | { error: unknown; errorResponse: NextResponse }> { try { return { text: await response.text() }; } catch (error) { @@ -71,6 +72,7 @@ async function readResponseText( } return { + error, errorResponse: NextResponse.json( { error: responseReadError.message, @@ -89,9 +91,13 @@ async function rewriteSseStream( controller: ReadableStreamDefaultController, doneReceived: () => boolean, serializeError: (error: ResponseReadError) => string, - onFinally: () => void + onFinally: () => void, + capture?: RequestLogCapture | null ) { const decoder = new TextDecoder(); + // Accumulate the raw upstream text for request logging while the stream is + // being processed anyway, so it doesn't have to be processed a second time. + const capturedChunks: string[] | null = capture ? [] : null; try { while (true) { const { done, value } = await reader.read(); @@ -103,17 +109,25 @@ async function rewriteSseStream( controller.enqueue('data: [DONE]\n\n'); } controller.close(); + if (capturedChunks) { + capturedChunks.push(decoder.decode()); + capture?.setBody(capturedChunks.join('')); + } return; } - parser.feed(decoder.decode(value, { stream: true })); + const chunk = decoder.decode(value, { stream: true }); + capturedChunks?.push(chunk); + parser.feed(chunk); } } catch (error) { const responseReadError = getResponseReadError(error); if (!responseReadError) { + capture?.setReadError(error); throw error; } errorExceptInTest('[rewriteModelResponse] emitting stream error event', responseReadError); + capture?.setReadError(error); controller.enqueue(serializeError(responseReadError)); controller.close(); } finally { @@ -135,7 +149,11 @@ function rewriteUsage(usage: OpenRouterUsage, removeCost: boolean) { } } -export async function rewriteModelResponse_ChatCompletions(response: Response, removeCost = true) { +export async function rewriteModelResponse_ChatCompletions( + response: Response, + removeCost = true, + capture?: RequestLogCapture | null +) { const headers = getOutputHeaders(response); if (headers.get('content-type')?.includes('application/json')) { @@ -143,8 +161,10 @@ export async function rewriteModelResponse_ChatCompletions(response: Response, r // disturbed or locked" errors that occur when `.clone().json()` fails. const textResult = await readResponseText(response, headers); if ('errorResponse' in textResult) { + capture?.setReadError(textResult.error); return textResult.errorResponse; } + capture?.setBody(textResult.text); const { text } = textResult; let json: OpenAI.ChatCompletion; try { @@ -174,6 +194,7 @@ export async function rewriteModelResponse_ChatCompletions(response: Response, r const reader = response.body?.getReader(); if (!reader) { controller.close(); + capture?.setBody(''); return; } @@ -237,9 +258,14 @@ export async function rewriteModelResponse_ChatCompletions(response: Response, r }, }) + '\n\n', - progress.stop + progress.stop, + capture ); }, + cancel() { + // The client disconnected before the stream completed. + capture?.setReadError(new Error('response stream was cancelled')); + }, }); return new NextResponse(stream, { @@ -274,14 +300,20 @@ function rewriteMessagesUsage(usage: MessagesApiUsage, removeCost: boolean) { } } -export async function rewriteModelResponse_Messages(response: Response, removeCost = true) { +export async function rewriteModelResponse_Messages( + response: Response, + removeCost = true, + capture?: RequestLogCapture | null +) { const headers = getOutputHeaders(response); if (headers.get('content-type')?.includes('application/json')) { const textResult = await readResponseText(response, headers); if ('errorResponse' in textResult) { + capture?.setReadError(textResult.error); return textResult.errorResponse; } + capture?.setBody(textResult.text); const { text } = textResult; let json: Anthropic.Messages.Message & { usage?: MessagesApiUsage }; try { @@ -311,6 +343,7 @@ export async function rewriteModelResponse_Messages(response: Response, removeCo const reader = response.body?.getReader(); if (!reader) { controller.close(); + capture?.setBody(''); return; } @@ -376,9 +409,14 @@ export async function rewriteModelResponse_Messages(response: Response, removeCo }, }) + '\n\n', - progress.stop + progress.stop, + capture ); }, + cancel() { + // The client disconnected before the stream completed. + capture?.setReadError(new Error('response stream was cancelled')); + }, }); return new NextResponse(stream, { @@ -394,14 +432,20 @@ type ResponsesApiEvent = { response?: OpenAI.Responses.Response & { usage?: OpenRouterUsage | null }; }; -export async function rewriteModelResponse_Responses(response: Response, removeCost = true) { +export async function rewriteModelResponse_Responses( + response: Response, + removeCost = true, + capture?: RequestLogCapture | null +) { const headers = getOutputHeaders(response); if (headers.get('content-type')?.includes('application/json')) { const textResult = await readResponseText(response, headers); if ('errorResponse' in textResult) { + capture?.setReadError(textResult.error); return textResult.errorResponse; } + capture?.setBody(textResult.text); const { text } = textResult; let json: OpenAI.Responses.Response & { usage?: OpenRouterUsage | null }; try { @@ -431,6 +475,7 @@ export async function rewriteModelResponse_Responses(response: Response, removeC const reader = response.body?.getReader(); if (!reader) { controller.close(); + capture?.setBody(''); return; } @@ -488,9 +533,14 @@ export async function rewriteModelResponse_Responses(response: Response, removeC }, }) + '\n\n', - progress.stop + progress.stop, + capture ); }, + cancel() { + // The client disconnected before the stream completed. + capture?.setReadError(new Error('response stream was cancelled')); + }, }); return new NextResponse(stream, { @@ -505,25 +555,29 @@ export async function rewriteModelResponse( model: string, providerId: ProviderId, kind: GatewayRequest['kind'], - organizationId?: string + organizationId?: string, + capture?: RequestLogCapture | null ): Promise { const isFreeModelRequiringCostRemoval = (providerId === 'openrouter' || providerId === 'vercel') && isKiloExclusiveFreeModel(model); - if (!isFreeModelRequiringCostRemoval && organizationId !== KILO_ORGANIZATION_ID) { + // When request logging is enabled the response has to be processed anyway + // so the body can be captured for the request log in a single pass, so the + // rewrite is not skipped in that case. + if (!isFreeModelRequiringCostRemoval && organizationId !== KILO_ORGANIZATION_ID && !capture) { console.debug('[rewriteModelResponse] skipping rewrite for %s', model); return null; } console.debug('[rewriteModelResponse] rewriting response for %s', model); if (kind === 'chat_completions') { - return rewriteModelResponse_ChatCompletions(response, isFreeModelRequiringCostRemoval); + return rewriteModelResponse_ChatCompletions(response, isFreeModelRequiringCostRemoval, capture); } if (kind === 'responses') { - return rewriteModelResponse_Responses(response, isFreeModelRequiringCostRemoval); + return rewriteModelResponse_Responses(response, isFreeModelRequiringCostRemoval, capture); } if (kind === 'messages') { - return rewriteModelResponse_Messages(response, isFreeModelRequiringCostRemoval); + return rewriteModelResponse_Messages(response, isFreeModelRequiringCostRemoval, capture); } console.error('[rewriteModelResponse] implementation error: unrecognized API kind %s', kind); From dfd96c5fa5ef28b650689216f1fb90c927e8b4f0 Mon Sep 17 00:00:00 2001 From: chrarnoldus <12196001+chrarnoldus@users.noreply.github.com> Date: Fri, 24 Jul 2026 15:19:00 +0000 Subject: [PATCH 2/5] style: format rewriteModelResponse.test.ts Co-authored-by: kiloconnect[bot] <240665456+kiloconnect[bot]@users.noreply.github.com> --- apps/web/src/lib/rewriteModelResponse.test.ts | 31 ++++++++++++------- 1 file changed, 19 insertions(+), 12 deletions(-) diff --git a/apps/web/src/lib/rewriteModelResponse.test.ts b/apps/web/src/lib/rewriteModelResponse.test.ts index 64809da6e9..9ddf308660 100644 --- a/apps/web/src/lib/rewriteModelResponse.test.ts +++ b/apps/web/src/lib/rewriteModelResponse.test.ts @@ -581,19 +581,22 @@ describe('request log capture', () => { expect(capture.setReadError).not.toHaveBeenCalled(); }); - test.each(rewriters)('%s: captures an empty body when upstream has no body', async (_name, rewrite) => { - const capture = makeCapture(); + test.each(rewriters)( + '%s: captures an empty body when upstream has no body', + async (_name, rewrite) => { + const capture = makeCapture(); - const result = await rewrite( - new Response(null, { headers: { 'content-type': 'text/event-stream' } }), - true, - capture - ); - await readOutputStream(result); + const result = await rewrite( + new Response(null, { headers: { 'content-type': 'text/event-stream' } }), + true, + capture + ); + await readOutputStream(result); - expect(capture.setBody).toHaveBeenCalledWith(''); - expect(capture.setReadError).not.toHaveBeenCalled(); - }); + expect(capture.setBody).toHaveBeenCalledWith(''); + expect(capture.setReadError).not.toHaveBeenCalled(); + } + ); test.each(rewriters)('%s: records a read error when the stream fails', async (_name, rewrite) => { const capture = makeCapture(); @@ -618,7 +621,11 @@ describe('request log capture', () => { async (_name, rewrite) => { const capture = makeCapture(); - const result = await rewrite(failingResponse('application/json', 'TimeoutError'), true, capture); + const result = await rewrite( + failingResponse('application/json', 'TimeoutError'), + true, + capture + ); expect(result.status).toBe(503); expect(capture.setReadError).toHaveBeenCalledTimes(1); From fd68f42c6f9f46b139ee155b0effe20cbd057f53 Mon Sep 17 00:00:00 2001 From: chrarnoldus <12196001+chrarnoldus@users.noreply.github.com> Date: Mon, 27 Jul 2026 09:00:40 +0000 Subject: [PATCH 3/5] refactor(ai-gateway): move request logging into rewriteModelResponse handleRequestLogging is folded into rewriteModelResponse: the wrapper creates the request log capture itself from the logging params, so the route makes a single call. Paths where the upstream response is replaced by a readable error use logUnrewrittenResponse. The KILO_ORGANIZATION_ID check in the skip condition is removed since the Kilo organization always has request logging enabled and therefore always gets a capture. Co-authored-by: kiloconnect[bot] <240665456+kiloconnect[bot]@users.noreply.github.com> --- .../api/openrouter/[...path]/route.test.ts | 20 +- .../src/app/api/openrouter/[...path]/route.ts | 57 ++---- .../lib/ai-gateway/handleRequestLogging.ts | 114 ----------- apps/web/src/lib/rewriteModelResponse.test.ts | 45 ++++- apps/web/src/lib/rewriteModelResponse.ts | 180 ++++++++++++++++-- 5 files changed, 238 insertions(+), 178 deletions(-) delete mode 100644 apps/web/src/lib/ai-gateway/handleRequestLogging.ts diff --git a/apps/web/src/app/api/openrouter/[...path]/route.test.ts b/apps/web/src/app/api/openrouter/[...path]/route.test.ts index e633a13cb3..1a6b962849 100644 --- a/apps/web/src/app/api/openrouter/[...path]/route.test.ts +++ b/apps/web/src/app/api/openrouter/[...path]/route.test.ts @@ -14,7 +14,7 @@ import { fetchEfficientAutoDecision } from '@/lib/ai-gateway/auto-routing-decisi import { logMicrodollarUsage } from '@/lib/ai-gateway/processUsage'; import { applyResolvedAutoModel } from '@/lib/ai-gateway/auto-model/resolution'; import { getDirectByokModel } from '@/lib/ai-gateway/providers/direct-byok'; -import { handleRequestLogging } from '@/lib/ai-gateway/handleRequestLogging'; +import { rewriteModelResponse } from '@/lib/rewriteModelResponse'; jest.mock('next/server', () => { return { @@ -55,9 +55,13 @@ jest.mock('@/lib/ai-gateway/o11y/api-metrics.server', () => ({ getToolsAvailable: jest.fn(() => false), getToolsUsed: jest.fn(() => false), })); -jest.mock('@/lib/ai-gateway/handleRequestLogging', () => ({ - handleRequestLogging: jest.fn(), -})); +jest.mock('@/lib/rewriteModelResponse', () => { + const actual = jest.requireActual('@/lib/rewriteModelResponse'); + return { + ...actual, + rewriteModelResponse: jest.fn(async () => null), + }; +}); jest.mock('@/lib/ai-gateway/llm-proxy-helpers', () => { const actual = jest.requireActual('@/lib/ai-gateway/llm-proxy-helpers'); return { @@ -96,7 +100,7 @@ const mockedFetchEfficientAutoDecision = jest.mocked(fetchEfficientAutoDecision) const mockedLogMicrodollarUsage = jest.mocked(logMicrodollarUsage); const mockedApplyResolvedAutoModel = jest.mocked(applyResolvedAutoModel); const mockedGetDirectByokModel = jest.mocked(getDirectByokModel); -const mockedHandleRequestLogging = jest.mocked(handleRequestLogging); +const mockedRewriteModelResponse = jest.mocked(rewriteModelResponse); const provider = { id: 'openrouter', @@ -255,7 +259,11 @@ describe('POST /api/openrouter/v1/chat/completions rules-engine actions', () => ); expect(response.status).toBe(200); - expect(mockedHandleRequestLogging).toHaveBeenCalledWith( + expect(mockedRewriteModelResponse).toHaveBeenCalledWith( + expect.anything(), + expect.anything(), + expect.anything(), + expect.anything(), expect.objectContaining({ vercel_request_id: 'iad1::iad1::request-id' }) ); }); diff --git a/apps/web/src/app/api/openrouter/[...path]/route.ts b/apps/web/src/app/api/openrouter/[...path]/route.ts index 049c5262b4..467e5a67c5 100644 --- a/apps/web/src/app/api/openrouter/[...path]/route.ts +++ b/apps/web/src/app/api/openrouter/[...path]/route.ts @@ -56,7 +56,7 @@ import { import { ProxyErrorType } from '@/lib/proxy-error-types'; import { getBalanceAndOrgSettings } from '@/lib/organizations/organization-usage'; import { isDataCollectionExplicitlyDisallowed } from '@/lib/ai-gateway/providers/openrouter/types'; -import { rewriteModelResponse } from '@/lib/rewriteModelResponse'; +import { rewriteModelResponse, logUnrewrittenResponse } from '@/lib/rewriteModelResponse'; import { createAnonymousContext, isAnonymousContext, @@ -69,7 +69,6 @@ import { checkPromotionLimit, } from '@/lib/free-model-rate-limiter'; import { PROMOTION_MAX_REQUESTS, PROMOTION_WINDOW_HOURS } from '@/lib/constants'; -import { handleRequestLogging } from '@/lib/ai-gateway/handleRequestLogging'; import { classifyAbuse, awaitClassifyAbuse, @@ -957,19 +956,13 @@ export async function POST(request: NextRequest): Promise { - if (user?.google_user_email.endsWith('@kilo.ai')) return true; - if (user?.google_user_email.endsWith('@kilocode.ai')) return true; - if (organizationId === KILO_ORGANIZATION_ID) return true; - return isDynamicallyOptedIntoRequestLogging({ - accountId: user?.id ?? null, - organizationId, - }); -} - -export async function handleRequestLogging(params: { - status: number; - user: User | null; - organization_id: string | null; - session_id: string | null; - vercel_request_id: string | null; - provider: string; - model: string; - request: GatewayRequest; -}): Promise { - const { status, user, organization_id, session_id, vercel_request_id, provider, model, request } = - params; - if (!(await isLoggingEnabledForUser(user, organization_id))) { - return null; - } - - let resolveCaptured: (result: CapturedResponseBody) => void = () => {}; - const captured = new Promise(resolve => { - resolveCaptured = resolve; - }); - let isSettled = false; - const settleOnce = (result: CapturedResponseBody) => { - if (!isSettled) { - isSettled = true; - resolveCaptured(result); - } - }; - - after(async () => { - // Wait until the response pipeline has processed the response body. This - // resolves when the response stream completes (or fails), which happens - // before after() callbacks are awaited. - const result = await captured; - const response = 'text' in result ? result.text : undefined; - const responseReadError = 'readError' in result ? result.readError : undefined; - if (responseReadError !== undefined) { - logExceptInTest( - `[handleRequestLogging] failed to read response body (user=${user?.id}, status=${status}, model=${model}): ${responseReadError}` - ); - } - try { - const error = - response !== undefined - ? detectToolCallArgumentErrors(response, request) - : { response_body_read_error: responseReadError }; - const apiRequestLogId = await db - .insert(api_request_log) - .values({ - kilo_user_id: user?.id, - organization_id: organization_id, - session_id, - vercel_request_id, - status_code: status, - model, - provider, - request: request.body, - response, - error, - }) - .returning({ id: api_request_log.id }); - logExceptInTest( - '[handleRequestLogging] Inserted into api_request_log', - apiRequestLogId[0].id - ); - } catch (e) { - const cause = e instanceof Error ? e.cause : undefined; - logExceptInTest( - `[handleRequestLogging] failed to insert api_request_log (user=${user?.id}, status=${status}, model=${model}) cause (truncated): ${String(cause).substring(0, 4000)} error (truncated): ${String(e).substring(0, 4000)}` - ); - } - }); - - return { - setBody: text => settleOnce({ text }), - setReadError: error => settleOnce({ readError: String(error).substring(0, 4000) }), - }; -} diff --git a/apps/web/src/lib/rewriteModelResponse.test.ts b/apps/web/src/lib/rewriteModelResponse.test.ts index 9ddf308660..6b1148e48d 100644 --- a/apps/web/src/lib/rewriteModelResponse.test.ts +++ b/apps/web/src/lib/rewriteModelResponse.test.ts @@ -1,12 +1,29 @@ -import { describe, test, expect, jest } from '@jest/globals'; +import { describe, test, expect, jest, beforeEach } from '@jest/globals'; import { rewriteModelResponse_ChatCompletions, rewriteModelResponse_Messages, rewriteModelResponse_Responses, rewriteModelResponse, + type RequestLogging, } from './rewriteModelResponse'; +import { isDynamicallyOptedIntoRequestLogging } from '@/lib/ai-gateway/request-logging-opt-ins'; import { KILO_ORGANIZATION_ID } from '@/lib/organizations/constants'; +jest.mock('next/server', () => ({ + ...(jest.requireActual('next/server') as Record), + after: jest.fn(), +})); + +jest.mock('@/lib/ai-gateway/request-logging-opt-ins', () => ({ + isDynamicallyOptedIntoRequestLogging: jest.fn(async () => false), +})); + +const mockedOptIn = jest.mocked(isDynamicallyOptedIntoRequestLogging); + +beforeEach(() => { + mockedOptIn.mockClear(); +}); + function jsonResponse(body: unknown, status = 200): Response { return new Response(JSON.stringify(body), { status, @@ -478,6 +495,17 @@ describe('rewriteModelResponse_Responses', () => { }); }); +function makeLogging(overrides?: Partial): RequestLogging { + return { + user: null, + organization_id: null, + session_id: null, + vercel_request_id: null, + request: { body: {} } as unknown as RequestLogging['request'], + ...overrides, + }; +} + describe('rewriteModelResponse', () => { test('rewrites paid-model Kilo organization traffic without stripping cost', async () => { const result = await rewriteModelResponse( @@ -492,7 +520,7 @@ describe('rewriteModelResponse', () => { 'openai/gpt-5', 'openrouter', 'chat_completions', - KILO_ORGANIZATION_ID + makeLogging({ organization_id: KILO_ORGANIZATION_ID }) ); expect(result).not.toBeNull(); @@ -511,7 +539,7 @@ describe('rewriteModelResponse', () => { 'openai/gpt-5', 'openrouter', 'chat_completions', - '00000000-0000-0000-0000-000000000000' + makeLogging({ organization_id: '00000000-0000-0000-0000-000000000000' }) ); expect(result).toBeNull(); @@ -525,7 +553,8 @@ describe('rewriteModelResponse', () => { }), 'google/gemma-4-26b-a4b-it:free', 'openrouter', - 'chat_completions' + 'chat_completions', + makeLogging() ); expect(result).not.toBeNull(); @@ -535,19 +564,17 @@ describe('rewriteModelResponse', () => { }); }); - test('processes responses it would normally skip when a capture is provided', async () => { - const capture = makeCapture(); + test('processes responses it would normally skip when request logging is enabled', async () => { + mockedOptIn.mockResolvedValueOnce(true); const result = await rewriteModelResponse( jsonResponse({ model: 'openai/gpt-5' }), 'openai/gpt-5', 'openrouter', 'chat_completions', - '00000000-0000-0000-0000-000000000000', - capture + makeLogging({ organization_id: '00000000-0000-0000-0000-000000000000' }) ); expect(result).not.toBeNull(); - expect(capture.setBody).toHaveBeenCalledWith(JSON.stringify({ model: 'openai/gpt-5' })); }); }); diff --git a/apps/web/src/lib/rewriteModelResponse.ts b/apps/web/src/lib/rewriteModelResponse.ts index 88c468bd41..9f577a17f6 100644 --- a/apps/web/src/lib/rewriteModelResponse.ts +++ b/apps/web/src/lib/rewriteModelResponse.ts @@ -1,17 +1,153 @@ +import { api_request_log, type User } from '@kilocode/db/schema'; import { isKiloExclusiveFreeModel } from '@/lib/ai-gateway/models'; -import type { RequestLogCapture } from '@/lib/ai-gateway/handleRequestLogging'; +import { detectToolCallArgumentErrors } from '@/lib/ai-gateway/api-request-log-errors'; import type { GatewayRequest } from '@/lib/ai-gateway/providers/openrouter/types'; import type { ProviderId } from '@/lib/ai-gateway/providers/types'; import { getOutputHeaders } from '@/lib/ai-gateway/llm-proxy-helpers'; import type { ChatCompletionChunk, OpenRouterUsage } from '@/lib/ai-gateway/processUsage.types'; +import { isDynamicallyOptedIntoRequestLogging } from '@/lib/ai-gateway/request-logging-opt-ins'; +import { db } from '@/lib/drizzle'; import { KILO_ORGANIZATION_ID } from '@/lib/organizations/constants'; import { errorExceptInTest, logExceptInTest } from '@/lib/utils.server'; import type { EventSourceMessage } from 'eventsource-parser'; import { createParser } from 'eventsource-parser'; -import { NextResponse } from 'next/server'; +import { after, NextResponse } from 'next/server'; import type OpenAI from 'openai'; import type Anthropic from '@anthropic-ai/sdk'; +/** + * Handle passed to the response pipeline so the upstream response body can be + * captured for request logging while the response is being processed anyway. + * This way the event stream is only processed once, instead of once for + * logging and once for rewriting. + */ +export type RequestLogCapture = { + /** Record the full upstream response body. Called at most once. */ + setBody(text: string): void; + /** Record that the upstream response body could not be read. Called at most once. */ + setReadError(error: unknown): void; +}; + +/** Parameters needed to write the request to api_request_log. */ +export type RequestLogging = { + user: User | null; + organization_id: string | null; + session_id: string | null; + vercel_request_id: string | null; + request: GatewayRequest; +}; + +type CapturedResponseBody = { text: string } | { readError: string }; + +async function isLoggingEnabledForUser( + user: User | null, + organizationId: string | null +): Promise { + if (user?.google_user_email.endsWith('@kilo.ai')) return true; + if (user?.google_user_email.endsWith('@kilocode.ai')) return true; + if (organizationId === KILO_ORGANIZATION_ID) return true; + return isDynamicallyOptedIntoRequestLogging({ + accountId: user?.id ?? null, + organizationId, + }); +} + +async function createRequestLogCapture( + response: Response, + model: string, + provider: string, + logging: RequestLogging +): Promise { + const { user, organization_id, session_id, vercel_request_id, request } = logging; + if (!(await isLoggingEnabledForUser(user, organization_id))) { + return null; + } + const status = response.status; + + let resolveCaptured: (result: CapturedResponseBody) => void = () => {}; + const captured = new Promise(resolve => { + resolveCaptured = resolve; + }); + let isSettled = false; + const settleOnce = (result: CapturedResponseBody) => { + if (!isSettled) { + isSettled = true; + resolveCaptured(result); + } + }; + + after(async () => { + // Wait until the response pipeline has processed the response body. This + // resolves when the response stream completes (or fails), which happens + // before after() callbacks are awaited. + const result = await captured; + const responseText = 'text' in result ? result.text : undefined; + const responseReadError = 'readError' in result ? result.readError : undefined; + if (responseReadError !== undefined) { + logExceptInTest( + `[rewriteModelResponse] failed to read response body (user=${user?.id}, status=${status}, model=${model}): ${responseReadError}` + ); + } + try { + const error = + responseText !== undefined + ? detectToolCallArgumentErrors(responseText, request) + : { response_body_read_error: responseReadError }; + const apiRequestLogId = await db + .insert(api_request_log) + .values({ + kilo_user_id: user?.id, + organization_id, + session_id, + vercel_request_id, + status_code: status, + model, + provider, + request: request.body, + response: responseText, + error, + }) + .returning({ id: api_request_log.id }); + logExceptInTest( + '[rewriteModelResponse] Inserted into api_request_log', + apiRequestLogId[0].id + ); + } catch (e) { + const cause = e instanceof Error ? e.cause : undefined; + logExceptInTest( + `[rewriteModelResponse] failed to insert api_request_log (user=${user?.id}, status=${status}, model=${model}) cause (truncated): ${String(cause).substring(0, 4000)} error (truncated): ${String(e).substring(0, 4000)}` + ); + } + }); + + return { + setBody: text => settleOnce({ text }), + setReadError: error => settleOnce({ readError: String(error).substring(0, 4000) }), + }; +} + +/** + * Logs the request and response body for paths where the upstream response is + * not passed through rewriteModelResponse (e.g. when an upstream error is + * replaced by a more readable one). Reads the response body once. + */ +export async function logUnrewrittenResponse( + response: Response, + model: string, + providerId: ProviderId, + logging: RequestLogging +): Promise { + const capture = await createRequestLogCapture(response, model, providerId, logging); + if (!capture) { + return; + } + try { + capture.setBody(await response.text()); + } catch (error) { + capture.setReadError(error); + } +} + type ResponseReadError = { errorType: 'timeout' | 'upstream_disconnect'; message: string; @@ -555,31 +691,49 @@ export async function rewriteModelResponse( model: string, providerId: ProviderId, kind: GatewayRequest['kind'], - organizationId?: string, - capture?: RequestLogCapture | null + logging: RequestLogging ): Promise { + const capture = await createRequestLogCapture(response, model, providerId, logging); const isFreeModelRequiringCostRemoval = (providerId === 'openrouter' || providerId === 'vercel') && isKiloExclusiveFreeModel(model); // When request logging is enabled the response has to be processed anyway // so the body can be captured for the request log in a single pass, so the // rewrite is not skipped in that case. - if (!isFreeModelRequiringCostRemoval && organizationId !== KILO_ORGANIZATION_ID && !capture) { + if (!isFreeModelRequiringCostRemoval && !capture) { console.debug('[rewriteModelResponse] skipping rewrite for %s', model); return null; } console.debug('[rewriteModelResponse] rewriting response for %s', model); - if (kind === 'chat_completions') { - return rewriteModelResponse_ChatCompletions(response, isFreeModelRequiringCostRemoval, capture); - } - if (kind === 'responses') { - return rewriteModelResponse_Responses(response, isFreeModelRequiringCostRemoval, capture); - } - if (kind === 'messages') { - return rewriteModelResponse_Messages(response, isFreeModelRequiringCostRemoval, capture); + try { + if (kind === 'chat_completions') { + return await rewriteModelResponse_ChatCompletions( + response, + isFreeModelRequiringCostRemoval, + capture + ); + } + if (kind === 'responses') { + return await rewriteModelResponse_Responses( + response, + isFreeModelRequiringCostRemoval, + capture + ); + } + if (kind === 'messages') { + return await rewriteModelResponse_Messages( + response, + isFreeModelRequiringCostRemoval, + capture + ); + } + } catch (error) { + capture?.setReadError(error); + throw error; } console.error('[rewriteModelResponse] implementation error: unrecognized API kind %s', kind); + capture?.setReadError(new Error('response was not processed')); return null; } From cf5cd24447795f93df654c745c0b7b29105d021d Mon Sep 17 00:00:00 2001 From: chrarnoldus <12196001+chrarnoldus@users.noreply.github.com> Date: Mon, 27 Jul 2026 09:23:08 +0000 Subject: [PATCH 4/5] refactor(ai-gateway): rename RequestLogging, drop obvious comments and wrapper try/catch - Rename RequestLogging to RequestLoggingParams. - Remove comments that just restate the code. - Remove the try/catch around the rewrite dispatch: it only guarded against unexpected throws from the JSON body read, which cannot happen in practice. - Fix jest.mock not being hoisted in rewriteModelResponse.test.ts: use the global jest instead of the @jest/globals import, matching the other test files. Co-authored-by: kiloconnect[bot] <240665456+kiloconnect[bot]@users.noreply.github.com> --- .../src/app/api/openrouter/[...path]/route.ts | 5 -- apps/web/src/lib/rewriteModelResponse.test.ts | 8 +-- apps/web/src/lib/rewriteModelResponse.ts | 54 +++++-------------- 3 files changed, 17 insertions(+), 50 deletions(-) diff --git a/apps/web/src/app/api/openrouter/[...path]/route.ts b/apps/web/src/app/api/openrouter/[...path]/route.ts index 467e5a67c5..ce1391116b 100644 --- a/apps/web/src/app/api/openrouter/[...path]/route.ts +++ b/apps/web/src/app/api/openrouter/[...path]/route.ts @@ -973,8 +973,6 @@ export async function POST(request: NextRequest): Promise { }); }); -function makeLogging(overrides?: Partial): RequestLogging { +function makeLogging(overrides?: Partial): RequestLoggingParams { return { user: null, organization_id: null, session_id: null, vercel_request_id: null, - request: { body: {} } as unknown as RequestLogging['request'], + request: { body: {} } as unknown as RequestLoggingParams['request'], ...overrides, }; } diff --git a/apps/web/src/lib/rewriteModelResponse.ts b/apps/web/src/lib/rewriteModelResponse.ts index 9f577a17f6..0923a05bc0 100644 --- a/apps/web/src/lib/rewriteModelResponse.ts +++ b/apps/web/src/lib/rewriteModelResponse.ts @@ -22,14 +22,11 @@ import type Anthropic from '@anthropic-ai/sdk'; * logging and once for rewriting. */ export type RequestLogCapture = { - /** Record the full upstream response body. Called at most once. */ setBody(text: string): void; - /** Record that the upstream response body could not be read. Called at most once. */ setReadError(error: unknown): void; }; -/** Parameters needed to write the request to api_request_log. */ -export type RequestLogging = { +export type RequestLoggingParams = { user: User | null; organization_id: string | null; session_id: string | null; @@ -56,7 +53,7 @@ async function createRequestLogCapture( response: Response, model: string, provider: string, - logging: RequestLogging + logging: RequestLoggingParams ): Promise { const { user, organization_id, session_id, vercel_request_id, request } = logging; if (!(await isLoggingEnabledForUser(user, organization_id))) { @@ -126,16 +123,12 @@ async function createRequestLogCapture( }; } -/** - * Logs the request and response body for paths where the upstream response is - * not passed through rewriteModelResponse (e.g. when an upstream error is - * replaced by a more readable one). Reads the response body once. - */ +/** For paths where the upstream response is not passed through rewriteModelResponse. */ export async function logUnrewrittenResponse( response: Response, model: string, providerId: ProviderId, - logging: RequestLogging + logging: RequestLoggingParams ): Promise { const capture = await createRequestLogCapture(response, model, providerId, logging); if (!capture) { @@ -399,7 +392,6 @@ export async function rewriteModelResponse_ChatCompletions( ); }, cancel() { - // The client disconnected before the stream completed. capture?.setReadError(new Error('response stream was cancelled')); }, }); @@ -550,7 +542,6 @@ export async function rewriteModelResponse_Messages( ); }, cancel() { - // The client disconnected before the stream completed. capture?.setReadError(new Error('response stream was cancelled')); }, }); @@ -674,7 +665,6 @@ export async function rewriteModelResponse_Responses( ); }, cancel() { - // The client disconnected before the stream completed. capture?.setReadError(new Error('response stream was cancelled')); }, }); @@ -691,7 +681,7 @@ export async function rewriteModelResponse( model: string, providerId: ProviderId, kind: GatewayRequest['kind'], - logging: RequestLogging + logging: RequestLoggingParams ): Promise { const capture = await createRequestLogCapture(response, model, providerId, logging); const isFreeModelRequiringCostRemoval = @@ -706,34 +696,16 @@ export async function rewriteModelResponse( } console.debug('[rewriteModelResponse] rewriting response for %s', model); - try { - if (kind === 'chat_completions') { - return await rewriteModelResponse_ChatCompletions( - response, - isFreeModelRequiringCostRemoval, - capture - ); - } - if (kind === 'responses') { - return await rewriteModelResponse_Responses( - response, - isFreeModelRequiringCostRemoval, - capture - ); - } - if (kind === 'messages') { - return await rewriteModelResponse_Messages( - response, - isFreeModelRequiringCostRemoval, - capture - ); - } - } catch (error) { - capture?.setReadError(error); - throw error; + if (kind === 'chat_completions') { + return rewriteModelResponse_ChatCompletions(response, isFreeModelRequiringCostRemoval, capture); + } + if (kind === 'responses') { + return rewriteModelResponse_Responses(response, isFreeModelRequiringCostRemoval, capture); + } + if (kind === 'messages') { + return rewriteModelResponse_Messages(response, isFreeModelRequiringCostRemoval, capture); } console.error('[rewriteModelResponse] implementation error: unrecognized API kind %s', kind); - capture?.setReadError(new Error('response was not processed')); return null; } From a75758583e4ed09c3dc0582aa18299616c8eda74 Mon Sep 17 00:00:00 2001 From: chrarnoldus <12196001+chrarnoldus@users.noreply.github.com> Date: Mon, 27 Jul 2026 09:36:03 +0000 Subject: [PATCH 5/5] fix(ai-gateway): settle request log capture when the body read throws readResponseText rethrows errors that are not recognized abort/timeout errors. Settle the capture before rethrowing so the after() callback awaiting it does not hang and the request is still logged without a response body. Addresses Kilobot review. Co-authored-by: kiloconnect[bot] <240665456+kiloconnect[bot]@users.noreply.github.com> --- apps/web/src/lib/rewriteModelResponse.ts | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/apps/web/src/lib/rewriteModelResponse.ts b/apps/web/src/lib/rewriteModelResponse.ts index 0923a05bc0..f9c2a40ecc 100644 --- a/apps/web/src/lib/rewriteModelResponse.ts +++ b/apps/web/src/lib/rewriteModelResponse.ts @@ -190,13 +190,17 @@ function getResponseReadError(error: unknown): ResponseReadError | null { async function readResponseText( response: Response, - headers: Headers + headers: Headers, + capture?: RequestLogCapture | null ): Promise<{ text: string } | { error: unknown; errorResponse: NextResponse }> { try { return { text: await response.text() }; } catch (error) { const responseReadError = getResponseReadError(error); if (!responseReadError) { + // Settle the capture so the after() callback awaiting it does not hang + // and the request is still logged (without a response body). + capture?.setReadError(error); throw error; } @@ -288,7 +292,7 @@ export async function rewriteModelResponse_ChatCompletions( if (headers.get('content-type')?.includes('application/json')) { // Read the body text once to avoid "Response body object should not be // disturbed or locked" errors that occur when `.clone().json()` fails. - const textResult = await readResponseText(response, headers); + const textResult = await readResponseText(response, headers, capture); if ('errorResponse' in textResult) { capture?.setReadError(textResult.error); return textResult.errorResponse; @@ -436,7 +440,7 @@ export async function rewriteModelResponse_Messages( const headers = getOutputHeaders(response); if (headers.get('content-type')?.includes('application/json')) { - const textResult = await readResponseText(response, headers); + const textResult = await readResponseText(response, headers, capture); if ('errorResponse' in textResult) { capture?.setReadError(textResult.error); return textResult.errorResponse; @@ -567,7 +571,7 @@ export async function rewriteModelResponse_Responses( const headers = getOutputHeaders(response); if (headers.get('content-type')?.includes('application/json')) { - const textResult = await readResponseText(response, headers); + const textResult = await readResponseText(response, headers, capture); if ('errorResponse' in textResult) { capture?.setReadError(textResult.error); return textResult.errorResponse; @@ -707,5 +711,6 @@ export async function rewriteModelResponse( } console.error('[rewriteModelResponse] implementation error: unrecognized API kind %s', kind); + capture?.setReadError(new Error('response was not processed')); return null; }