perf(server): decouple proxy traces from responses (#1014)

This commit is contained in:
程序员阿江(Relakkes)
2026-07-17 22:49:48 +08:00
parent 15e554a565
commit b6bbbd5ff8
2 changed files with 196 additions and 29 deletions
+163 -2
View File
@@ -9,7 +9,11 @@ import * as os from 'os'
import { ProviderService } from '../services/providerService.js'
import { handleProvidersApi } from '../api/providers.js'
import { handleProxyRequest } from '../proxy/handler.js'
import { clearTraceCaptureStateForTests, traceCaptureService } from '../services/traceCaptureService.js'
import {
clearTraceCaptureStateForTests,
setTraceAppendBeforeWriteHookForTests,
traceCaptureService,
} from '../services/traceCaptureService.js'
import type { CreateProviderInput } from '../types/provider.js'
// ─── Test helpers ─────────────────────────────────────────────────────────────
@@ -89,6 +93,46 @@ async function readProvidersConfig(): Promise<Record<string, unknown>> {
return JSON.parse(raw) as Record<string, unknown>
}
async function waitForCompletedProxyTrace(sessionId: string) {
for (let attempt = 0; attempt < 100; attempt += 1) {
const trace = await traceCaptureService.getSessionTrace(sessionId)
if (
trace.calls.some((call) => call.response) &&
trace.events.some((event) => event.phase === 'upstream_fetch_completed')
) {
return trace
}
await new Promise((resolve) => setTimeout(resolve, 5))
}
return traceCaptureService.getSessionTrace(sessionId)
}
function blockNextTraceAppend() {
let releaseWrite: () => void = () => {}
const blockedWrite = new Promise<void>((resolve) => {
releaseWrite = resolve
})
let signalBlocked: () => void = () => {}
const writeBlocked = new Promise<void>((resolve) => {
signalBlocked = resolve
})
setTraceAppendBeforeWriteHookForTests(async () => {
setTraceAppendBeforeWriteHookForTests(null)
signalBlocked()
await blockedWrite
})
return { releaseWrite, writeBlocked }
}
async function settlesBeforeBlockedTraceWrite<T>(promise: Promise<T>): Promise<T | null> {
return Promise.race([
promise,
new Promise<null>((resolve) => setTimeout(() => resolve(null), 50)),
])
}
// =============================================================================
// ProviderService
// =============================================================================
@@ -1352,7 +1396,7 @@ describe('ProviderService', () => {
})
const res = await handleProxyRequest(req, new URL(req.url))
const trace = await traceCaptureService.getSessionTrace('session-proxy-trace')
const trace = await waitForCompletedProxyTrace('session-proxy-trace')
expect(res.status).toBe(200)
expect(trace.summary.apiCalls).toBe(1)
@@ -1374,6 +1418,123 @@ describe('ProviderService', () => {
}
})
test('returns a non-streaming proxy response before trace persistence finishes', async () => {
const originalFetch = globalThis.fetch
globalThis.fetch = mock(async () => new Response(JSON.stringify({
id: 'chatcmpl-trace-background',
object: 'chat.completion',
created: 0,
model: 'gpt-4',
choices: [{ index: 0, message: { role: 'assistant', content: 'background trace ok' }, finish_reason: 'stop' }],
usage: { prompt_tokens: 5, completion_tokens: 3, total_tokens: 8 },
}), {
status: 200,
headers: { 'Content-Type': 'application/json' },
})) as typeof fetch
const svc = new ProviderService()
const provider = await svc.addProvider(sampleInput({ apiFormat: 'openai_chat' }))
await svc.activateProvider(provider.id)
const { releaseWrite, writeBlocked } = blockNextTraceAppend()
let released = false
let responsePromise: Promise<Response> | undefined
try {
const req = new Request('http://localhost:3456/proxy/v1/messages', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'X-Claude-Code-Session-Id': 'session-non-stream-background-trace',
},
body: JSON.stringify({
model: 'gpt-4',
max_tokens: 64,
messages: [{ role: 'user', content: 'return before trace persistence' }],
}),
})
responsePromise = handleProxyRequest(req, new URL(req.url))
await writeBlocked
const response = await settlesBeforeBlockedTraceWrite(responsePromise)
expect(response).not.toBeNull()
expect(response?.status).toBe(200)
await expect(response?.json()).resolves.toMatchObject({
content: [{ text: 'background trace ok' }],
})
releaseWrite()
released = true
const trace = await waitForCompletedProxyTrace('session-non-stream-background-trace')
expect(trace.calls[0]?.response?.body.preview).toContain('chatcmpl-trace-background')
expect(trace.events.at(-1)?.phase).toBe('upstream_fetch_completed')
} finally {
if (!released) releaseWrite()
await responsePromise?.catch(() => undefined)
globalThis.fetch = originalFetch
}
})
test('delivers streaming EOF before trace persistence finishes', async () => {
const originalFetch = globalThis.fetch
const encoder = new TextEncoder()
globalThis.fetch = mock(async () => new Response(new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(encoder.encode([
'data: {"id":"chatcmpl-stream-trace","object":"chat.completion.chunk","model":"gpt-4","choices":[{"index":0,"delta":{"role":"assistant","content":"streamed"},"finish_reason":null}]}',
'',
'data: {"id":"chatcmpl-stream-trace","object":"chat.completion.chunk","model":"gpt-4","choices":[{"index":0,"delta":{},"finish_reason":"stop"}],"usage":{"prompt_tokens":5,"completion_tokens":2,"total_tokens":7}}',
'',
'data: [DONE]',
'',
].join('\n')))
controller.close()
},
}), {
status: 200,
headers: { 'Content-Type': 'text/event-stream' },
})) as typeof fetch
const svc = new ProviderService()
const provider = await svc.addProvider(sampleInput({ apiFormat: 'openai_chat' }))
await svc.activateProvider(provider.id)
const { releaseWrite, writeBlocked } = blockNextTraceAppend()
let released = false
let bodyPromise: Promise<string> | undefined
try {
const req = new Request('http://localhost:3456/proxy/v1/messages', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'X-Claude-Code-Session-Id': 'session-stream-background-trace',
},
body: JSON.stringify({
model: 'gpt-4',
max_tokens: 64,
stream: true,
messages: [{ role: 'user', content: 'finish before trace persistence' }],
}),
})
const response = await handleProxyRequest(req, new URL(req.url))
bodyPromise = response.text()
await writeBlocked
const body = await settlesBeforeBlockedTraceWrite(bodyPromise)
expect(body).not.toBeNull()
expect(body).toContain('message_stop')
releaseWrite()
released = true
const trace = await waitForCompletedProxyTrace('session-stream-background-trace')
expect(trace.calls[0]?.response?.body.preview).toContain('message_stop')
expect(trace.events.at(-1)?.phase).toBe('upstream_fetch_completed')
} finally {
if (!released) releaseWrite()
await bodyPromise?.catch(() => undefined)
globalThis.fetch = originalFetch
}
})
test('strips leading billing attribution instead of injecting it for OpenAI-compatible upstreams', async () => {
const originalFetch = globalThis.fetch
const calls: Array<{ body: Record<string, unknown> }> = []
+33 -27
View File
@@ -301,7 +301,7 @@ async function handleOpenaiChat(
}, networkSettings.aiRequestTimeoutMs, isStream)
} catch (err) {
if (traceContext) {
await recordProxyTrace({
recordProxyTraceInBackground({
callId: traceCallId,
context: traceContext,
model: body.model,
@@ -327,7 +327,7 @@ async function handleOpenaiChat(
},
}
if (traceContext) {
await recordProxyTrace({
recordProxyTraceInBackground({
context: traceContext,
callId: traceCallId,
model: body.model,
@@ -351,7 +351,7 @@ async function handleOpenaiChat(
if (isStream) {
if (!upstream.body) {
if (traceContext) {
await recordProxyTrace({
recordProxyTraceInBackground({
callId: traceCallId,
context: traceContext,
model: body.model,
@@ -402,7 +402,7 @@ async function handleOpenaiChat(
const responseBody = await upstream.json()
const anthropicResponse = openaiChatToAnthropic(responseBody, body.model)
if (traceContext) {
await recordProxyTrace({
recordProxyTraceInBackground({
callId: traceCallId,
context: traceContext,
model: body.model,
@@ -470,7 +470,7 @@ async function handleOpenaiResponses(
}, networkSettings.aiRequestTimeoutMs, isStream)
} catch (err) {
if (traceContext) {
await recordProxyTrace({
recordProxyTraceInBackground({
callId: traceCallId,
context: traceContext,
model: body.model,
@@ -496,7 +496,7 @@ async function handleOpenaiResponses(
},
}
if (traceContext) {
await recordProxyTrace({
recordProxyTraceInBackground({
context: traceContext,
callId: traceCallId,
model: body.model,
@@ -520,7 +520,7 @@ async function handleOpenaiResponses(
if (isStream) {
if (!upstream.body) {
if (traceContext) {
await recordProxyTrace({
recordProxyTraceInBackground({
callId: traceCallId,
context: traceContext,
model: body.model,
@@ -571,7 +571,7 @@ async function handleOpenaiResponses(
const responseBody = await upstream.json()
const anthropicResponse = openaiResponsesToAnthropic(responseBody, body.model)
if (traceContext) {
await recordProxyTrace({
recordProxyTraceInBackground({
callId: traceCallId,
context: traceContext,
model: body.model,
@@ -672,6 +672,27 @@ function startProxyTraceCall({
return callId
}
type RecordProxyTraceInput = {
callId?: string
context: ProxyTraceContext
model: string
upstreamUrl: string
upstreamRequest: unknown
requestHeaders?: Record<string, string>
startedAt: string
startedAtMs: number
responseStatus?: number
upstreamResponseBody?: unknown
anthropicResponseBody?: unknown
responseBodySnapshot?: TraceBodySnapshot
responseHeaders?: Headers
error?: unknown
}
function recordProxyTraceInBackground(input: RecordProxyTraceInput): void {
void recordProxyTrace(input).catch(() => {})
}
async function recordProxyTrace({
callId,
context,
@@ -687,22 +708,7 @@ async function recordProxyTrace({
responseBodySnapshot,
responseHeaders,
error,
}: {
callId?: string
context: ProxyTraceContext
model: string
upstreamUrl: string
upstreamRequest: unknown
requestHeaders?: Record<string, string>
startedAt: string
startedAtMs: number
responseStatus?: number
upstreamResponseBody?: unknown
anthropicResponseBody?: unknown
responseBodySnapshot?: TraceBodySnapshot
responseHeaders?: Headers
error?: unknown
}): Promise<void> {
}: RecordProxyTraceInput): Promise<void> {
const completedAt = new Date().toISOString()
const requestBody = createProxyTraceRequestBody(context, upstreamRequest)
const responseBody = anthropicResponseBody === undefined && upstreamResponseBody === undefined
@@ -797,11 +803,11 @@ function captureTraceStream(
captureChunk(value)
controller.enqueue(value)
}
await finalize()
controller.close()
void finalize()
} catch (err) {
await finalize(err)
controller.error(err)
void finalize(err)
} finally {
reader?.releaseLock()
reader = null
@@ -811,7 +817,7 @@ function captureTraceStream(
const error = reason instanceof Error
? reason
: new Error(reason ? `Stream cancelled: ${String(reason)}` : 'Stream cancelled')
await finalize(error)
void finalize(error)
await reader?.cancel(reason).catch(() => undefined)
},
})