diff --git a/packages/agent/CHANGELOG.md b/packages/agent/CHANGELOG.md index 10e1a6e21..df8d5fc24 100644 --- a/packages/agent/CHANGELOG.md +++ b/packages/agent/CHANGELOG.md @@ -1,6 +1,11 @@ # Changelog ## [Unreleased] +### Added + +- Added `responseHeaders` to `ChatUsageEvent` and `ManualChatTelemetryOptions` so telemetry hooks receive captured lowercase upstream response headers for each chat span +- Added automatic gateway/proxy detection from response headers (`litellm`, `helicone`, `portkey`, `openrouter`) and stamped `pi.gen_ai.gateway.*` span attributes for detected routing metadata +- Added exported `detectGatewayFromHeaders` API for header-based gateway detection ## [15.1.0] - 2026-05-15 ### Breaking Changes diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index 59785a981..066e66e90 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -699,10 +699,22 @@ async function streamAssistantResponse( }, }); + // Wrap the user-supplied onResponse so we always observe response headers + // for telemetry (`ChatUsageEvent.headers`, gateway auto-detection) without + // stealing them from the configured hook. + let capturedHeaders: Readonly> | undefined; + const userOnResponse = config.onResponse; + const captureOnResponse: AgentLoopConfig["onResponse"] = (response, modelInfo) => { + capturedHeaders = response.headers; + return userOnResponse?.(response, modelInfo); + }; + const finishChat = async (message: AssistantMessage): Promise => { await finishChatSpan(telemetry, chatSpan, message, { stepNumber: chatStepNumber, serviceTier: config.serviceTier, + responseHeaders: capturedHeaders, + baseUrl: config.model.baseUrl, }); }; @@ -716,6 +728,7 @@ async function streamAssistantResponse( reasoning: effectiveReasoning, temperature: effectiveTemperature, signal: requestSignal, + onResponse: captureOnResponse, }); let partialMessage: AssistantMessage | null = null; @@ -823,7 +836,11 @@ async function streamAssistantResponse( return trailing; }); } catch (err) { - failChatSpan(telemetry, chatSpan, { errorObject: err }); + failChatSpan(telemetry, chatSpan, { + errorObject: err, + responseHeaders: capturedHeaders, + baseUrl: config.model.baseUrl, + }); throw err; } } diff --git a/packages/agent/src/telemetry.ts b/packages/agent/src/telemetry.ts index 4ee4d20e4..a1a41f094 100644 --- a/packages/agent/src/telemetry.ts +++ b/packages/agent/src/telemetry.ts @@ -139,6 +139,14 @@ export const enum PiGenAIAttr { HandoffFromAgentId = "pi.gen_ai.handoff.from_agent.id", HandoffToAgentName = "pi.gen_ai.handoff.to_agent.name", HandoffToAgentId = "pi.gen_ai.handoff.to_agent.id", + // Gateway / proxy (LiteLLM, Helicone, Portkey, …) — populated when a known + // gateway header pattern is detected on the upstream response. The base + // `gen_ai.provider.name` continues to track the *upstream* provider (e.g. + // `anthropic`) that the gateway routed to. + GatewayName = "pi.gen_ai.gateway.name", + GatewayEndpoint = "pi.gen_ai.gateway.endpoint", + GatewayCallId = "pi.gen_ai.gateway.call_id", + GatewayRoutedTo = "pi.gen_ai.gateway.routed_to", } /** GenAI operation names — values for {@link GenAIAttr.OperationName}. */ @@ -223,6 +231,17 @@ export interface ChatUsageEvent { readonly cost: CostEstimate | undefined; /** Resolved dynamic attributes for this chat span (from `resolveAttributes`). */ readonly attributes: Attributes | undefined; + /** + * Response headers captured from the upstream HTTP response, with keys + * lowercased (mirrors {@link ProviderResponseMetadata.headers}). `undefined` + * when the provider transport did not surface headers (non-HTTP providers, + * mocked streams, requests that aborted before headers arrived). + * + * Use this to reconcile gateway-issued ids (e.g. `x-litellm-call-id`) with + * downstream billing / spend dashboards. Known gateway patterns are also + * auto-stamped on the chat span as `pi.gen_ai.gateway.*` attributes. + */ + readonly headers: Readonly> | undefined; } export type TelemetryContentCapture = boolean | "none" | "summary" | "full"; @@ -1078,11 +1097,17 @@ export async function finishChatSpan( telemetry: AgentTelemetry | undefined, span: Span | undefined, message: AssistantMessage, - options: { readonly stepNumber: number; readonly serviceTier?: ServiceTier }, + options: { + readonly stepNumber: number; + readonly serviceTier?: ServiceTier; + readonly responseHeaders?: Readonly>; + readonly baseUrl?: string; + }, ): Promise { if (!span) return; applyChatResponseAttributes(span, message); applyUsageAttributes(span, message.usage); + applyGatewayAttributes(span, options.responseHeaders, options.baseUrl); const cost = applyCostEstimate(telemetry, span, message, options.serviceTier, options.stepNumber); if (telemetry) { await emitChatUsage(telemetry, span, { @@ -1092,6 +1117,7 @@ export async function finishChatSpan( stepNumber: options.stepNumber, usage: message.usage, applied: cost, + headers: options.responseHeaders, }).catch(err => { emitTelemetryWarning(telemetry, { code: "on_chat_usage_failed", @@ -1125,9 +1151,15 @@ export async function finishChatSpan( export function failChatSpan( telemetry: AgentTelemetry | undefined, span: Span | undefined, - options: { readonly errorObject: unknown; readonly errorType?: string }, + options: { + readonly errorObject: unknown; + readonly errorType?: string; + readonly responseHeaders?: Readonly>; + readonly baseUrl?: string; + }, ): void { if (!span) return; + applyGatewayAttributes(span, options.responseHeaders, options.baseUrl); const err = options.errorObject; if (err instanceof Error) { span.recordException(err); @@ -1172,6 +1204,74 @@ function applyUsageAttributes(span: Span, usage: Usage | undefined): void { } } +/** + * Result of {@link detectGatewayFromHeaders}. `callId` and `routedTo` are + * populated only when the gateway surfaces them; consumers should treat + * `undefined` as "unknown for this gateway" rather than "no value". + */ +export interface GatewayHeaderDetection { + readonly name: string; + readonly callId: string | undefined; + readonly routedTo: string | undefined; +} + +/** + * Identify a known LLM gateway / proxy from response headers (LiteLLM, + * Helicone, Portkey). Returns `undefined` when no recognizable pattern is + * present so direct-API traffic stays unaffected. + * + * Header keys are matched case-insensitively against the lowercased map that + * {@link ProviderResponseMetadata.headers} produces. + */ +export function detectGatewayFromHeaders( + headers: Readonly> | undefined, +): GatewayHeaderDetection | undefined { + if (!headers) return undefined; + const litellmCallId = headers["x-litellm-call-id"]; + if (litellmCallId) { + return { + name: "litellm", + callId: litellmCallId, + routedTo: headers["x-litellm-model-id"] ?? headers["x-litellm-model-group"], + }; + } + const heliconeId = headers["helicone-id"]; + if (heliconeId) { + return { name: "helicone", callId: heliconeId, routedTo: headers["helicone-target-provider"] }; + } + const portkeyId = headers["x-portkey-trace-id"] ?? headers["x-portkey-request-id"]; + if (portkeyId) { + return { + name: "portkey", + callId: portkeyId, + routedTo: headers["x-portkey-llm-provider"] ?? headers["x-portkey-provider"], + }; + } + const openRouterGenerationId = headers["x-generation-id"]; + if (openRouterGenerationId?.startsWith("gen-")) { + // OpenRouter does not surface the upstream provider in response headers + // (only the body's `provider` field carries it), so `routedTo` is left + // undefined here. The `gen-` prefix on `x-generation-id` is OpenRouter- + // specific and disambiguates from other proxies that also expose a + // `x-generation-id` header. + return { name: "openrouter", callId: openRouterGenerationId, routedTo: undefined }; + } + return undefined; +} + +function applyGatewayAttributes( + span: Span, + headers: Readonly> | undefined, + baseUrl: string | undefined, +): void { + const gateway = detectGatewayFromHeaders(headers); + if (!gateway) return; + span.setAttribute(PiGenAIAttr.GatewayName, gateway.name); + if (baseUrl) span.setAttribute(PiGenAIAttr.GatewayEndpoint, baseUrl); + if (gateway.callId) span.setAttribute(PiGenAIAttr.GatewayCallId, gateway.callId); + if (gateway.routedTo) span.setAttribute(PiGenAIAttr.GatewayRoutedTo, gateway.routedTo); +} + interface AppliedCostEstimate { readonly costUsd: number | undefined; readonly inputUsd: number | undefined; @@ -1314,6 +1414,7 @@ async function emitChatUsage( readonly stepNumber: number | undefined; readonly usage: Usage | undefined; readonly applied: AppliedCostEstimate; + readonly headers: Readonly> | undefined; }, ): Promise { const hook = telemetry.config.onChatUsage; @@ -1332,6 +1433,7 @@ async function emitChatUsage( telemetry, buildTelemetryAttributeContext(telemetry, "chat", { stepNumber: input.stepNumber }), ), + headers: input.headers, }; try { await hook(event); @@ -1403,6 +1505,7 @@ export interface ManualChatTelemetryOptions { readonly responseText?: string; readonly responseToolCalls?: readonly ManualChatToolCallTelemetry[]; readonly attributes?: Attributes; + readonly responseHeaders?: Readonly>; readonly endSpan?: boolean; } @@ -1427,6 +1530,7 @@ export async function recordManualChatTelemetry( const finishReason = mapStopReason(options.finishReason); if (finishReason) span.setAttribute(GenAIAttr.ResponseFinishReasons, [finishReason]); applyUsageAttributes(span, options.usage); + applyGatewayAttributes(span, options.responseHeaders, options.model.baseUrl); if (telemetry) { const applied = applyCostEstimateForUsage(telemetry, span, { model: options.responseModel ?? options.model.id, @@ -1442,6 +1546,7 @@ export async function recordManualChatTelemetry( stepNumber: options.stepNumber, usage: options.usage, applied, + headers: options.responseHeaders, }).catch(err => { emitTelemetryWarning(telemetry, { code: "on_chat_usage_failed", diff --git a/packages/agent/test/otel.test.ts b/packages/agent/test/otel.test.ts index 815f87bfe..df9892441 100644 --- a/packages/agent/test/otel.test.ts +++ b/packages/agent/test/otel.test.ts @@ -10,6 +10,7 @@ import { agentLoop } from "@oh-my-pi/pi-agent-core/agent-loop"; import { type AgentTelemetryConfig, type ChatUsageEvent, + detectGatewayFromHeaders, GenAIAttr, GenAIOperation, OpenAIAttr, @@ -783,3 +784,212 @@ describe("agent-loop OTEL instrumentation", () => { } }); }); + +describe("detectGatewayFromHeaders", () => { + it("identifies LiteLLM via x-litellm-call-id and resolves routed_to from x-litellm-model-id", () => { + const detection = detectGatewayFromHeaders({ + "x-litellm-call-id": "call-abc-123", + "x-litellm-model-id": "anthropic/claude-sonnet-4-7", + "x-litellm-version": "1.59.2", + }); + expect(detection).toEqual({ + name: "litellm", + callId: "call-abc-123", + routedTo: "anthropic/claude-sonnet-4-7", + }); + }); + + it("falls back to x-litellm-model-group when model-id is absent", () => { + const detection = detectGatewayFromHeaders({ + "x-litellm-call-id": "call-xyz", + "x-litellm-model-group": "claude-fast", + }); + expect(detection?.routedTo).toBe("claude-fast"); + }); + + it("identifies Helicone via helicone-id", () => { + const detection = detectGatewayFromHeaders({ + "helicone-id": "req_42", + "helicone-target-provider": "openai", + }); + expect(detection).toEqual({ name: "helicone", callId: "req_42", routedTo: "openai" }); + }); + + it("identifies Portkey via x-portkey-trace-id with provider routing", () => { + const detection = detectGatewayFromHeaders({ + "x-portkey-trace-id": "trace_99", + "x-portkey-llm-provider": "anthropic", + }); + expect(detection).toEqual({ name: "portkey", callId: "trace_99", routedTo: "anthropic" }); + }); + + it("identifies OpenRouter via x-generation-id with gen- prefix", () => { + expect(detectGatewayFromHeaders({ "x-generation-id": "gen-1234567890" })).toEqual({ + name: "openrouter", + callId: "gen-1234567890", + routedTo: undefined, + }); + }); + + it("ignores x-generation-id without the OpenRouter gen- prefix", () => { + expect(detectGatewayFromHeaders({ "x-generation-id": "1234567890" })).toBeUndefined(); + }); + + it("falls back to x-portkey-request-id when trace id is absent", () => { + expect(detectGatewayFromHeaders({ "x-portkey-request-id": "rq_1" })?.callId).toBe("rq_1"); + }); + + it("returns undefined when no known gateway header is present", () => { + expect(detectGatewayFromHeaders({})).toBeUndefined(); + expect( + detectGatewayFromHeaders({ + "content-type": "application/json", + "x-request-id": "rid", + }), + ).toBeUndefined(); + expect(detectGatewayFromHeaders(undefined)).toBeUndefined(); + }); + + it("prefers LiteLLM detection over Helicone when both header families are present", () => { + const detection = detectGatewayFromHeaders({ + "x-litellm-call-id": "ll-1", + "helicone-id": "he-1", + }); + expect(detection?.name).toBe("litellm"); + }); +}); + +describe("ChatUsageEvent.headers and pi.gen_ai.gateway.* span attributes", () => { + it("forwards captured response headers to onChatUsage", async () => { + const mock = createMockModel({ + ...MOCK_IDENT, + responses: [ + { + content: ["ok"], + usage: { input: 10, output: 5, totalTokens: 15 }, + responseHeaders: { + "X-Request-Id": "upstream-req-77", + "content-type": "application/json", + }, + }, + ], + }); + const events: ChatUsageEvent[] = []; + const config: AgentLoopConfig = { + model: mock.model, + convertToLlm: identityConverter, + telemetry: { onChatUsage: event => void events.push(event) }, + }; + const ctx: AgentContext = { systemPrompt: [], messages: [], tools: [] }; + await runAndDrain(agentLoop([createUserMessage("hi")], ctx, config, undefined, mock.stream)); + + expect(events).toHaveLength(1); + // Header keys are normalized to lowercase to match `ProviderResponseMetadata.headers`. + expect(events[0]?.headers).toEqual({ + "x-request-id": "upstream-req-77", + "content-type": "application/json", + }); + }); + + it("leaves ChatUsageEvent.headers undefined when the provider does not surface headers", async () => { + const mock = createMockModel({ + ...MOCK_IDENT, + responses: [{ content: ["ok"], usage: { input: 1, output: 1, totalTokens: 2 } }], + }); + const events: ChatUsageEvent[] = []; + const config: AgentLoopConfig = { + model: mock.model, + convertToLlm: identityConverter, + telemetry: { onChatUsage: event => void events.push(event) }, + }; + const ctx: AgentContext = { systemPrompt: [], messages: [], tools: [] }; + await runAndDrain(agentLoop([createUserMessage("hi")], ctx, config, undefined, mock.stream)); + + expect(events).toHaveLength(1); + expect(events[0]?.headers).toBeUndefined(); + }); + + it("auto-stamps pi.gen_ai.gateway.* on the chat span when LiteLLM headers are present", async () => { + const mock = createMockModel({ + ...MOCK_IDENT, + responses: [ + { + content: ["ok"], + usage: { input: 4, output: 2, totalTokens: 6 }, + responseHeaders: { + "x-litellm-call-id": "ll-call-abc", + "x-litellm-model-id": "anthropic/claude-sonnet-4-7", + "x-litellm-version": "1.59.2", + }, + }, + ], + }); + const config: AgentLoopConfig = { + model: mock.model, + convertToLlm: identityConverter, + telemetry: {}, + }; + const ctx: AgentContext = { systemPrompt: [], messages: [], tools: [] }; + await runAndDrain(agentLoop([createUserMessage("hi")], ctx, config, undefined, mock.stream)); + + const chat = findSpan(exporter.getFinishedSpans(), "chat mock-model"); + expect(chat?.attributes[PiGenAIAttr.GatewayName]).toBe("litellm"); + expect(chat?.attributes[PiGenAIAttr.GatewayCallId]).toBe("ll-call-abc"); + expect(chat?.attributes[PiGenAIAttr.GatewayRoutedTo]).toBe("anthropic/claude-sonnet-4-7"); + expect(chat?.attributes[PiGenAIAttr.GatewayEndpoint]).toBe(mock.model.baseUrl); + }); + + it("does not stamp gateway attributes when headers carry no known pattern", async () => { + const mock = createMockModel({ + ...MOCK_IDENT, + responses: [ + { + content: ["ok"], + usage: { input: 4, output: 2, totalTokens: 6 }, + responseHeaders: { "x-request-id": "rid-1" }, + }, + ], + }); + const config: AgentLoopConfig = { + model: mock.model, + convertToLlm: identityConverter, + telemetry: {}, + }; + const ctx: AgentContext = { systemPrompt: [], messages: [], tools: [] }; + await runAndDrain(agentLoop([createUserMessage("hi")], ctx, config, undefined, mock.stream)); + + const chat = findSpan(exporter.getFinishedSpans(), "chat mock-model"); + expect(chat?.attributes[PiGenAIAttr.GatewayName]).toBeUndefined(); + expect(chat?.attributes[PiGenAIAttr.GatewayCallId]).toBeUndefined(); + expect(chat?.attributes[PiGenAIAttr.GatewayEndpoint]).toBeUndefined(); + }); + + it("still invokes the user-supplied onResponse alongside header capture", async () => { + const seen: Array> = []; + const mock = createMockModel({ + ...MOCK_IDENT, + responses: [ + { + content: ["ok"], + usage: { input: 1, output: 1, totalTokens: 2 }, + responseHeaders: { "x-litellm-call-id": "ll-2" }, + }, + ], + }); + const config: AgentLoopConfig = { + model: mock.model, + convertToLlm: identityConverter, + onResponse: response => { + seen.push({ ...response.headers }); + }, + telemetry: {}, + }; + const ctx: AgentContext = { systemPrompt: [], messages: [], tools: [] }; + await runAndDrain(agentLoop([createUserMessage("hi")], ctx, config, undefined, mock.stream)); + + expect(seen).toHaveLength(1); + expect(seen[0]?.["x-litellm-call-id"]).toBe("ll-2"); + const chat = findSpan(exporter.getFinishedSpans(), "chat mock-model"); + expect(chat?.attributes[PiGenAIAttr.GatewayName]).toBe("litellm"); + }); +}); diff --git a/packages/ai/CHANGELOG.md b/packages/ai/CHANGELOG.md index 2fd454a41..f507c9431 100644 --- a/packages/ai/CHANGELOG.md +++ b/packages/ai/CHANGELOG.md @@ -1,13 +1,14 @@ # Changelog ## [Unreleased] - ### Breaking Changes - Rejected draft-07 tuple and dependency keywords (`items` arrays, `dependencies`, `additionalItems`) in JSON Schema validation ### Added +- Added `responseHeaders`, `responseStatus`, and `responseRequestId` fields to `MockResponse` so mock providers can provide synthetic `ProviderResponseMetadata` +- Added `onResponse` metadata emission for mocks that sends lowercased headers and a default status of 200 before streaming when response headers are configured - Added recursive strict-mode sanitization for array `prefixItems` entries so tuple schemas now enforce object constraints per item ### Changed diff --git a/packages/ai/src/providers/mock.ts b/packages/ai/src/providers/mock.ts index 5dae1b980..7c2228526 100644 --- a/packages/ai/src/providers/mock.ts +++ b/packages/ai/src/providers/mock.ts @@ -89,6 +89,16 @@ export interface MockResponse { throw?: string | Error; /** Delay before any event is emitted. Honors the call's AbortSignal. */ delayMs?: number; + /** + * If set, the mock synthesizes a {@link ProviderResponseMetadata} and fires + * `options.onResponse` once before streaming events. Headers are forwarded + * verbatim (keys lowercased to match real provider plumbing). + */ + responseHeaders?: Readonly>; + /** HTTP status code paired with {@link responseHeaders}. Defaults to 200. */ + responseStatus?: number; + /** Pre-set requestId surfaced via {@link ProviderResponseMetadata.requestId}. */ + responseRequestId?: string; } /** Handler resolved per call: static script or function. */ @@ -301,6 +311,26 @@ async function runMock( return; } + if (response.responseHeaders && options?.onResponse) { + const headers: Record = {}; + for (const [key, value] of Object.entries(response.responseHeaders)) { + headers[key.toLowerCase()] = value; + } + try { + await options.onResponse( + { + status: response.responseStatus ?? 200, + headers, + ...(response.responseRequestId !== undefined ? { requestId: response.responseRequestId } : {}), + }, + model, + ); + } catch (err) { + stream.fail(err); + return; + } + } + if (response.delayMs && response.delayMs > 0) { try { await sleep(response.delayMs, options?.signal);