diff --git a/packages/core/src/core/loggingContentGenerator/loggingContentGenerator.test.ts b/packages/core/src/core/loggingContentGenerator/loggingContentGenerator.test.ts index c2965a52361..17b1d9286d2 100644 --- a/packages/core/src/core/loggingContentGenerator/loggingContentGenerator.test.ts +++ b/packages/core/src/core/loggingContentGenerator/loggingContentGenerator.test.ts @@ -15,6 +15,7 @@ import type { ContentGenerator } from '../contentGenerator.js'; import { AuthType } from '../contentGenerator.js'; import { LoggingContentGenerator } from './index.js'; import { OpenAIContentConverter } from '../openaiContentGenerator/converter.js'; +import { openaiRequestCaptureContext } from '../openaiContentGenerator/requestCaptureContext.js'; import { logApiRequest, logApiResponse, @@ -547,6 +548,290 @@ describe('LoggingContentGenerator', () => { ]); }); + it('logs the captured wire request including provider-injected fields (generateContent)', async () => { + const wireRequest: OpenAI.Chat.ChatCompletionCreateParams = { + model: 'deepseek-v4-pro', + messages: [{ role: 'user', content: 'hi' }], + temperature: 0.5, + max_tokens: 1024, + // Provider-injected fields the synthetic reconstruction would drop: + reasoning_effort: 'max', + extra_body: { thinking: { type: 'enabled' }, enable_thinking: true }, + metadata: { dashscope_user_id: 'abc' }, + } as unknown as OpenAI.Chat.ChatCompletionCreateParams; + + const wrapped = createWrappedGenerator( + vi.fn().mockImplementation(async () => { + openaiRequestCaptureContext.getStore()?.(wireRequest); + return createResponse('resp-cap', 'deepseek-v4-pro', [{ text: 'ok' }]); + }), + vi.fn(), + ); + + const generator = new LoggingContentGenerator(wrapped, createConfig(), { + model: 'deepseek-v4-pro', + authType: AuthType.USE_OPENAI, + enableOpenAILogging: true, + openAILoggingDir: 'logs', + }); + + const request = { + model: 'deepseek-v4-pro', + contents: [{ role: 'user', parts: [{ text: 'hi' }] }], + } as unknown as GenerateContentParameters; + + await generator.generateContent(request, 'prompt-cap'); + + const openaiLoggerInstance = vi.mocked(OpenAILogger).mock.results[0] + ?.value as { logInteraction: ReturnType }; + expect(openaiLoggerInstance.logInteraction).toHaveBeenCalledTimes(1); + const [loggedRequest] = openaiLoggerInstance.logInteraction.mock + .calls[0] as [OpenAI.Chat.ChatCompletionCreateParams]; + // The logger must observe the actual wire request, not a stripped reconstruction. + expect(loggedRequest).toBe(wireRequest); + expect(loggedRequest).toMatchObject({ + reasoning_effort: 'max', + extra_body: { thinking: { type: 'enabled' }, enable_thinking: true }, + metadata: { dashscope_user_id: 'abc' }, + }); + }); + + it('logs the captured wire request for streaming requests (generateContentStream)', async () => { + const wireRequest: OpenAI.Chat.ChatCompletionCreateParams = { + model: 'glm-5.1', + messages: [{ role: 'user', content: 'hi' }], + stream: true, + stream_options: { include_usage: true }, + extra_body: { thinking: { type: 'enabled' } }, + } as unknown as OpenAI.Chat.ChatCompletionCreateParams; + + const chunk = createResponse('resp-stream-cap', 'glm-5.1', [ + { text: 'ok' }, + ]); + + const wrapped = createWrappedGenerator( + vi.fn(), + vi.fn().mockImplementation(async () => { + openaiRequestCaptureContext.getStore()?.(wireRequest); + return (async function* () { + yield chunk; + })(); + }), + ); + + const generator = new LoggingContentGenerator(wrapped, createConfig(), { + model: 'glm-5.1', + authType: AuthType.USE_OPENAI, + enableOpenAILogging: true, + openAILoggingDir: 'logs', + }); + + const request = { + model: 'glm-5.1', + contents: [{ role: 'user', parts: [{ text: 'hi' }] }], + } as unknown as GenerateContentParameters; + + const stream = await generator.generateContentStream( + request, + 'prompt-stream-cap', + ); + for await (const _ of stream) { + // drain + } + + const openaiLoggerInstance = vi.mocked(OpenAILogger).mock.results[0] + ?.value as { logInteraction: ReturnType }; + expect(openaiLoggerInstance.logInteraction).toHaveBeenCalledTimes(1); + const [loggedRequest] = openaiLoggerInstance.logInteraction.mock + .calls[0] as [OpenAI.Chat.ChatCompletionCreateParams]; + expect(loggedRequest).toBe(wireRequest); + expect(loggedRequest).toMatchObject({ + stream: true, + stream_options: { include_usage: true }, + extra_body: { thinking: { type: 'enabled' } }, + }); + }); + + it('falls back to synthetic request when the wrapped generator does not capture', async () => { + const wrapped = createWrappedGenerator( + vi + .fn() + .mockResolvedValue( + createResponse('resp-fallback', 'test-model', [{ text: 'ok' }]), + ), + vi.fn(), + ); + + const generator = new LoggingContentGenerator(wrapped, createConfig(), { + model: 'test-model', + authType: AuthType.USE_OPENAI, + enableOpenAILogging: true, + openAILoggingDir: 'logs', + }); + + const request = { + model: 'test-model', + contents: [{ role: 'user', parts: [{ text: 'hi' }] }], + config: { temperature: 0.4 }, + } as unknown as GenerateContentParameters; + + await generator.generateContent(request, 'prompt-fallback'); + + const openaiLoggerInstance = vi.mocked(OpenAILogger).mock.results[0] + ?.value as { logInteraction: ReturnType }; + const [loggedRequest] = openaiLoggerInstance.logInteraction.mock + .calls[0] as [OpenAI.Chat.ChatCompletionCreateParams]; + expect(loggedRequest).toEqual( + expect.objectContaining({ + model: 'test-model', + temperature: 0.4, + }), + ); + }); + + it('does not propagate logging-side throws (success and error paths)', async () => { + const successResponse = createResponse('resp-safe', 'test-model', [ + { text: 'ok' }, + ]); + const successWrapped = createWrappedGenerator( + vi.fn().mockResolvedValue(successResponse), + vi.fn(), + ); + const successGen = new LoggingContentGenerator( + successWrapped, + createConfig(), + { + model: 'test-model', + authType: AuthType.USE_OPENAI, + enableOpenAILogging: true, + openAILoggingDir: 'logs', + }, + ); + + // No capture fires, so resolve() falls through to the synthetic builder. + // Force the synthetic build to throw, then verify the API result still surfaces. + convertGeminiRequestToOpenAISpy.mockImplementationOnce(() => { + throw new Error('synth-fail-success'); + }); + + const request = { + model: 'test-model', + contents: [{ role: 'user', parts: [{ text: 'hi' }] }], + } as unknown as GenerateContentParameters; + + await expect( + successGen.generateContent(request, 'prompt-safe-success'), + ).resolves.toBe(successResponse); + + const apiError = new Error('api-boom'); + const errorWrapped = createWrappedGenerator( + vi.fn().mockRejectedValue(apiError), + vi.fn(), + ); + const errorGen = new LoggingContentGenerator(errorWrapped, createConfig(), { + model: 'test-model', + authType: AuthType.USE_OPENAI, + enableOpenAILogging: true, + openAILoggingDir: 'logs', + }); + convertGeminiRequestToOpenAISpy.mockImplementationOnce(() => { + throw new Error('synth-fail-error'); + }); + + await expect( + errorGen.generateContent(request, 'prompt-safe-error'), + ).rejects.toThrow('api-boom'); + }); + + it('does not propagate logging-side throws on a successful stream', async () => { + const chunk1 = createResponse('resp-stream-safe-1', 'test-model', [ + { text: 'hello' }, + ]); + const chunk2 = createResponse('resp-stream-safe-2', 'test-model', [ + { text: ' world' }, + ]); + const wrapped = createWrappedGenerator( + vi.fn(), + vi.fn().mockResolvedValue( + (async function* () { + yield chunk1; + yield chunk2; + })(), + ), + ); + const generator = new LoggingContentGenerator(wrapped, createConfig(), { + model: 'test-model', + authType: AuthType.USE_OPENAI, + enableOpenAILogging: true, + openAILoggingDir: 'logs', + }); + const openaiLoggerInstance = vi.mocked(OpenAILogger).mock.results[0] + ?.value as { logInteraction: ReturnType }; + openaiLoggerInstance.logInteraction.mockRejectedValueOnce( + new Error('log-fail-on-stream-success'), + ); + + const request = { + model: 'test-model', + contents: [{ role: 'user', parts: [{ text: 'hi' }] }], + } as unknown as GenerateContentParameters; + + const stream = await generator.generateContentStream( + request, + 'prompt-stream-safe-success', + ); + const seen: GenerateContentResponse[] = []; + for await (const item of stream) { + seen.push(item); + } + // All chunks must reach the consumer; the logger throw must not surface. + expect(seen).toHaveLength(2); + expect(openaiLoggerInstance.logInteraction).toHaveBeenCalledTimes(1); + }); + + it('does not let logging-side throws replace the original stream error', async () => { + const chunk = createResponse('resp-stream-err', 'test-model', [ + { text: 'partial' }, + ]); + const apiError = new Error('stream-api-fail'); + const wrapped = createWrappedGenerator( + vi.fn(), + vi.fn().mockResolvedValue( + (async function* () { + yield chunk; + throw apiError; + })(), + ), + ); + const generator = new LoggingContentGenerator(wrapped, createConfig(), { + model: 'test-model', + authType: AuthType.USE_OPENAI, + enableOpenAILogging: true, + openAILoggingDir: 'logs', + }); + const openaiLoggerInstance = vi.mocked(OpenAILogger).mock.results[0] + ?.value as { logInteraction: ReturnType }; + openaiLoggerInstance.logInteraction.mockRejectedValueOnce( + new Error('log-fail-on-stream-error'), + ); + + const request = { + model: 'test-model', + contents: [{ role: 'user', parts: [{ text: 'hi' }] }], + } as unknown as GenerateContentParameters; + + const stream = await generator.generateContentStream( + request, + 'prompt-stream-safe-error', + ); + await expect(async () => { + for await (const _item of stream) { + // drain + } + }).rejects.toThrow('stream-api-fail'); + expect(openaiLoggerInstance.logInteraction).toHaveBeenCalledTimes(1); + }); + it.each(['prompt_suggestion', 'forked_query', 'speculation'])( 'skips logApiRequest and OpenAI logging for internal promptId %s (generateContent)', async (promptId) => { diff --git a/packages/core/src/core/loggingContentGenerator/loggingContentGenerator.ts b/packages/core/src/core/loggingContentGenerator/loggingContentGenerator.ts index b6ad0229a97..51643d98aae 100644 --- a/packages/core/src/core/loggingContentGenerator/loggingContentGenerator.ts +++ b/packages/core/src/core/loggingContentGenerator/loggingContentGenerator.ts @@ -39,6 +39,7 @@ import type { InputModalities, } from '../contentGenerator.js'; import { OpenAIContentConverter } from '../openaiContentGenerator/converter.js'; +import { openaiRequestCaptureContext } from '../openaiContentGenerator/requestCaptureContext.js'; import type { RequestContext } from '../openaiContentGenerator/types.js'; import { OpenAILogger } from '../../utils/openaiLogger.js'; import { @@ -166,11 +167,13 @@ export class LoggingContentGenerator implements ContentGenerator { userPromptId, ); } - const openaiRequest = isInternal - ? undefined - : await this.buildOpenAIRequestForLogging(req); + + const session = this.startCaptureSession(isInternal); + try { - const response = await this.wrapped.generateContent(req, userPromptId); + const response = await session.wrap(() => + this.wrapped.generateContent(req, userPromptId), + ); const durationMs = Date.now() - startTime; const responseText = isInternal ? undefined @@ -184,14 +187,30 @@ export class LoggingContentGenerator implements ContentGenerator { responseText, ); if (!isInternal) { - await this.logOpenAIInteraction(openaiRequest, response); + // Logging must not discard the successful API response if + // resolve() or logOpenAIInteraction throws. + try { + await this.logOpenAIInteraction(await session.resolve(req), response); + } catch { + // swallow logging-side error + } } return response; } catch (error) { const durationMs = Date.now() - startTime; this._logApiError('', durationMs, error, req.model, userPromptId); if (!isInternal) { - await this.logOpenAIInteraction(openaiRequest, undefined, error); + // Logging must not replace the original API error if + // resolve() or logOpenAIInteraction throws. + try { + await this.logOpenAIInteraction( + await session.resolve(req), + undefined, + error, + ); + } catch { + // swallow logging-side error + } } throw error; } @@ -210,31 +229,70 @@ export class LoggingContentGenerator implements ContentGenerator { userPromptId, ); } - const openaiRequest = isInternal - ? undefined - : await this.buildOpenAIRequestForLogging(req); + + const session = this.startCaptureSession(isInternal); let stream: AsyncGenerator; try { - stream = await this.wrapped.generateContentStream(req, userPromptId); + stream = await session.wrap(() => + this.wrapped.generateContentStream(req, userPromptId), + ); } catch (error) { const durationMs = Date.now() - startTime; this._logApiError('', durationMs, error, req.model, userPromptId); if (!isInternal) { - await this.logOpenAIInteraction(openaiRequest, undefined, error); + // Logging must not replace the original API error if + // resolve() or logOpenAIInteraction throws. + try { + await this.logOpenAIInteraction( + await session.resolve(req), + undefined, + error, + ); + } catch { + // swallow logging-side error + } } throw error; } + let resolvedRequest: OpenAI.Chat.ChatCompletionCreateParams | undefined; + if (!isInternal) { + try { + resolvedRequest = await session.resolve(req); + } catch { + // Resolve must not abort the stream that the SDK already returned. + } + } return this.loggingStreamWrapper( stream, startTime, userPromptId, req.model, - openaiRequest, + resolvedRequest, ); } + private startCaptureSession(isInternal: boolean): { + wrap: (fn: () => Promise) => Promise; + resolve: ( + req: GenerateContentParameters, + ) => Promise; + } { + let captured: OpenAI.Chat.ChatCompletionCreateParams | undefined; + const skipCapture = isInternal || !this.openaiLogger; + return { + wrap: (fn: () => Promise): Promise => + skipCapture + ? fn() + : openaiRequestCaptureContext.run((built) => { + captured = built; + }, fn), + resolve: async (req) => + captured ?? (await this.buildOpenAIRequestForLogging(req)), + }; + } + private async *loggingStreamWrapper( stream: AsyncGenerator, startTime: number, @@ -282,7 +340,12 @@ export class LoggingContentGenerator implements ContentGenerator { this.extractResponseText(consolidatedResponse), ); if (!isInternal) { - await this.logOpenAIInteraction(openaiRequest, consolidatedResponse); + // Logging must not turn a fully-yielded stream into a thrown error. + try { + await this.logOpenAIInteraction(openaiRequest, consolidatedResponse); + } catch { + // swallow logging-side error + } } } catch (error) { const durationMs = Date.now() - startTime; @@ -294,7 +357,12 @@ export class LoggingContentGenerator implements ContentGenerator { userPromptId, ); if (!isInternal) { - await this.logOpenAIInteraction(openaiRequest, undefined, error); + // Logging must not replace the original stream/API error. + try { + await this.logOpenAIInteraction(openaiRequest, undefined, error); + } catch { + // swallow logging-side error + } } throw error; } diff --git a/packages/core/src/core/openaiContentGenerator/pipeline.test.ts b/packages/core/src/core/openaiContentGenerator/pipeline.test.ts index d19ff27d787..23defe9ca7f 100644 --- a/packages/core/src/core/openaiContentGenerator/pipeline.test.ts +++ b/packages/core/src/core/openaiContentGenerator/pipeline.test.ts @@ -12,6 +12,7 @@ import { GenerateContentResponse, Type, FinishReason } from '@google/genai'; import type { ErrorHandler, PipelineConfig } from './types.js'; import { ContentGenerationPipeline, StreamContentError } from './pipeline.js'; import { OpenAIContentConverter } from './converter.js'; +import { openaiRequestCaptureContext } from './requestCaptureContext.js'; import { StreamingToolCallParser } from './streamingToolCallParser.js'; import type { Config } from '../../config/config.js'; import type { ContentGeneratorConfig, AuthType } from '../contentGenerator.js'; @@ -1934,4 +1935,156 @@ describe('ContentGenerationPipeline', () => { expect(responses[0]).toBe(finalGeminiResponse); }); }); + + describe('openaiRequestCaptureContext integration', () => { + it('forwards the provider-enhanced request to the active capture', async () => { + const request: GenerateContentParameters = { + model: 'test-model', + contents: [{ parts: [{ text: 'Hello' }], role: 'user' }], + }; + + const mockMessages = [ + { role: 'user', content: 'Hello' }, + ] as OpenAI.Chat.ChatCompletionMessageParam[]; + (mockConverter.convertGeminiRequestToOpenAI as Mock).mockReturnValue( + mockMessages, + ); + (mockConverter.convertOpenAIResponseToGemini as Mock).mockReturnValue( + new GenerateContentResponse(), + ); + (mockClient.chat.completions.create as Mock).mockResolvedValue({ + id: 'r', + choices: [], + created: 0, + model: 'test-model', + } as unknown as OpenAI.Chat.ChatCompletion); + + // Provider injects extra_body and metadata, mimicking real DashScope behavior. + (mockProvider.buildRequest as Mock).mockImplementation((req) => ({ + ...req, + extra_body: { thinking: { type: 'enabled' } }, + metadata: { user_id: 'abc' }, + })); + + let captured: OpenAI.Chat.ChatCompletionCreateParams | undefined; + await openaiRequestCaptureContext.run( + (built) => { + captured = built; + }, + () => pipeline.execute(request, 'p'), + ); + + expect(captured).toBeDefined(); + // The captured request must be the same object passed to the SDK. + expect(mockClient.chat.completions.create).toHaveBeenCalledWith( + captured, + expect.anything(), + ); + expect(captured).toEqual( + expect.objectContaining({ + model: 'test-model', + messages: mockMessages, + extra_body: { thinking: { type: 'enabled' } }, + metadata: { user_id: 'abc' }, + }), + ); + }); + + it('captures the streaming request including stream/stream_options', async () => { + const request: GenerateContentParameters = { + model: 'test-model', + contents: [{ parts: [{ text: 'Hello' }], role: 'user' }], + }; + + const mockMessages = [ + { role: 'user', content: 'Hello' }, + ] as OpenAI.Chat.ChatCompletionMessageParam[]; + (mockConverter.convertGeminiRequestToOpenAI as Mock).mockReturnValue( + mockMessages, + ); + (mockConverter.convertOpenAIChunkToGemini as Mock).mockReturnValue( + new GenerateContentResponse(), + ); + + const fakeStream = (async function* () { + // empty stream + })(); + (mockClient.chat.completions.create as Mock).mockResolvedValue( + fakeStream, + ); + + (mockProvider.buildRequest as Mock).mockImplementation((req) => ({ + ...req, + extra_body: { enable_thinking: true }, + })); + + let captured: OpenAI.Chat.ChatCompletionCreateParams | undefined; + await openaiRequestCaptureContext.run( + (built) => { + captured = built; + }, + async () => { + const stream = await pipeline.executeStream(request, 'p'); + for await (const _ of stream) { + // drain + } + }, + ); + + expect(captured).toBeDefined(); + expect(captured).toEqual( + expect.objectContaining({ + stream: true, + stream_options: { include_usage: true }, + extra_body: { enable_thinking: true }, + }), + ); + }); + + it('isolates concurrent captures', async () => { + const request: GenerateContentParameters = { + model: 'test-model', + contents: [{ parts: [{ text: 'Hello' }], role: 'user' }], + }; + + (mockConverter.convertGeminiRequestToOpenAI as Mock).mockReturnValue([]); + (mockConverter.convertOpenAIResponseToGemini as Mock).mockReturnValue( + new GenerateContentResponse(), + ); + (mockClient.chat.completions.create as Mock).mockResolvedValue({ + id: 'r', + choices: [], + created: 0, + model: 'test-model', + } as unknown as OpenAI.Chat.ChatCompletion); + + let n = 0; + (mockProvider.buildRequest as Mock).mockImplementation((req) => ({ + ...req, + extra_body: { call_index: ++n }, + })); + + const runOne = async () => { + let captured: OpenAI.Chat.ChatCompletionCreateParams | undefined; + await openaiRequestCaptureContext.run( + (built) => { + captured = built; + }, + () => pipeline.execute(request, 'p'), + ); + return captured; + }; + + const [a, b] = await Promise.all([runOne(), runOne()]); + expect(a).toBeDefined(); + expect(b).toBeDefined(); + // Each call's capture must have received its own object — + // the outer AsyncLocalStorage stores must not bleed across awaits. + const aExtra = (a as unknown as { extra_body: { call_index: number } }) + .extra_body; + const bExtra = (b as unknown as { extra_body: { call_index: number } }) + .extra_body; + expect(aExtra).not.toEqual(bExtra); + }); + }); }); diff --git a/packages/core/src/core/openaiContentGenerator/pipeline.ts b/packages/core/src/core/openaiContentGenerator/pipeline.ts index 7fc8e0f92a9..21487e655f6 100644 --- a/packages/core/src/core/openaiContentGenerator/pipeline.ts +++ b/packages/core/src/core/openaiContentGenerator/pipeline.ts @@ -13,6 +13,7 @@ import { import type { ContentGeneratorConfig } from '../contentGenerator.js'; import { OpenAIContentConverter } from './converter.js'; import { isDeepSeekHostname } from './provider/deepseek.js'; +import { openaiRequestCaptureContext } from './requestCaptureContext.js'; import { StreamingToolCallParser } from './streamingToolCallParser.js'; import { TaggedThinkingParser } from './taggedThinkingParser.js'; import type { PipelineConfig, RequestContext } from './types.js'; @@ -509,6 +510,11 @@ export class ContentGenerationPipeline { isStreaming, ); + // Position is load-bearing: capture must run after buildRequest (post + // provider enhancement, post disable-reasoning) and before the SDK call + // so the logger sees the exact bytes sent on the wire. + openaiRequestCaptureContext.getStore()?.(openaiRequest); + const result = await executor(openaiRequest, context); return result; } catch (error) { diff --git a/packages/core/src/core/openaiContentGenerator/requestCaptureContext.ts b/packages/core/src/core/openaiContentGenerator/requestCaptureContext.ts new file mode 100644 index 00000000000..dfceefeecae --- /dev/null +++ b/packages/core/src/core/openaiContentGenerator/requestCaptureContext.ts @@ -0,0 +1,23 @@ +/** + * @license + * Copyright 2025 Qwen + * SPDX-License-Identifier: Apache-2.0 + */ + +import { AsyncLocalStorage } from 'node:async_hooks'; +import type OpenAI from 'openai'; + +/** + * Lets the `LoggingContentGenerator` decorator observe the exact OpenAI + * request that `ContentGenerationPipeline` built and handed to the SDK, + * including provider-injected fields (`extra_body`, `metadata`, + * `stream_options`, `samplingParams` pass-through keys, etc.) — without + * which the logger would have to reconstruct a parallel request and + * silently miss those fields. + */ +export type OpenAIRequestCapture = ( + request: OpenAI.Chat.ChatCompletionCreateParams, +) => void; + +export const openaiRequestCaptureContext = + new AsyncLocalStorage();