Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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,
});
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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]!;
Expand All @@ -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,
);
});
29 changes: 8 additions & 21 deletions packages/server-utils/src/ai/flue/index.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,5 @@
import type { Span } from '@sentry/core';
import {
_INTERNAL_shouldSkipAiProviderWrapping,
_INTERNAL_skipAiProviderWrapping,
continueTrace,
getActiveSpan,
LRUMap,
Expand All @@ -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 {
Expand All @@ -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`.
*
Expand All @@ -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`
Expand Down
21 changes: 5 additions & 16 deletions packages/server-utils/src/ai/flue/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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;
}

Expand Down Expand Up @@ -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;
}

Expand Down
97 changes: 44 additions & 53 deletions packages/server-utils/src/ai/flue/utils.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import type { LRUMap, Span } from '@sentry/core';
import {
captureException,
isObjectLike,
SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN,
SPAN_STATUS_ERROR,
startInactiveSpan,
Expand All @@ -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,
Expand All @@ -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.
Expand Down Expand Up @@ -163,62 +161,32 @@ 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' });
}
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<string, number> = {};
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
Expand Down Expand Up @@ -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) {
Expand All @@ -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.
Expand All @@ -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));
Expand Down
Loading
Loading