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 0b1b913c2b..ce1391116b 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,16 +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: { - clonedResponse: Response; - user: User | null; - organization_id: string | null; - session_id: string | null; - vercel_request_id: string | null; - provider: string; - model: string; - request: GatewayRequest; -}) { - const { - clonedResponse, - user, - organization_id, - session_id, - vercel_request_id, - provider, - model, - request, - } = params; - if (!(await isLoggingEnabledForUser(user, organization_id))) { - return; - } - 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); - logExceptInTest( - `[handleRequestLogging] failed to read response body (user=${user?.id}, status=${clonedResponse.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: clonedResponse.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=${clonedResponse.status}, model=${model}) cause (truncated): ${String(cause).substring(0, 4000)} error (truncated): ${String(e).substring(0, 4000)}` - ); - } - }); -} diff --git a/apps/web/src/lib/rewriteModelResponse.test.ts b/apps/web/src/lib/rewriteModelResponse.test.ts index dc204fd350..62f111c99e 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, beforeEach } from '@jest/globals'; import { rewriteModelResponse_ChatCompletions, rewriteModelResponse_Messages, rewriteModelResponse_Responses, rewriteModelResponse, + type RequestLoggingParams, } 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): RequestLoggingParams { + return { + user: null, + organization_id: null, + session_id: null, + vercel_request_id: null, + request: { body: {} } as unknown as RequestLoggingParams['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(); @@ -534,4 +563,114 @@ describe('rewriteModelResponse', () => { usage: {}, }); }); + + 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', + makeLogging({ organization_id: '00000000-0000-0000-0000-000000000000' }) + ); + + expect(result).not.toBeNull(); + }); +}); + +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..f9c2a40ecc 100644 --- a/apps/web/src/lib/rewriteModelResponse.ts +++ b/apps/web/src/lib/rewriteModelResponse.ts @@ -1,16 +1,146 @@ +import { api_request_log, type User } from '@kilocode/db/schema'; import { isKiloExclusiveFreeModel } from '@/lib/ai-gateway/models'; +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 = { + setBody(text: string): void; + setReadError(error: unknown): void; +}; + +export type RequestLoggingParams = { + 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: RequestLoggingParams +): 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) }), + }; +} + +/** For paths where the upstream response is not passed through rewriteModelResponse. */ +export async function logUnrewrittenResponse( + response: Response, + model: string, + providerId: ProviderId, + logging: RequestLoggingParams +): 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; @@ -60,17 +190,22 @@ function getResponseReadError(error: unknown): ResponseReadError | null { async function readResponseText( response: Response, - headers: Headers -): Promise<{ text: string } | { errorResponse: NextResponse }> { + 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; } return { + error, errorResponse: NextResponse.json( { error: responseReadError.message, @@ -89,9 +224,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 +242,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,16 +282,22 @@ 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')) { // 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; } + capture?.setBody(textResult.text); const { text } = textResult; let json: OpenAI.ChatCompletion; try { @@ -174,6 +327,7 @@ export async function rewriteModelResponse_ChatCompletions(response: Response, r const reader = response.body?.getReader(); if (!reader) { controller.close(); + capture?.setBody(''); return; } @@ -237,9 +391,13 @@ export async function rewriteModelResponse_ChatCompletions(response: Response, r }, }) + '\n\n', - progress.stop + progress.stop, + capture ); }, + cancel() { + capture?.setReadError(new Error('response stream was cancelled')); + }, }); return new NextResponse(stream, { @@ -274,14 +432,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); + const textResult = await readResponseText(response, headers, capture); 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 +475,7 @@ export async function rewriteModelResponse_Messages(response: Response, removeCo const reader = response.body?.getReader(); if (!reader) { controller.close(); + capture?.setBody(''); return; } @@ -376,9 +541,13 @@ export async function rewriteModelResponse_Messages(response: Response, removeCo }, }) + '\n\n', - progress.stop + progress.stop, + capture ); }, + cancel() { + capture?.setReadError(new Error('response stream was cancelled')); + }, }); return new NextResponse(stream, { @@ -394,14 +563,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); + const textResult = await readResponseText(response, headers, capture); 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 +606,7 @@ export async function rewriteModelResponse_Responses(response: Response, removeC const reader = response.body?.getReader(); if (!reader) { controller.close(); + capture?.setBody(''); return; } @@ -488,9 +664,13 @@ export async function rewriteModelResponse_Responses(response: Response, removeC }, }) + '\n\n', - progress.stop + progress.stop, + capture ); }, + cancel() { + capture?.setReadError(new Error('response stream was cancelled')); + }, }); return new NextResponse(stream, { @@ -505,27 +685,32 @@ export async function rewriteModelResponse( model: string, providerId: ProviderId, kind: GatewayRequest['kind'], - organizationId?: string + logging: RequestLoggingParams ): Promise { + const capture = await createRequestLogCapture(response, model, providerId, logging); 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 && !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); + capture?.setReadError(new Error('response was not processed')); return null; }