From 203f215ee57c4a41798c76430ff8f01cd3a959a9 Mon Sep 17 00:00:00 2001 From: JPeer264 Date: Fri, 2 Oct 2026 14:14:18 +0200 Subject: [PATCH] fix(server-utils): Map Flue messages and token usage to the gen_ai conventions Flue records its turn messages in pi-ai's shape, with a `toolResult` role and `toolCall` parts, and Sentry renders neither, so tool calls and tool results were missing from the conversation views. Its token counts come from pi-ai too, whose `input` leaves the cached tokens out while the conventions count them in, so input tokens and input cost were too low whenever prompt caching was active. A shared pi-ai mapper now translates messages, usage and the finish reason to the conventions, and the tool span records the result the model receives instead of Flue's internal result wrapper. Co-Authored-By: Claude Opus 5.5 Co-Authored-By: Claude Fable 5.1 --- .../tracing/flue/instrument-with-pii.mjs | 10 ++ .../suites/tracing/flue/test.ts | 66 ++++++- packages/server-utils/src/ai/flue/index.ts | 29 +-- packages/server-utils/src/ai/flue/types.ts | 21 +-- packages/server-utils/src/ai/flue/utils.ts | 97 +++++----- .../server-utils/src/ai/pi-ai/messages.ts | 108 +++++++++++ .../server-utils/src/ai/pi-ai/providers.ts | 27 +++ packages/server-utils/src/ai/pi-ai/usage.ts | 75 ++++++++ .../test/ai/lib/tracing/flue.test.ts | 169 +++++++++++++++++- .../test/ai/lib/utils/pi-ai-messages.test.ts | 141 +++++++++++++++ .../test/ai/lib/utils/pi-ai-usage.test.ts | 62 +++++++ 11 files changed, 710 insertions(+), 95 deletions(-) create mode 100644 dev-packages/node-integration-tests/suites/tracing/flue/instrument-with-pii.mjs create mode 100644 packages/server-utils/src/ai/pi-ai/messages.ts create mode 100644 packages/server-utils/src/ai/pi-ai/providers.ts create mode 100644 packages/server-utils/src/ai/pi-ai/usage.ts create mode 100644 packages/server-utils/test/ai/lib/utils/pi-ai-messages.test.ts create mode 100644 packages/server-utils/test/ai/lib/utils/pi-ai-usage.test.ts diff --git a/dev-packages/node-integration-tests/suites/tracing/flue/instrument-with-pii.mjs b/dev-packages/node-integration-tests/suites/tracing/flue/instrument-with-pii.mjs new file mode 100644 index 000000000000..1c507ba84b9c --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/flue/instrument-with-pii.mjs @@ -0,0 +1,10 @@ +import * as Sentry from '@sentry/node'; +import { loggingTransport } from '@sentry-internal/node-integration-tests'; + +Sentry.init({ + dsn: 'https://public@dsn.ingest.sentry.io/1337', + release: '1.0', + tracesSampleRate: 1.0, + dataCollection: { genAI: { inputs: true, outputs: true } }, + transport: loggingTransport, +}); diff --git a/dev-packages/node-integration-tests/suites/tracing/flue/test.ts b/dev-packages/node-integration-tests/suites/tracing/flue/test.ts index 98ad59dcc4df..f4cf35579ae4 100644 --- a/dev-packages/node-integration-tests/suites/tracing/flue/test.ts +++ b/dev-packages/node-integration-tests/suites/tracing/flue/test.ts @@ -2,8 +2,13 @@ import { GEN_AI_AGENT_NAME, GEN_AI_CONVERSATION_ID, GEN_AI_COST_TOTAL_TOKENS, + GEN_AI_INPUT_MESSAGES, GEN_AI_OPERATION_NAME, + GEN_AI_OUTPUT_MESSAGES, GEN_AI_RESPONSE_FINISH_REASONS, + GEN_AI_SYSTEM_INSTRUCTIONS, + GEN_AI_TOOL_CALL_ARGUMENTS, + GEN_AI_TOOL_CALL_RESULT, GEN_AI_TOOL_NAME, GEN_AI_USAGE_INPUT_TOKENS, GEN_AI_USAGE_OUTPUT_TOKENS, @@ -87,9 +92,10 @@ conditionalTest({ min: 22 })('Flue integration', () => { expect(chat.attributes[GEN_AI_COST_TOTAL_TOKENS]?.value).toEqual(expect.any(Number)); } + // pi-ai's `toolUse` is reported as the conventions' `tool_call`. expect(chats.map(chat => chat.attributes[GEN_AI_RESPONSE_FINISH_REASONS]?.value).sort()).toEqual([ '["stop"]', - '["toolUse"]', + '["tool_call"]', ]); const tool = tools[0]!; @@ -111,4 +117,62 @@ conditionalTest({ min: 22 })('Flue integration', () => { }, FLUE_DEPENDENCIES, ); + + createEsmAndCjsTests( + __dirname, + 'scenario.mjs', + 'instrument-with-pii.mjs', + (createRunner, test, mode) => { + if (mode === 'cjs') { + return; + } + + test('records the messages of the turns in the gen_ai conventions shape', async () => { + await createRunner() + .expect({ + span: container => { + const chats = container.items + .filter(span => span.name === 'chat faux-model') + .sort((a, b) => a.start_timestamp - b.start_timestamp); + expect(chats).toHaveLength(2); + + // Flue sends pi-ai's shape (`toolResult`, `toolCall`); Sentry renders the conventions. + expect(JSON.parse(String(chats[0]!.attributes[GEN_AI_INPUT_MESSAGES]?.value))).toEqual([ + { role: 'user', parts: [{ type: 'text', content: 'What is the weather in Berlin?' }] }, + ]); + expect(JSON.parse(String(chats[0]!.attributes[GEN_AI_OUTPUT_MESSAGES]?.value))).toEqual([ + { + role: 'assistant', + parts: [{ type: 'tool_call', id: 'call_1', name: 'get_weather', arguments: '{"city":"Berlin"}' }], + finish_reason: 'tool_call', + }, + ]); + + const answerInput = JSON.parse(String(chats[1]!.attributes[GEN_AI_INPUT_MESSAGES]?.value)); + expect(answerInput.map((message: { role: string }) => message.role)).toEqual([ + 'user', + 'assistant', + 'tool', + ]); + expect(answerInput[2].parts[0]).toMatchObject({ + type: 'tool_call_response', + id: 'call_1', + name: 'get_weather', + }); + expect(answerInput[2].parts[0].result).toContain('sunny in Berlin'); + expect(chats[1]!.attributes[GEN_AI_SYSTEM_INSTRUCTIONS]?.value).toContain('You are a helpful assistant.'); + + // The tool span records the result the model receives, not Flue's internal wrapper. + const tool = container.items.find(span => span.name === 'execute_tool get_weather')!; + expect(tool.attributes[GEN_AI_TOOL_CALL_ARGUMENTS]?.value).toBe('{"city":"Berlin"}'); + expect(tool.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).toContain('sunny in Berlin'); + expect(tool.attributes[GEN_AI_TOOL_CALL_RESULT]?.value).not.toContain('details'); + }, + }) + .start() + .completed(); + }); + }, + FLUE_DEPENDENCIES, + ); }); diff --git a/packages/server-utils/src/ai/flue/index.ts b/packages/server-utils/src/ai/flue/index.ts index 67e59c9fcdc6..ffbab2f25659 100644 --- a/packages/server-utils/src/ai/flue/index.ts +++ b/packages/server-utils/src/ai/flue/index.ts @@ -1,7 +1,5 @@ import type { Span } from '@sentry/core'; import { - _INTERNAL_shouldSkipAiProviderWrapping, - _INTERNAL_skipAiProviderWrapping, continueTrace, getActiveSpan, LRUMap, @@ -10,11 +8,9 @@ import { withActiveSpan, } from '@sentry/core'; import { GEN_AI_AGENT_NAME, GEN_AI_CONVERSATION_ID, GEN_AI_OPERATION_NAME } from '@sentry/conventions/attributes'; -import { ANTHROPIC_AI_INTEGRATION_NAME } from '../anthropic-ai/constants'; import type { GenAiOptions } from '../core/utils'; import { getGenAiSpanOp, resolveAIRecordingOptions } from '../core/utils'; -import { GOOGLE_GENAI_INTEGRATION_NAME } from '../google-genai/constants'; -import { OPENAI_INTEGRATION_NAME } from '../openai/constants'; +import { skipPiAiProviderIntegrations } from '../pi-ai/providers'; import { FLUE_INSTRUMENTATION_KEY, FLUE_OPERATION, FLUE_ORIGIN, MAX_TRACKED_FLUE_SPANS } from './constants'; import type { SpanTracker } from './utils'; import { @@ -29,8 +25,6 @@ import type { FlueInstrumentation } from './types'; export type FlueOptions = GenAiOptions; -const SKIPPED_PROVIDERS = [OPENAI_INTEGRATION_NAME, ANTHROPIC_AI_INTEGRATION_NAME, GOOGLE_GENAI_INTEGRATION_NAME]; - /** * Build the object to hand to `instrument()` from `@flue/runtime`. * @@ -45,20 +39,13 @@ const SKIPPED_PROVIDERS = [OPENAI_INTEGRATION_NAME, ANTHROPIC_AI_INTEGRATION_NAM * which fall back to the current client's `dataCollection.genAI` settings and are read per event. */ export function createFlueInstrumentation(options: FlueOptions = {}): FlueInstrumentation { - // Flue drives the providers through `@earendil-works/pi-ai`, which bundles the `openai`, - // `@anthropic-ai/sdk` and `@google/genai` clients. Left alone they instrument the same call this - // reports as a turn, emitting a second `gen_ai.chat` beside ours. - // - // Applied on first use rather than here, for two reasons. Constructing the object proves nothing - // — if `instrument()` rejects it, suppressing the provider integrations would leave the app with - // no `gen_ai.chat` spans at all. And the registry is reset per client (`_setupIntegrations` - // clears it, and Cloudflare calls `init()` per request), so a one-shot call at module scope is - // wiped by the next `init()` and every later request double-reports. - const skipProviders = (): void => { - if (!SKIPPED_PROVIDERS.every(provider => _INTERNAL_shouldSkipAiProviderWrapping(provider))) { - _INTERNAL_skipAiProviderWrapping(SKIPPED_PROVIDERS); - } - }; + // Flue drives the providers through `@earendil-works/pi-ai`. The provider integrations are skipped + // on first use rather than here, for two reasons. Constructing the object proves nothing: if + // `instrument()` rejects it, suppressing the provider integrations would leave the app with no + // `gen_ai.chat` spans at all. And the registry is reset per client (`_setupIntegrations` clears it, + // and Cloudflare calls `init()` per request), so a one-shot call at module scope is wiped by the + // next `init()` and every later request double-reports. + const skipProviders = skipPiAiProviderIntegrations; // Keyed by the agent operation's own id, which is what the observations carry. That keeps // concurrent runs apart and gives a delegated subagent its own span: Flue nests a second `agent` diff --git a/packages/server-utils/src/ai/flue/types.ts b/packages/server-utils/src/ai/flue/types.ts index bdbd58455518..30850dea9949 100644 --- a/packages/server-utils/src/ai/flue/types.ts +++ b/packages/server-utils/src/ai/flue/types.ts @@ -6,21 +6,7 @@ * `FlueExecutionContext` as of `@flue/runtime` 2.x. */ -/** Token counts and Flue-computed costs on a settled turn. */ -export interface FlueUsage { - input?: number; - output?: number; - cacheRead?: number; - cacheWrite?: number; - totalTokens?: number; - cost?: { - input?: number; - output?: number; - cacheRead?: number; - cacheWrite?: number; - total?: number; - }; -} +import type { PiAiUsage } from '../pi-ai/usage'; /** Mirrors `ModelRequestInfo`. */ export interface FlueModelRequestInfo { @@ -51,7 +37,8 @@ export interface FlueModelResponse { responseId?: string; responseModel?: string; output?: unknown; - usage?: FlueUsage; + /** pi-ai's usage, which Flue passes through with its own cost figures. */ + usage?: PiAiUsage; finishReason?: string; } @@ -88,6 +75,8 @@ export interface FlueObservation { request?: FlueModelRequest; args?: unknown; result?: unknown; + /** The tool result as the model sees it, on successful `tool` events. */ + effectiveResult?: unknown; response?: FlueModelResponse; } diff --git a/packages/server-utils/src/ai/flue/utils.ts b/packages/server-utils/src/ai/flue/utils.ts index a1b463cbfbd7..6c7d3ea34ccf 100644 --- a/packages/server-utils/src/ai/flue/utils.ts +++ b/packages/server-utils/src/ai/flue/utils.ts @@ -1,6 +1,7 @@ import type { LRUMap, Span } from '@sentry/core'; import { captureException, + isObjectLike, SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN, SPAN_STATUS_ERROR, startInactiveSpan, @@ -9,11 +10,6 @@ import { } from '@sentry/core'; import { GEN_AI_CONVERSATION_ID, - GEN_AI_COST_CACHE_CREATION_INPUT_TOKENS, - GEN_AI_COST_CACHE_READ_INPUT_TOKENS, - GEN_AI_COST_INPUT_TOKENS, - GEN_AI_COST_OUTPUT_TOKENS, - GEN_AI_COST_TOTAL_TOKENS, GEN_AI_INPUT_MESSAGES, GEN_AI_OPERATION_NAME, GEN_AI_OUTPUT_MESSAGES, @@ -30,17 +26,19 @@ import { GEN_AI_TOOL_CALL_RESULT, GEN_AI_TOOL_DEFINITIONS, GEN_AI_TOOL_NAME, - GEN_AI_USAGE_CACHE_CREATION_INPUT_TOKENS, - GEN_AI_USAGE_CACHE_READ_INPUT_TOKENS, - GEN_AI_USAGE_INPUT_TOKENS, - GEN_AI_USAGE_OUTPUT_TOKENS, - GEN_AI_USAGE_TOTAL_TOKENS, SERVER_ADDRESS, SERVER_PORT, } from '@sentry/conventions/attributes'; import { getGenAiSpanOp } from '../core/utils'; +import { + piAiAssistantMessageToGenAiMessage, + piAiContentToString, + piAiFinishReason, + piAiMessagesToGenAiMessages, +} from '../pi-ai/messages'; +import { setPiAiUsageAttributes } from '../pi-ai/usage'; import { FLUE_ORIGIN, MAX_TRACKED_FLUE_SPANS } from './constants'; -import type { FlueErrorInfo, FlueModelRequestInfo, FlueObservation, FlueUsage } from './types'; +import type { FlueErrorInfo, FlueModelRequestInfo, FlueObservation } from './types'; /** * Flue persists the incoming W3C `traceparent` at admission and replays it on the agent operation. @@ -163,22 +161,25 @@ export function endTurnSpan(observation: FlueObservation, turnSpans: SpanTracker setRequestAttributes(span, observation.request); - const { responseId, finishReason } = observation.response ?? {}; + const responseId = observation.response?.responseId; if (responseId) { span.setAttribute(GEN_AI_RESPONSE_ID, responseId); } + const finishReason = piAiFinishReason(observation.response?.finishReason); if (finishReason) { // Serialized, not a raw array: the conventions declare this attribute's value type as `string`, // and that is what `ai/core`, Mastra, OpenAI and Vercel AI all write. span.setAttribute(GEN_AI_RESPONSE_FINISH_REASONS, stringify([finishReason])); } - const output = observation.response?.output; - if (recordOutputs && output !== undefined) { - span.setAttribute(GEN_AI_OUTPUT_MESSAGES, stringify(output)); + const output = recordOutputs + ? piAiAssistantMessageToGenAiMessage(observation.response?.output, finishReason) + : undefined; + if (output) { + span.setAttribute(GEN_AI_OUTPUT_MESSAGES, stringify([output])); } - setUsageAttributes(span, observation.response?.usage, observation.isError); + setPiAiUsageAttributes(span, observation.response?.usage, observation.isError); if (observation.isError) { span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' }); @@ -186,39 +187,6 @@ export function endTurnSpan(observation: FlueObservation, turnSpans: SpanTracker span.end(); } -/** - * Flue reports token counts and its own computed costs on the same `usage` object, so both are set - * here. The cost figures have no equivalent in the provider SDKs' own instrumentation. - */ -export function setUsageAttributes(span: Span, usage: FlueUsage | undefined, isError?: boolean): void { - // A turn that failed before the provider billed anything reports every counter as 0. Writing - // those is noise that reads as a real zero-cost call, so skip the block entirely. - if (!usage || (isError && !usage.totalTokens)) { - return; - } - - const attributes: Record = {}; - const set = (key: string, value: number | undefined): void => { - if (typeof value === 'number') { - attributes[key] = value; - } - }; - - set(GEN_AI_USAGE_INPUT_TOKENS, usage.input); - set(GEN_AI_USAGE_OUTPUT_TOKENS, usage.output); - set(GEN_AI_USAGE_TOTAL_TOKENS, usage.totalTokens); - set(GEN_AI_USAGE_CACHE_READ_INPUT_TOKENS, usage.cacheRead); - set(GEN_AI_USAGE_CACHE_CREATION_INPUT_TOKENS, usage.cacheWrite); - - set(GEN_AI_COST_INPUT_TOKENS, usage.cost?.input); - set(GEN_AI_COST_OUTPUT_TOKENS, usage.cost?.output); - set(GEN_AI_COST_TOTAL_TOKENS, usage.cost?.total); - set(GEN_AI_COST_CACHE_READ_INPUT_TOKENS, usage.cost?.cacheRead); - set(GEN_AI_COST_CACHE_CREATION_INPUT_TOKENS, usage.cost?.cacheWrite); - - span.setAttributes(attributes); -} - /** * Tool spans hang off the agent invocation rather than the turn, matching how Flue's own * OpenTelemetry adapter projects them: siblings of `chat`, correlated to model output by tool call @@ -257,8 +225,11 @@ export function endToolSpan(observation: FlueObservation, toolSpans: SpanTracker } toolSpans.remove(toolCallId); - if (recordOutputs && observation.result !== undefined) { - span.setAttribute(GEN_AI_TOOL_CALL_RESULT, stringify(observation.result)); + if (recordOutputs) { + const result = toolResultForModel(observation); + if (result) { + span.setAttribute(GEN_AI_TOOL_CALL_RESULT, result); + } } if (observation.isError) { @@ -268,6 +239,25 @@ export function endToolSpan(observation: FlueObservation, toolSpans: SpanTracker span.end(); } +/** + * The result the model receives: Flue's `effectiveResult` when it has one, else the content blocks + * of its harness-level `result` (`{ content, details }`, whose `details` payload is tool-specific). + */ +function toolResultForModel(observation: FlueObservation): string | undefined { + const { effectiveResult } = observation; + if (effectiveResult !== undefined) { + if (typeof effectiveResult === 'string') { + return effectiveResult; + } + return Array.isArray(effectiveResult) ? piAiContentToString(effectiveResult) : stringify(effectiveResult); + } + const { result } = observation; + if (isObjectLike(result) && Array.isArray(result.content)) { + return piAiContentToString(result.content); + } + return result === undefined ? undefined : stringify(result); +} + /** * `turn_request` is the only event carrying the request's content — the settled `turn` reports * metadata alone — so input messages, system prompt and tool definitions are read from it. @@ -283,8 +273,9 @@ export function recordRequestContent(observation: FlueObservation, turnSpans: Sp if (input.systemPrompt) { span.setAttribute(GEN_AI_SYSTEM_INSTRUCTIONS, input.systemPrompt); } - if (input.messages) { - span.setAttribute(GEN_AI_INPUT_MESSAGES, stringify(input.messages)); + const messages = piAiMessagesToGenAiMessages(input.messages); + if (messages.length) { + span.setAttribute(GEN_AI_INPUT_MESSAGES, stringify(messages)); } if (input.tools?.length) { span.setAttribute(GEN_AI_TOOL_DEFINITIONS, stringify(input.tools)); diff --git a/packages/server-utils/src/ai/pi-ai/messages.ts b/packages/server-utils/src/ai/pi-ai/messages.ts new file mode 100644 index 000000000000..5e3408d6f354 --- /dev/null +++ b/packages/server-utils/src/ai/pi-ai/messages.ts @@ -0,0 +1,108 @@ +import { isObjectLike, stringify } from '@sentry/core'; + +/** A message in the gen_ai conventions shape, see https://develop.sentry.dev/sdk/telemetry/traces/modules/ai-agents/. */ +export interface GenAiMessage { + role: string; + parts: GenAiMessagePart[]; + finish_reason?: string; +} + +export type GenAiMessagePart = Record & { type: string }; + +/** + * Map pi-ai messages (`@earendil-works/pi-ai`, which Flue and pi-durable send) to the gen_ai + * conventions. Sentry renders neither pi-ai's `toolResult` role nor its `toolCall` parts. + * + * System messages are left out: their prompt belongs in `gen_ai.system_instructions`. + */ +export function piAiMessagesToGenAiMessages(messages: unknown): GenAiMessage[] { + if (!Array.isArray(messages)) { + return []; + } + + const mapped: GenAiMessage[] = []; + for (const message of messages) { + if (!isObjectLike(message) || typeof message.role !== 'string' || message.role === 'system') { + continue; + } + + if (message.role === 'toolResult') { + mapped.push({ + role: 'tool', + parts: [ + { + type: 'tool_call_response', + id: message.toolCallId, + name: message.toolName, + result: piAiContentToString(message.content), + }, + ], + }); + continue; + } + + const parts = piAiContentToParts(message.content); + if (parts.length) { + mapped.push({ role: message.role, parts }); + } + } + return mapped; +} + +/** The `gen_ai.output.messages` entry of a pi-ai assistant message, or `undefined` when it has no content. */ +export function piAiAssistantMessageToGenAiMessage(message: unknown, finishReason?: string): GenAiMessage | undefined { + const parts = isObjectLike(message) ? piAiContentToParts(message.content) : []; + if (!parts.length) { + return undefined; + } + return { role: 'assistant', parts, ...(finishReason ? { finish_reason: finishReason } : {}) }; +} + +/** pi-ai names a tool-calling stop `toolUse`; the conventions call it `tool_call`. */ +export function piAiFinishReason(stopReason: unknown): string | undefined { + if (typeof stopReason !== 'string' || !stopReason) { + return undefined; + } + return stopReason === 'toolUse' ? 'tool_call' : stopReason; +} + +/** Text-only content as its text, joined the way pi-ai joins it; anything else as its mapped parts. */ +export function piAiContentToString(content: unknown): string | undefined { + const parts = piAiContentToParts(content); + return parts.every(part => part.type === 'text') ? parts.map(part => part.content).join('\n') : stringify(parts); +} + +function piAiContentToParts(content: unknown): GenAiMessagePart[] { + if (typeof content === 'string') { + return content ? [{ type: 'text', content }] : []; + } + if (!Array.isArray(content)) { + return []; + } + return content.map(piAiPartToGenAiPart).filter((part): part is GenAiMessagePart => part !== undefined); +} + +/** An image is reported by its media type only: its `data` is base64, which the conventions require to be dropped. */ +function piAiPartToGenAiPart(part: unknown): GenAiMessagePart | undefined { + if (!isObjectLike(part)) { + return undefined; + } + + switch (part.type) { + case 'text': + // An empty text block makes Sentry render the message as "(no value)" instead of its tool calls. + return typeof part.text === 'string' && part.text.trim() ? { type: 'text', content: part.text } : undefined; + case 'thinking': + // A redacted thinking block carries an encrypted payload, not text. + return typeof part.thinking === 'string' && part.thinking && !part.redacted + ? { type: 'reasoning', content: part.thinking } + : undefined; + case 'toolCall': + return { type: 'tool_call', id: part.id, name: part.name, arguments: stringify(part.arguments ?? {}, String) }; + case 'image': + return { type: 'blob', mime_type: part.mimeType }; + default: + // Part kinds we don't know yet render as JSON rather than disappearing. + return { type: 'object', content: part }; + } +} diff --git a/packages/server-utils/src/ai/pi-ai/providers.ts b/packages/server-utils/src/ai/pi-ai/providers.ts new file mode 100644 index 000000000000..524afe79d573 --- /dev/null +++ b/packages/server-utils/src/ai/pi-ai/providers.ts @@ -0,0 +1,27 @@ +import { _INTERNAL_shouldSkipAiProviderWrapping, _INTERNAL_skipAiProviderWrapping } from '@sentry/core'; +import { ANTHROPIC_AI_INTEGRATION_NAME } from '../anthropic-ai/constants'; +import { GOOGLE_GENAI_INTEGRATION_NAME } from '../google-genai/constants'; +import { OPENAI_INTEGRATION_NAME } from '../openai/constants'; + +// pi-ai sends its requests through the `openai`, `@anthropic-ai/sdk` and `@google/genai` clients. +// Left alone, those integrations report the same request a second time beside the `chat` span of +// the framework that sent it. Bedrock requests go through `@aws-sdk/client-bedrock-runtime`, which +// `awsIntegration` still reports; that one has no skip yet. +const PI_AI_PROVIDER_INTEGRATIONS = [ + OPENAI_INTEGRATION_NAME, + ANTHROPIC_AI_INTEGRATION_NAME, + GOOGLE_GENAI_INTEGRATION_NAME, +]; + +/** + * Stop the provider SDK integrations from reporting the requests pi-ai sends. The skip is + * process-wide: the provider SDKs do not know which of their calls pi-ai made. + * + * The registry is reset per client, so callers apply this on the first request instead of once at + * setup, or the next `init()` would undo it. + */ +export function skipPiAiProviderIntegrations(): void { + if (!PI_AI_PROVIDER_INTEGRATIONS.every(provider => _INTERNAL_shouldSkipAiProviderWrapping(provider))) { + _INTERNAL_skipAiProviderWrapping(PI_AI_PROVIDER_INTEGRATIONS); + } +} diff --git a/packages/server-utils/src/ai/pi-ai/usage.ts b/packages/server-utils/src/ai/pi-ai/usage.ts new file mode 100644 index 000000000000..4bfb8869e72e --- /dev/null +++ b/packages/server-utils/src/ai/pi-ai/usage.ts @@ -0,0 +1,75 @@ +import type { Span } from '@sentry/core'; +import { + GEN_AI_COST_CACHE_CREATION_INPUT_TOKENS, + GEN_AI_COST_CACHE_READ_INPUT_TOKENS, + GEN_AI_COST_INPUT_TOKENS, + GEN_AI_COST_OUTPUT_TOKENS, + GEN_AI_COST_TOTAL_TOKENS, + GEN_AI_USAGE_CACHE_CREATION_INPUT_TOKENS, + GEN_AI_USAGE_CACHE_READ_INPUT_TOKENS, + GEN_AI_USAGE_INPUT_TOKENS, + GEN_AI_USAGE_OUTPUT_TOKENS, + GEN_AI_USAGE_REASONING_OUTPUT_TOKENS, + GEN_AI_USAGE_TOTAL_TOKENS, +} from '@sentry/conventions/attributes'; + +/** + * Token counts and costs of a pi-ai response. `input` and `cost.input` leave the cache tokens out; + * `reasoning` is part of `output`. + */ +export interface PiAiUsage { + input?: number; + output?: number; + cacheRead?: number; + cacheWrite?: number; + reasoning?: number; + totalTokens?: number; + cost?: { + input?: number; + output?: number; + cacheRead?: number; + cacheWrite?: number; + total?: number; + }; +} + +/** + * Set the token and cost attributes of a pi-ai response. The conventions count cached tokens in + * `gen_ai.usage.input_tokens` and their cost in `gen_ai.cost.input_tokens`, so both are summed here. + * + * `unbilled` is true for a response the provider may not have billed: a failed request, or one the + * provider parked to answer later. Its counters are all 0, which would read as a real zero-cost + * call, so nothing is set for it unless it reports tokens. + */ +export function setPiAiUsageAttributes(span: Span, usage: PiAiUsage | undefined, unbilled?: boolean): void { + if (!usage || (unbilled && !usage.totalTokens)) { + return; + } + + const attributes: Record = {}; + const set = (key: string, value: number | undefined): void => { + if (typeof value === 'number') { + attributes[key] = value; + } + }; + + set(GEN_AI_USAGE_INPUT_TOKENS, sum(usage.input, usage.cacheRead, usage.cacheWrite)); + set(GEN_AI_USAGE_OUTPUT_TOKENS, usage.output); + set(GEN_AI_USAGE_TOTAL_TOKENS, usage.totalTokens); + set(GEN_AI_USAGE_CACHE_READ_INPUT_TOKENS, usage.cacheRead); + set(GEN_AI_USAGE_CACHE_CREATION_INPUT_TOKENS, usage.cacheWrite); + set(GEN_AI_USAGE_REASONING_OUTPUT_TOKENS, usage.reasoning); + + set(GEN_AI_COST_INPUT_TOKENS, sum(usage.cost?.input, usage.cost?.cacheRead, usage.cost?.cacheWrite)); + set(GEN_AI_COST_OUTPUT_TOKENS, usage.cost?.output); + set(GEN_AI_COST_TOTAL_TOKENS, usage.cost?.total); + set(GEN_AI_COST_CACHE_READ_INPUT_TOKENS, usage.cost?.cacheRead); + set(GEN_AI_COST_CACHE_CREATION_INPUT_TOKENS, usage.cost?.cacheWrite); + + span.setAttributes(attributes); +} + +function sum(...values: (number | undefined)[]): number | undefined { + const numbers = values.filter((value): value is number => typeof value === 'number'); + return numbers.length ? numbers.reduce((total, value) => total + value, 0) : undefined; +} diff --git a/packages/server-utils/test/ai/lib/tracing/flue.test.ts b/packages/server-utils/test/ai/lib/tracing/flue.test.ts index 5817e1033871..57ad09690028 100644 --- a/packages/server-utils/test/ai/lib/tracing/flue.test.ts +++ b/packages/server-utils/test/ai/lib/tracing/flue.test.ts @@ -338,6 +338,55 @@ describe('createFlueInstrumentation', () => { expect(json?.data['gen_ai.cost.total_tokens']).toBe(0.001199); }); + // pi-ai keeps the cached tokens out of `input`; the conventions count them in. + it('counts cached tokens into the input tokens and their cost into the input cost', async () => { + const usage = { + input: 37, + output: 2, + cacheRead: 52, + cacheWrite: 38, + totalTokens: 129, + cost: { input: 0.000037, output: 0.00001, cacheRead: 0.0000052, cacheWrite: 0.0000475, total: 0.0000997 }, + }; + + await withAgent(() => { + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', operationId: 'op_1' }, {}); + instrumentation.observe(turn({ response: { ...turn().response, usage } }), {}); + }); + + const json = findSpan('chat claude-haiku-4.5'); + expect(json?.data['gen_ai.usage.input_tokens']).toBe(127); + expect(json?.data['gen_ai.usage.cache_read.input_tokens']).toBe(52); + expect(json?.data['gen_ai.usage.cache_creation.input_tokens']).toBe(38); + expect(json?.data['gen_ai.cost.input_tokens']).toBe(0.000037 + 0.0000052 + 0.0000475); + expect(json?.data['gen_ai.cost.total_tokens']).toBe(0.0000997); + }); + + it('reports a tool-calling turn with the finish reason of the conventions', async () => { + await withAgent(() => { + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', operationId: 'op_1' }, {}); + instrumentation.observe( + turn({ + response: { + ...turn().response, + finishReason: 'toolUse', + output: { + role: 'assistant', + content: [{ type: 'toolCall', id: 'c1', name: 'get_weather', arguments: { city: 'Berlin' } }], + }, + }, + }), + {}, + ); + }); + + const json = findSpan('chat claude-haiku-4.5'); + expect(json?.data['gen_ai.response.finish_reasons']).toBe('["tool_call"]'); + expect(json?.data['gen_ai.output.messages']).toBe( + '[{"role":"assistant","parts":[{"type":"tool_call","id":"c1","name":"get_weather","arguments":"{\\"city\\":\\"Berlin\\"}"}],"finish_reason":"tool_call"}]', + ); + }); + it('records the model-call tuning and provider endpoint', async () => { await withAgent(() => { instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', operationId: 'op_1', purpose: 'agent' }, {}); @@ -404,7 +453,15 @@ describe('createFlueInstrumentation', () => { await instr.interceptor(AGENT_OP, AGENT_CTX, async () => { instr.observe({ type: 'turn_start', turnId: 'turn_1', operationId: 'op_1' }, {}); instr.observe(requestContent, {}); - instr.observe(turn({ response: { ...turn().response, output: { role: 'assistant' } } }), {}); + instr.observe( + turn({ + response: { + ...turn().response, + output: { role: 'assistant', content: [{ type: 'text', text: 'It is sunny.' }] }, + }, + }), + {}, + ); instr.observe( { type: 'tool_start', @@ -415,8 +472,16 @@ describe('createFlueInstrumentation', () => { }, {}, ); + // `result` is the harness-level shape; `effectiveResult` is what the model receives. instr.observe( - { type: 'tool', toolCallId: 'c1', toolName: 'get_weather', result: 'sunny', operationId: 'op_1' }, + { + type: 'tool', + toolCallId: 'c1', + toolName: 'get_weather', + result: { content: [{ type: 'text', text: '"sunny"' }], details: { output: 'sunny' } }, + effectiveResult: 'sunny', + operationId: 'op_1', + }, {}, ); }); @@ -427,8 +492,11 @@ describe('createFlueInstrumentation', () => { const chat = findSpan('chat claude-haiku-4.5'); expect(chat?.data['gen_ai.system_instructions']).toBe('You are helpful.'); - expect(chat?.data['gen_ai.input.messages']).toContain('"role":"user"'); - expect(chat?.data['gen_ai.output.messages']).toContain('"role":"assistant"'); + // Mapped from pi-ai's message shape to the conventions, which is what Sentry renders. + expect(chat?.data['gen_ai.input.messages']).toBe('[{"role":"user","parts":[{"type":"text","content":"hi"}]}]'); + expect(chat?.data['gen_ai.output.messages']).toBe( + '[{"role":"assistant","parts":[{"type":"text","content":"It is sunny."}],"finish_reason":"stop"}]', + ); expect(chat?.data['gen_ai.tool.definitions']).toContain('get_weather'); const tool = findSpan('execute_tool get_weather'); @@ -436,6 +504,99 @@ describe('createFlueInstrumentation', () => { expect(tool?.data['gen_ai.tool.call.result']).toBe('sunny'); }); + it('records the content blocks of a tool result that has no effective result', async () => { + await withAgent(() => { + instrumentation.observe( + { type: 'tool_start', toolCallId: 'c1', toolName: 'get_weather', operationId: 'op_1' }, + {}, + ); + instrumentation.observe( + { + type: 'tool', + toolCallId: 'c1', + toolName: 'get_weather', + result: { + content: [ + { type: 'text', text: 'sunny' }, + { type: 'image', data: 'iVBORw0KGgo=', mimeType: 'image/png' }, + ], + details: { customTool: 'get_weather' }, + }, + operationId: 'op_1', + }, + {}, + ); + }); + + expect(findSpan('execute_tool get_weather')?.data['gen_ai.tool.call.result']).toBe( + '[{"type":"text","content":"sunny"},{"type":"blob","mime_type":"image/png"}]', + ); + }); + + it('maps an effective result made of several content blocks', async () => { + await withAgent(() => { + instrumentation.observe( + { type: 'tool_start', toolCallId: 'c1', toolName: 'get_weather', operationId: 'op_1' }, + {}, + ); + instrumentation.observe( + { + type: 'tool', + toolCallId: 'c1', + toolName: 'get_weather', + effectiveResult: [ + { type: 'text', text: 'sunny' }, + { type: 'image', data: '[image data omitted from event]', mimeType: 'image/png' }, + ], + operationId: 'op_1', + }, + {}, + ); + }); + + expect(findSpan('execute_tool get_weather')?.data['gen_ai.tool.call.result']).toBe( + '[{"type":"text","content":"sunny"},{"type":"blob","mime_type":"image/png"}]', + ); + }); + + it('maps a tool result in the request and a tool call in the response', async () => { + await withAgent(() => { + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', operationId: 'op_1' }, {}); + instrumentation.observe( + { + ...requestContent, + request: { + ...requestContent.request, + input: { + messages: [ + { role: 'user', content: 'Weather in Berlin?' }, + { + role: 'assistant', + content: [{ type: 'toolCall', id: 'c1', name: 'get_weather', arguments: { city: 'Berlin' } }], + }, + { + role: 'toolResult', + toolCallId: 'c1', + toolName: 'get_weather', + content: [{ type: 'text', text: 'sunny' }], + isError: false, + }, + ], + }, + }, + }, + {}, + ); + instrumentation.observe(turn(), {}); + }); + + expect( + JSON.parse(String(findSpan('chat claude-haiku-4.5')?.data['gen_ai.input.messages'])).map( + (message: { role: string }) => message.role, + ), + ).toEqual(['user', 'assistant', 'tool']); + }); + it('omits inputs when recordInputs is false but keeps outputs', async () => { await record(createFlueInstrumentation({ recordInputs: false })); diff --git a/packages/server-utils/test/ai/lib/utils/pi-ai-messages.test.ts b/packages/server-utils/test/ai/lib/utils/pi-ai-messages.test.ts new file mode 100644 index 000000000000..d75389725cf5 --- /dev/null +++ b/packages/server-utils/test/ai/lib/utils/pi-ai-messages.test.ts @@ -0,0 +1,141 @@ +import { describe, expect, it } from 'vitest'; +import { + piAiAssistantMessageToGenAiMessage, + piAiContentToString, + piAiFinishReason, + piAiMessagesToGenAiMessages, +} from '../../../../src/ai/pi-ai/messages'; + +describe('convert pi-ai messages to gen_ai messages', () => { + it('maps user, assistant and tool result messages', () => { + expect( + piAiMessagesToGenAiMessages([ + { role: 'user', content: 'Weather in Vienna?', timestamp: 1 }, + { + role: 'assistant', + content: [ + { type: 'thinking', thinking: 'The user wants the weather.' }, + { type: 'text', text: 'Let me check.' }, + { type: 'toolCall', id: 'call_1', name: 'get_weather', arguments: { city: 'Vienna' } }, + ], + stopReason: 'toolUse', + }, + { + role: 'toolResult', + toolCallId: 'call_1', + toolName: 'get_weather', + content: [{ type: 'text', text: 'Sunny in Vienna.' }], + isError: false, + }, + ]), + ).toStrictEqual([ + { role: 'user', parts: [{ type: 'text', content: 'Weather in Vienna?' }] }, + { + role: 'assistant', + parts: [ + { type: 'reasoning', content: 'The user wants the weather.' }, + { type: 'text', content: 'Let me check.' }, + { type: 'tool_call', id: 'call_1', name: 'get_weather', arguments: '{"city":"Vienna"}' }, + ], + }, + { + role: 'tool', + parts: [{ type: 'tool_call_response', id: 'call_1', name: 'get_weather', result: 'Sunny in Vienna.' }], + }, + ]); + }); + + it('reports images by media type and drops their base64 data', () => { + expect( + piAiMessagesToGenAiMessages([ + { + role: 'user', + content: [ + { type: 'text', text: 'What is this?' }, + { type: 'image', data: 'iVBORw0KGgo=', mimeType: 'image/png' }, + ], + }, + ]), + ).toStrictEqual([ + { + role: 'user', + parts: [ + { type: 'text', content: 'What is this?' }, + { type: 'blob', mime_type: 'image/png' }, + ], + }, + ]); + }); + + it('leaves out system messages, redacted thinking and messages without content', () => { + expect( + piAiMessagesToGenAiMessages([ + { role: 'system', content: '', sections: { preamble: 'You are helpful.' } }, + { role: 'assistant', content: [{ type: 'thinking', thinking: 'c2VjcmV0', redacted: true }] }, + { role: 'user', content: '' }, + 'not a message', + ]), + ).toStrictEqual([]); + }); + + it('drops empty text blocks, which would hide the tool calls next to them', () => { + expect( + piAiMessagesToGenAiMessages([ + { + role: 'assistant', + content: [ + { type: 'text', text: '' }, + { type: 'text', text: ' \n' }, + { type: 'toolCall', id: 'call_1', name: 'get_weather', arguments: {} }, + ], + }, + { role: 'assistant', content: [{ type: 'text', text: '' }] }, + ]), + ).toStrictEqual([ + { role: 'assistant', parts: [{ type: 'tool_call', id: 'call_1', name: 'get_weather', arguments: '{}' }] }, + ]); + }); + + it('maps the tool-calling stop reason to the conventions and passes the others through', () => { + expect(piAiFinishReason('toolUse')).toBe('tool_call'); + expect(piAiFinishReason('stop')).toBe('stop'); + expect(piAiFinishReason('aborted')).toBe('aborted'); + expect(piAiFinishReason(undefined)).toBeUndefined(); + expect(piAiFinishReason('')).toBeUndefined(); + }); + + it('keeps unknown part kinds as objects', () => { + expect(piAiMessagesToGenAiMessages([{ role: 'user', content: [{ type: 'audio', url: 'a.mp3' }] }])).toStrictEqual([ + { role: 'user', parts: [{ type: 'object', content: { type: 'audio', url: 'a.mp3' } }] }, + ]); + }); + + it('maps an assistant message to an output message with its finish reason', () => { + expect( + piAiAssistantMessageToGenAiMessage( + { + role: 'assistant', + content: [{ type: 'toolCall', id: 'call_1', name: 'get_weather', arguments: { city: 'Vienna' } }], + }, + 'tool_call', + ), + ).toStrictEqual({ + role: 'assistant', + parts: [{ type: 'tool_call', id: 'call_1', name: 'get_weather', arguments: '{"city":"Vienna"}' }], + finish_reason: 'tool_call', + }); + expect(piAiAssistantMessageToGenAiMessage({ role: 'assistant', content: [] }, 'stop')).toBeUndefined(); + }); + + it('renders text-only content as text and other content as mapped parts', () => { + expect( + piAiContentToString([ + { type: 'text', text: 'line 1' }, + { type: 'text', text: 'line 2' }, + ]), + ).toBe('line 1\nline 2'); + expect(piAiContentToString([{ type: 'image', data: 'iVBORw0KGgo=', mimeType: 'image/png' }])).toBe( + '[{"type":"blob","mime_type":"image/png"}]', + ); + }); +}); diff --git a/packages/server-utils/test/ai/lib/utils/pi-ai-usage.test.ts b/packages/server-utils/test/ai/lib/utils/pi-ai-usage.test.ts new file mode 100644 index 000000000000..6d4b768ba616 --- /dev/null +++ b/packages/server-utils/test/ai/lib/utils/pi-ai-usage.test.ts @@ -0,0 +1,62 @@ +import { describe, expect, it } from 'vitest'; +import type { Span } from '@sentry/core'; +import { setPiAiUsageAttributes } from '../../../../src/ai/pi-ai/usage'; + +function spanRecorder(): { span: Span; attributes: Record } { + const attributes: Record = {}; + return { + span: { setAttributes: (values: Record) => Object.assign(attributes, values) } as unknown as Span, + attributes, + }; +} + +describe('setPiAiUsageAttributes', () => { + it('counts the cached tokens into the input tokens and their cost into the input cost', () => { + const { span, attributes } = spanRecorder(); + + setPiAiUsageAttributes(span, { + input: 37, + output: 20, + cacheRead: 52, + cacheWrite: 38, + reasoning: 8, + totalTokens: 147, + cost: { input: 0.000111, output: 0.0001, cacheRead: 0.0000156, cacheWrite: 0.0001425, total: 0.0003691 }, + }); + + expect(attributes).toStrictEqual({ + 'gen_ai.usage.input_tokens': 127, + 'gen_ai.usage.output_tokens': 20, + 'gen_ai.usage.total_tokens': 147, + 'gen_ai.usage.cache_read.input_tokens': 52, + 'gen_ai.usage.cache_creation.input_tokens': 38, + 'gen_ai.usage.reasoning.output_tokens': 8, + 'gen_ai.cost.input_tokens': 0.000111 + 0.0000156 + 0.0001425, + 'gen_ai.cost.output_tokens': 0.0001, + 'gen_ai.cost.total_tokens': 0.0003691, + 'gen_ai.cost.cache_read.input_tokens': 0.0000156, + 'gen_ai.cost.cache_creation.input_tokens': 0.0001425, + }); + }); + + it('writes only the counters the response has', () => { + const { span, attributes } = spanRecorder(); + + setPiAiUsageAttributes(span, { input: 10, output: 2, totalTokens: 12 }); + + expect(attributes).toStrictEqual({ + 'gen_ai.usage.input_tokens': 10, + 'gen_ai.usage.output_tokens': 2, + 'gen_ai.usage.total_tokens': 12, + }); + }); + + it('skips the zero counters of a request that failed before it was billed', () => { + const { span, attributes } = spanRecorder(); + + setPiAiUsageAttributes(span, { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, totalTokens: 0 }, true); + setPiAiUsageAttributes(span, undefined); + + expect(attributes).toStrictEqual({}); + }); +});