import { buildProbeContext } from '../probe/context.js' import { probe as defaultProbeRegistry } from '../probe/registry.js' import type { ProbeObservation } from '../probe/registry.js' import type { ProviderCallId, ProviderCallUsage } from '../types/bus/index.js' import type { TokenUsage } from '../types/common/index.js' import type { SessionId, TurnId } from '../types/ids/index.js' import type { ChatCompletionParams } from '../types/provider/chat.js' import type { LLMProvider } from '../types/provider/interface.js' import type { StreamChunk } from '../types/provider/stream.js' export interface ProviderInstrumentationOptions { /** Observation only — a provider wrapper records, it never refuses. */ readonly probes?: ProbeObservation /** The session whose calls this wrapper observes. */ readonly sessionId?: SessionId /** The turn whose calls this wrapper observes. */ readonly turnId?: TurnId } let providerCallCounter = 0 function nextCallId(): ProviderCallId { providerCallCounter += 1 return `pcall_${Date.now().toString(36)}${providerCallCounter.toString(36)}` as ProviderCallId } function extractStreamUsage(usage: TokenUsage | undefined): ProviderCallUsage | undefined { if (!usage) return undefined const u = usage as TokenUsage & Partial return { inputTokens: u.inputTokens ?? u.promptTokens, outputTokens: u.outputTokens ?? u.completionTokens, totalTokens: u.totalTokens, costUsd: u.costUsd, } } export function wrapProviderWithProbes( provider: LLMProvider, opts: ProviderInstrumentationOptions = {}, ): LLMProvider { const probes = opts.probes ?? defaultProbeRegistry const { sessionId, turnId } = opts const scope = { ...(sessionId !== undefined ? { sessionId } : {}), ...(turnId !== undefined ? { turnId } : {}), } const wrapped: LLMProvider = { id: provider.id, name: provider.name, listModels: provider.listModels?.bind(provider), probeCredential: provider.probeCredential?.bind(provider), healthCheck: provider.healthCheck?.bind(provider), async *chatStream(params: ChatCompletionParams): AsyncIterable { const callId = nextCallId() const ctx = buildProbeContext(scope) const startedAt = Date.now() probes.dispatch( { type: 'provider_call_start', providerId: provider.id, model: params.model, callId, ...scope, }, ctx, ) try { let lastUsage: TokenUsage | undefined for await (const chunk of provider.chatStream(params)) { if (chunk.usage) lastUsage = chunk.usage yield chunk } probes.dispatch( { type: 'provider_call_completed', providerId: provider.id, model: params.model, callId, ...scope, durationMs: Date.now() - startedAt, usage: extractStreamUsage(lastUsage), }, ctx, ) } catch (error) { probes.dispatch( { type: 'provider_call_failed', providerId: provider.id, model: params.model, callId, ...scope, durationMs: Date.now() - startedAt, error: error instanceof Error ? error.message : String(error), }, ctx, ) throw error } }, } return wrapped }