From 253fa9a4126a706c8d01d3f7251046ad6fe8e18e Mon Sep 17 00:00:00 2001 From: RulaKhaled Date: Fri, 11 Sep 2026 08:12:06 +0200 Subject: [PATCH 1/4] test(server-utils): Cover the Flue instrumentation Unit tests over `createFlueInstrumentation` for the span shapes, the conversation id lifted off the re-entered agent operation, the usage/cost mapping, the all-zero-usage guard on failed turns, tool spans, content recording and its `recordInputs`/`recordOutputs` gating, and dispose. The integration test drives a real agent through a tool call using `pi-ai`'s `faux` provider, so the run is deterministic and needs no provider key or mock server. ESM only: `@flue/runtime` has no `require` export condition, and it is installed per-suite because its `engines.node >= 22.19` would break `yarn install` on the Node 20 CI matrix. Co-Authored-By: Claude Opus 5 --- .../suites/tracing/flue/instrument.mjs | 11 + .../suites/tracing/flue/scenario.mjs | 38 +++ .../suites/tracing/flue/test.ts | 86 ++++++ .../test/ai/lib/tracing/flue.test.ts | 272 ++++++++++++++++++ 4 files changed, 407 insertions(+) create mode 100644 dev-packages/node-integration-tests/suites/tracing/flue/instrument.mjs create mode 100644 dev-packages/node-integration-tests/suites/tracing/flue/scenario.mjs create mode 100644 dev-packages/node-integration-tests/suites/tracing/flue/test.ts create mode 100644 packages/server-utils/test/ai/lib/tracing/flue.test.ts diff --git a/dev-packages/node-integration-tests/suites/tracing/flue/instrument.mjs b/dev-packages/node-integration-tests/suites/tracing/flue/instrument.mjs new file mode 100644 index 000000000000..42052d281304 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/flue/instrument.mjs @@ -0,0 +1,11 @@ +import * as Sentry from '@sentry/node'; +import { loggingTransport } from '@sentry-internal/node-integration-tests'; + +Sentry.init({ + traceLifecycle: 'static', + dsn: 'https://public@dsn.ingest.sentry.io/1337', + release: '1.0', + tracesSampleRate: 1.0, + dataCollection: { genAI: { inputs: false, outputs: false } }, + transport: loggingTransport, +}); diff --git a/dev-packages/node-integration-tests/suites/tracing/flue/scenario.mjs b/dev-packages/node-integration-tests/suites/tracing/flue/scenario.mjs new file mode 100644 index 000000000000..625abf964ee0 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/flue/scenario.mjs @@ -0,0 +1,38 @@ +import * as Sentry from '@sentry/node'; +import { __flueBindAgentModule, init, instrument, useModel, useTool } from '@flue/runtime'; +import { start } from '@flue/runtime/node'; +import { fauxAssistantMessage, fauxProvider, fauxToolCall } from '@earendil-works/pi-ai/providers/faux'; + +// `pi-ai`'s faux provider scripts model responses in-process, so the run is deterministic and needs +// no provider key or mock server. Two steps: a tool call, then the final answer. +instrument(Sentry.createFlueInstrumentation()); + +const faux = fauxProvider({ + provider: 'faux', + models: [{ id: 'faux-model', cost: { input: 1, output: 2, cacheRead: 0, cacheWrite: 0 } }], +}); +faux.setResponses([ + fauxAssistantMessage(fauxToolCall('get_weather', { city: 'Berlin' }, { id: 'call_1' }), { stopReason: 'toolUse' }), + fauxAssistantMessage('It is 21 degrees and sunny in Berlin.'), +]); + +function Hello() { + useModel('faux/faux-model'); + useTool({ + name: 'get_weather', + description: 'Get the current weather for a city.', + run: ({ city }) => `It is 21 degrees and sunny in ${city}.`, + }); + return 'You are a helpful assistant.'; +} +__flueBindAgentModule(Hello, { identity: 'Hello' }); + +await Sentry.startSpan({ name: 'flue-test', op: 'function' }, async () => { + const flue = await start({ agents: [Hello], providers: [faux.provider] }); + const agent = init(Hello, { id: 'e2e' }); + const receipt = await agent.dispatch('What is the weather in Berlin?'); + await agent.read(receipt); + await flue[Symbol.asyncDispose]?.(); +}); + +await Sentry.flush(2000); diff --git a/dev-packages/node-integration-tests/suites/tracing/flue/test.ts b/dev-packages/node-integration-tests/suites/tracing/flue/test.ts new file mode 100644 index 000000000000..ec7a26f71ad1 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/flue/test.ts @@ -0,0 +1,86 @@ +import { + GEN_AI_AGENT_NAME, + GEN_AI_CONVERSATION_ID, + GEN_AI_COST_TOTAL_TOKENS, + GEN_AI_OPERATION_NAME, + GEN_AI_TOOL_NAME, + GEN_AI_USAGE_INPUT_TOKENS, + GEN_AI_USAGE_OUTPUT_TOKENS, + GEN_AI_USAGE_TOTAL_TOKENS, +} from '@sentry/conventions/attributes'; +import { afterAll, expect } from 'vitest'; +import { conditionalTest } from '../../../utils'; +import { cleanupChildProcesses, createEsmAndCjsTests } from '../../../utils/runner'; + +// `@flue/runtime` declares `engines.node >= 22.19`, so it can't live in the package's root +// `devDependencies` (that would break `yarn install` on the 20.19 CI matrix). Install it per-suite +// instead, guarded by the `min: 22` skip below. +const FLUE_DEPENDENCIES = { + additionalDependencies: { + '@flue/runtime': '2.0.3', + '@earendil-works/pi-ai': '0.85.1', + }, +}; + +conditionalTest({ min: 22 })('Flue integration', () => { + afterAll(() => { + cleanupChildProcesses(); + }); + + createEsmAndCjsTests( + __dirname, + 'scenario.mjs', + 'instrument.mjs', + (createRunner, test, mode) => { + // `@flue/runtime` is ESM-only — its `exports` map has no `require` condition, so there is no + // CJS variant of this scenario to run. + if (mode === 'cjs') { + return; + } + + test('creates the invoke_agent / chat / execute_tool hierarchy', async () => { + await createRunner() + .expect({ transaction: { transaction: 'flue-test' } }) + .expect({ + span: container => { + const spans = container.items; + const names = spans.map(span => span.name); + + expect(names).toContain('invoke_agent Hello'); + expect(names).toContain('execute_tool get_weather'); + expect(names.filter(name => name?.startsWith('chat'))).toHaveLength(2); + + const agent = spans.find(span => span.name === 'invoke_agent Hello')!; + expect(agent.attributes['sentry.op'].value).toBe('gen_ai.invoke_agent'); + expect(agent.attributes['sentry.origin'].value).toBe('auto.ai.flue'); + expect(agent.attributes[GEN_AI_OPERATION_NAME].value).toBe('invoke_agent'); + expect(agent.attributes[GEN_AI_AGENT_NAME].value).toBe('Hello'); + expect(agent.attributes[GEN_AI_CONVERSATION_ID].value).toEqual(expect.any(String)); + + const chat = spans.find(span => span.name?.startsWith('chat'))!; + expect(chat.attributes['sentry.op'].value).toBe('gen_ai.chat'); + expect(chat.attributes['sentry.origin'].value).toBe('auto.ai.flue'); + expect(chat.attributes[GEN_AI_USAGE_INPUT_TOKENS].value).toEqual(expect.any(Number)); + expect(chat.attributes[GEN_AI_USAGE_OUTPUT_TOKENS].value).toEqual(expect.any(Number)); + expect(chat.attributes[GEN_AI_USAGE_TOTAL_TOKENS].value).toEqual(expect.any(Number)); + // Flue computes cost itself; no provider SDK reports it. + expect(chat.attributes[GEN_AI_COST_TOTAL_TOKENS].value).toEqual(expect.any(Number)); + + const tool = spans.find(span => span.name === 'execute_tool get_weather')!; + expect(tool.attributes['sentry.op'].value).toBe('gen_ai.execute_tool'); + expect(tool.attributes['sentry.origin'].value).toBe('auto.ai.flue'); + expect(tool.attributes[GEN_AI_TOOL_NAME].value).toBe('get_weather'); + + // Tool spans are siblings of `chat` under the agent invocation, matching how Flue's + // own OpenTelemetry adapter projects them. + expect(tool.parent_span_id).toBe(agent.span_id); + expect(chat.parent_span_id).toBe(agent.span_id); + }, + }) + .start() + .completed(); + }); + }, + FLUE_DEPENDENCIES, + ); +}); diff --git a/packages/server-utils/test/ai/lib/tracing/flue.test.ts b/packages/server-utils/test/ai/lib/tracing/flue.test.ts new file mode 100644 index 000000000000..97485285282e --- /dev/null +++ b/packages/server-utils/test/ai/lib/tracing/flue.test.ts @@ -0,0 +1,272 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import type { Span } from '@sentry/core'; +import { + _INTERNAL_clearAiProviderSkips, + _INTERNAL_shouldSkipAiProviderWrapping, + getMainCarrier, + setCurrentClient, + spanToStaticSpanJSON, +} from '@sentry/core'; +import { ANTHROPIC_AI_INTEGRATION_NAME } from '../../../../src/ai/anthropic-ai/constants'; +import { createFlueInstrumentation } from '../../../../src/ai/flue'; +import type { FlueInstrumentation, FlueObservation } from '../../../../src/ai/flue/types'; +import { OPENAI_INTEGRATION_NAME } from '../../../../src/ai/openai/constants'; +import { getDefaultTestClientOptions, TestClient } from '../../../mocks/client'; + +const AGENT_CTX = { agentName: 'Hello' }; +const INNER_CTX = { conversationId: 'conv_1' }; + +/** A settled turn as Flue reports it, with the field names `ModelRequestInfo`/`ModelResponse` use. */ +function turn(overrides: Partial = {}): FlueObservation { + return { + type: 'turn', + turnId: 'turn_1', + request: { requestedModel: 'claude-haiku-4.5', providerId: 'anthropic' }, + response: { + responseId: 'resp_1', + finishReason: 'stop', + usage: { + input: 924, + output: 57, + totalTokens: 981, + cacheRead: 0, + cacheWrite: 0, + cost: { input: 0.000924, output: 0.000275, total: 0.001199, cacheRead: 0, cacheWrite: 0 }, + }, + }, + ...overrides, + }; +} + +describe('createFlueInstrumentation', () => { + let endedSpans: Span[]; + let instrumentation: FlueInstrumentation; + + beforeEach(() => { + _INTERNAL_clearAiProviderSkips(); + getMainCarrier().__SENTRY__ = undefined; + const client = new TestClient( + getDefaultTestClientOptions({ + dsn: 'https://public@dsn.ingest.sentry.io/1337', + tracesSampleRate: 1, + traceLifecycle: 'stream', + }), + ); + setCurrentClient(client); + client.init(); + + endedSpans = []; + client.on('spanEnd', span => endedSpans.push(span)); + instrumentation = createFlueInstrumentation(); + }); + + afterEach(() => { + _INTERNAL_clearAiProviderSkips(); + getMainCarrier().__SENTRY__ = undefined; + }); + + /** Run `fn` inside an agent operation, the way Flue's interceptor would. */ + function withAgent(fn: () => Promise | T): Promise { + return instrumentation.interceptor({ type: 'agent' }, AGENT_CTX, async () => fn()); + } + + function findSpan(description: string): ReturnType | undefined { + return endedSpans.map(span => spanToStaticSpanJSON(span)).find(json => json.description === description); + } + + // Flue drives the providers through `pi-ai`, which bundles the `openai` / `@anthropic-ai/sdk` / + // `@google/genai` clients those integrations patch, so their spans duplicate the turn span. + it('skips raw provider wrapping as soon as the instrumentation is built', () => { + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(true); + expect(_INTERNAL_shouldSkipAiProviderWrapping(ANTHROPIC_AI_INTEGRATION_NAME)).toBe(true); + }); + + it('names agent spans `invoke_agent {name}` and sets the gen_ai op', async () => { + await withAgent(() => undefined); + + const json = findSpan('invoke_agent Hello'); + expect(json?.data['sentry.op']).toBe('gen_ai.invoke_agent'); + expect(json?.data['sentry.origin']).toBe('auto.ai.flue'); + expect(json?.data['gen_ai.operation.name']).toBe('invoke_agent'); + expect(json?.data['gen_ai.agent.name']).toBe('Hello'); + }); + + // The agent operation re-enters once, and only the inner context names the conversation. + it('lifts the conversation id off the re-entered agent operation', async () => { + await withAgent(() => instrumentation.interceptor({ type: 'agent' }, INNER_CTX, async () => undefined)); + + expect(findSpan('invoke_agent Hello')?.data['gen_ai.conversation.id']).toBe('conv_1'); + }); + + it('does not span operations other than `agent`', async () => { + await instrumentation.interceptor({ type: 'model', turnId: 'turn_1' }, {}, async () => undefined); + + expect(endedSpans).toHaveLength(0); + }); + + it('opens a chat span on turn_start and completes it from the settled turn', async () => { + await withAgent(() => { + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', conversationId: 'conv_1' }, {}); + instrumentation.observe(turn(), {}); + }); + + const json = findSpan('chat claude-haiku-4.5'); + expect(json?.data['sentry.op']).toBe('gen_ai.chat'); + expect(json?.data['sentry.origin']).toBe('auto.ai.flue'); + expect(json?.data['gen_ai.request.model']).toBe('claude-haiku-4.5'); + expect(json?.data['gen_ai.provider.name']).toBe('anthropic'); + expect(json?.data['gen_ai.response.id']).toBe('resp_1'); + expect(json?.data['gen_ai.response.finish_reasons']).toEqual(['stop']); + expect(json?.data['gen_ai.conversation.id']).toBe('conv_1'); + }); + + // Flue computes costs itself; the provider SDKs report none. + it('records token usage and Flue-computed cost on the chat span', async () => { + await withAgent(() => { + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1' }, {}); + instrumentation.observe(turn(), {}); + }); + + const json = findSpan('chat claude-haiku-4.5'); + expect(json?.data['gen_ai.usage.input_tokens']).toBe(924); + expect(json?.data['gen_ai.usage.output_tokens']).toBe(57); + expect(json?.data['gen_ai.usage.total_tokens']).toBe(981); + expect(json?.data['gen_ai.cost.total_tokens']).toBe(0.001199); + }); + + it('marks a failed turn as errored', async () => { + await withAgent(() => { + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1' }, {}); + instrumentation.observe(turn({ isError: true }), {}); + }); + + expect(findSpan('chat claude-haiku-4.5')?.status).toBe('internal_error'); + }); + + // A turn that fails before the provider bills anything reports every counter as 0; writing those + // reads as a real zero-cost call. + it('omits usage entirely when a failed turn produced no tokens', async () => { + const empty = { + input: 0, + output: 0, + totalTokens: 0, + cacheRead: 0, + cacheWrite: 0, + cost: { input: 0, output: 0, total: 0, cacheRead: 0, cacheWrite: 0 }, + }; + + await withAgent(() => { + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1' }, {}); + instrumentation.observe(turn({ isError: true, response: { usage: empty } }), {}); + }); + + const json = findSpan('chat claude-haiku-4.5'); + expect(json).toBeDefined(); + expect(json?.data['gen_ai.usage.total_tokens']).toBeUndefined(); + expect(json?.data['gen_ai.cost.total_tokens']).toBeUndefined(); + }); + + describe('content recording', () => { + const requestContent = { + type: 'turn_request', + turnId: 'turn_1', + request: { + requestedModel: 'claude-haiku-4.5', + input: { + systemPrompt: 'You are helpful.', + messages: [{ role: 'user', content: 'hi' }], + tools: [{ name: 'get_weather', description: 'weather', parameters: {} }], + }, + }, + } satisfies FlueObservation; + + async function record(instr: FlueInstrumentation): Promise { + await instr.interceptor({ type: 'agent' }, AGENT_CTX, async () => { + instr.observe({ type: 'turn_start', turnId: 'turn_1' }, {}); + instr.observe(requestContent, {}); + instr.observe(turn({ response: { ...turn().response, output: { role: 'assistant' } } }), {}); + instr.observe({ type: 'tool_start', toolCallId: 'c1', toolName: 'get_weather', args: { city: 'Berlin' } }, {}); + instr.observe({ type: 'tool', toolCallId: 'c1', toolName: 'get_weather', result: 'sunny' }, {}); + }); + } + + it('records messages, instructions, tool definitions, arguments and results by default', async () => { + await record(instrumentation); + + 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"'); + expect(chat?.data['gen_ai.tool.definitions']).toContain('get_weather'); + + const tool = findSpan('execute_tool get_weather'); + expect(tool?.data['gen_ai.tool.call.arguments']).toBe('{"city":"Berlin"}'); + expect(tool?.data['gen_ai.tool.call.result']).toBe('sunny'); + }); + + it('omits inputs when recordInputs is false but keeps outputs', async () => { + await record(createFlueInstrumentation({ recordInputs: false })); + + const chat = findSpan('chat claude-haiku-4.5'); + expect(chat?.data['gen_ai.input.messages']).toBeUndefined(); + expect(chat?.data['gen_ai.system_instructions']).toBeUndefined(); + expect(chat?.data['gen_ai.tool.definitions']).toBeUndefined(); + expect(chat?.data['gen_ai.output.messages']).toBeDefined(); + expect(findSpan('execute_tool get_weather')?.data['gen_ai.tool.call.arguments']).toBeUndefined(); + }); + + it('omits outputs when recordOutputs is false but keeps inputs', async () => { + await record(createFlueInstrumentation({ recordOutputs: false })); + + const chat = findSpan('chat claude-haiku-4.5'); + expect(chat?.data['gen_ai.output.messages']).toBeUndefined(); + expect(chat?.data['gen_ai.input.messages']).toBeDefined(); + expect(findSpan('execute_tool get_weather')?.data['gen_ai.tool.call.result']).toBeUndefined(); + }); + }); + + it('emits execute_tool spans keyed by tool call id', async () => { + await withAgent(() => { + instrumentation.observe({ type: 'tool_start', toolCallId: 'call_1', toolName: 'get_weather' }, {}); + instrumentation.observe({ type: 'tool', toolCallId: 'call_1', toolName: 'get_weather' }, {}); + }); + + const json = findSpan('execute_tool get_weather'); + expect(json?.data['sentry.op']).toBe('gen_ai.execute_tool'); + expect(json?.data['sentry.origin']).toBe('auto.ai.flue'); + expect(json?.data['gen_ai.operation.name']).toBe('execute_tool'); + expect(json?.data['gen_ai.tool.name']).toBe('get_weather'); + }); + + it('marks a failed tool call as errored', async () => { + await withAgent(() => { + instrumentation.observe({ type: 'tool_start', toolCallId: 'call_1', toolName: 'boom' }, {}); + instrumentation.observe({ type: 'tool', toolCallId: 'call_1', toolName: 'boom', isError: true }, {}); + }); + + expect(findSpan('execute_tool boom')?.status).toBe('internal_error'); + }); + + it('ignores a settled turn or tool it never opened a span for', async () => { + await withAgent(() => { + instrumentation.observe(turn({ turnId: 'never_started' }), {}); + instrumentation.observe({ type: 'tool', toolCallId: 'never_started', toolName: 'x' }, {}); + }); + + expect(endedSpans.map(span => spanToStaticSpanJSON(span).description)).toEqual(['invoke_agent Hello']); + }); + + it('ends spans still open at dispose', async () => { + await withAgent(() => { + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1' }, {}); + instrumentation.observe({ type: 'tool_start', toolCallId: 'call_1', toolName: 'get_weather' }, {}); + }); + // Never settled, so the span keeps the unqualified name it opened with. + expect(findSpan('chat')).toBeUndefined(); + + instrumentation.dispose(); + + expect(findSpan('chat')).toBeDefined(); + expect(findSpan('execute_tool get_weather')).toBeDefined(); + }); +}); From 4789aeb8d95781f7afe8d0d9c28932eed0fc5192 Mon Sep 17 00:00:00 2001 From: RulaKhaled Date: Fri, 11 Sep 2026 10:19:27 +0200 Subject: [PATCH 2/4] test(server-utils): Cover the Flue review fixes Concurrent agent runs, subagent delegation (adopted from Isaac's repro, moved into the suite and asserted through a helper so a span-order assumption cannot creep back), the agent name arriving via the observations, trace continuation from the replayed traceparent, the provider skip applying on first use and re-applying after a registry reset, the recording options following the current client, and the conventional request attributes. Co-Authored-By: Claude Opus 5 --- .../test/ai/lib/tracing/flue.test.ts | 276 ++++++++++++++++-- 1 file changed, 255 insertions(+), 21 deletions(-) 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 97485285282e..066d94b16f76 100644 --- a/packages/server-utils/test/ai/lib/tracing/flue.test.ts +++ b/packages/server-utils/test/ai/lib/tracing/flue.test.ts @@ -6,6 +6,7 @@ import { getMainCarrier, setCurrentClient, spanToStaticSpanJSON, + startSpan, } from '@sentry/core'; import { ANTHROPIC_AI_INTEGRATION_NAME } from '../../../../src/ai/anthropic-ai/constants'; import { createFlueInstrumentation } from '../../../../src/ai/flue'; @@ -13,7 +14,8 @@ import type { FlueInstrumentation, FlueObservation } from '../../../../src/ai/fl import { OPENAI_INTEGRATION_NAME } from '../../../../src/ai/openai/constants'; import { getDefaultTestClientOptions, TestClient } from '../../../mocks/client'; -const AGENT_CTX = { agentName: 'Hello' }; +const AGENT_OP = { type: 'agent', operationId: 'op_1' }; +const AGENT_CTX = { agentName: 'Hello', submissionId: 'sub_1' }; const INNER_CTX = { conversationId: 'conv_1' }; /** A settled turn as Flue reports it, with the field names `ModelRequestInfo`/`ModelResponse` use. */ @@ -21,7 +23,17 @@ function turn(overrides: Partial = {}): FlueObservation { return { type: 'turn', turnId: 'turn_1', - request: { requestedModel: 'claude-haiku-4.5', providerId: 'anthropic' }, + submissionId: 'sub_1', + operationId: 'op_1', + request: { + requestedModel: 'claude-haiku-4.5', + providerId: 'anthropic', + temperature: 0.7, + maxTokens: 1024, + reasoningLevel: 'high', + serverAddress: 'api.anthropic.com', + serverPort: 443, + }, response: { responseId: 'resp_1', finishReason: 'stop', @@ -67,7 +79,13 @@ describe('createFlueInstrumentation', () => { /** Run `fn` inside an agent operation, the way Flue's interceptor would. */ function withAgent(fn: () => Promise | T): Promise { - return instrumentation.interceptor({ type: 'agent' }, AGENT_CTX, async () => fn()); + return instrumentation.interceptor(AGENT_OP, AGENT_CTX, async () => fn()); + } + + function agentSpans(): ReturnType[] { + return endedSpans + .map(span => spanToStaticSpanJSON(span)) + .filter(json => json.data['sentry.op'] === 'gen_ai.invoke_agent'); } function findSpan(description: string): ReturnType | undefined { @@ -76,11 +94,30 @@ describe('createFlueInstrumentation', () => { // Flue drives the providers through `pi-ai`, which bundles the `openai` / `@anthropic-ai/sdk` / // `@google/genai` clients those integrations patch, so their spans duplicate the turn span. - it('skips raw provider wrapping as soon as the instrumentation is built', () => { + // Not at construction: if `instrument()` rejects the object, suppressing the provider + // integrations would leave the app with no `gen_ai.chat` spans at all. + it('skips raw provider wrapping on first use, not on construction', async () => { + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(false); + + await withAgent(() => undefined); + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(true); expect(_INTERNAL_shouldSkipAiProviderWrapping(ANTHROPIC_AI_INTEGRATION_NAME)).toBe(true); }); + // The registry is reset per client (`_setupIntegrations` clears it, and Cloudflare calls `init()` + // per request), so a one-shot call at construction is wiped by the next reset. + it('re-applies the provider skip after the registry is cleared', async () => { + await withAgent(() => undefined); + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(true); + + _INTERNAL_clearAiProviderSkips(); + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(false); + + await withAgent(() => undefined); + expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(true); + }); + it('names agent spans `invoke_agent {name}` and sets the gen_ai op', async () => { await withAgent(() => undefined); @@ -91,13 +128,113 @@ describe('createFlueInstrumentation', () => { expect(json?.data['gen_ai.agent.name']).toBe('Hello'); }); - // The agent operation re-enters once, and only the inner context names the conversation. - it('lifts the conversation id off the re-entered agent operation', async () => { - await withAgent(() => instrumentation.interceptor({ type: 'agent' }, INNER_CTX, async () => undefined)); + // The agent span opens before the conversation is known — the submission-scoped operation names + // the agent, and the conversation arrives on the observations that follow. + it('sets the conversation id on the agent span from the observations', async () => { + await withAgent(() => { + instrumentation.observe( + { type: 'turn_start', turnId: 'turn_1', operationId: 'op_1', conversationId: 'conv_1' }, + {}, + ); + }); expect(findSpan('invoke_agent Hello')?.data['gen_ai.conversation.id']).toBe('conv_1'); }); + // The re-entered agent operation carries no `submissionId` and must not open a second span. + it('does not open a second agent span for the re-entry', async () => { + await withAgent(() => instrumentation.interceptor(AGENT_OP, INNER_CTX, async () => undefined)); + + const agentSpans = endedSpans + .map(span => spanToStaticSpanJSON(span)) + .filter(json => json.data['sentry.op'] === 'gen_ai.invoke_agent'); + expect(agentSpans).toHaveLength(1); + }); + + // A durable submission is resumed later, with nothing linking it to the request that enqueued it. + // Flue replays that request's `traceparent`, so the agent span should continue from it. + describe('trace continuation', () => { + const TRACE_ID = 'aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa'; + const CARRIER = { traceparent: `00-${TRACE_ID}-bbbbbbbbbbbbbbbb-01` }; + + function agentTraceId(): string | undefined { + return endedSpans.map(span => spanToStaticSpanJSON(span)).find(json => json.description === 'invoke_agent Hello') + ?.trace_id; + } + + it('continues the trace from the replayed traceparent', async () => { + await instrumentation.interceptor(AGENT_OP, { ...AGENT_CTX, traceCarrier: CARRIER }, async () => undefined); + + expect(agentTraceId()).toBe(TRACE_ID); + }); + + it('ignores a malformed traceparent', async () => { + await instrumentation.interceptor( + AGENT_OP, + { ...AGENT_CTX, traceCarrier: { traceparent: 'not-a-traceparent' } }, + async () => undefined, + ); + + expect(agentTraceId()).not.toBe(TRACE_ID); + }); + + // An in-process dispatch is genuinely part of the surrounding trace; continuing the persisted + // one would detach it from the request it is actually running inside. + it('keeps the active trace when one is already running', async () => { + await startSpan({ name: 'incoming request' }, async () => { + await instrumentation.interceptor(AGENT_OP, { ...AGENT_CTX, traceCarrier: CARRIER }, async () => undefined); + }); + + expect(agentTraceId()).not.toBe(TRACE_ID); + }); + }); + + // The spanned operation carries no `agentName` — only the submission wrapper does, and that opens + // no span. The name arrives on the observations instead. + it('names the agent span from the observations', async () => { + await instrumentation.interceptor(AGENT_OP, {}, async () => { + instrumentation.observe({ type: 'agent_start', operationId: 'op_1', agentName: 'Hello' }, {}); + }); + + const json = findSpan('invoke_agent Hello'); + expect(json).toBeDefined(); + expect(json?.data['gen_ai.agent.name']).toBe('Hello'); + }); + + // Delegation nests a second agent operation inside the first, in its own session. Flue defers + // each to a microtask, so the helper mirrors that rather than calling the interceptor directly. + describe('subagent delegation', () => { + function runNested(operationId: string, ctx: Record, next: () => Promise): Promise { + return Promise.resolve().then(() => instrumentation.interceptor({ type: 'agent', operationId }, ctx, next)); + } + + it('opens one agent span per invocation rather than folding the delegate into its parent', async () => { + await runNested('op_parent', {}, async () => { + instrumentation.observe({ type: 'agent_start', operationId: 'op_parent', conversationId: 'conv_parent' }, {}); + // The tool's task delegation, then the delegate's own prompt. + return runNested('op_child', {}, async () => { + instrumentation.observe({ type: 'agent_start', operationId: 'op_child', conversationId: 'conv_child' }, {}); + }); + }); + + expect(agentSpans()).toHaveLength(2); + }); + + // Asserted on conversation rather than agent name: the observation stream reports the root + // agent's name for both operations, so the delegate's own name never reaches us. + it('keeps each invocation on its own conversation', async () => { + await runNested('op_parent', {}, async () => { + instrumentation.observe({ type: 'agent_start', operationId: 'op_parent', conversationId: 'conv_parent' }, {}); + return runNested('op_child', {}, async () => { + instrumentation.observe({ type: 'agent_start', operationId: 'op_child', conversationId: 'conv_child' }, {}); + }); + }); + + // The delegate ends first, so order is inner-to-outer. + expect(agentSpans().map(json => json.data['gen_ai.conversation.id'])).toEqual(['conv_child', 'conv_parent']); + }); + }); + it('does not span operations other than `agent`', async () => { await instrumentation.interceptor({ type: 'model', turnId: 'turn_1' }, {}, async () => undefined); @@ -106,7 +243,10 @@ describe('createFlueInstrumentation', () => { it('opens a chat span on turn_start and completes it from the settled turn', async () => { await withAgent(() => { - instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', conversationId: 'conv_1' }, {}); + instrumentation.observe( + { type: 'turn_start', turnId: 'turn_1', operationId: 'op_1', conversationId: 'conv_1' }, + {}, + ); instrumentation.observe(turn(), {}); }); @@ -123,7 +263,7 @@ describe('createFlueInstrumentation', () => { // Flue computes costs itself; the provider SDKs report none. it('records token usage and Flue-computed cost on the chat span', async () => { await withAgent(() => { - instrumentation.observe({ type: 'turn_start', turnId: 'turn_1' }, {}); + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', operationId: 'op_1' }, {}); instrumentation.observe(turn(), {}); }); @@ -134,9 +274,25 @@ describe('createFlueInstrumentation', () => { expect(json?.data['gen_ai.cost.total_tokens']).toBe(0.001199); }); + 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' }, {}); + instrumentation.observe(turn(), {}); + }); + + const json = findSpan('chat claude-haiku-4.5'); + expect(json?.data['gen_ai.request.temperature']).toBe(0.7); + expect(json?.data['gen_ai.request.max_tokens']).toBe(1024); + expect(json?.data['gen_ai.request.reasoning.level']).toBe('high'); + expect(json?.data['server.address']).toBe('api.anthropic.com'); + expect(json?.data['server.port']).toBe(443); + // Distinguishes a compaction turn from a user-facing one. + expect(json?.data['flue.turn.purpose']).toBe('agent'); + }); + it('marks a failed turn as errored', async () => { await withAgent(() => { - instrumentation.observe({ type: 'turn_start', turnId: 'turn_1' }, {}); + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', operationId: 'op_1' }, {}); instrumentation.observe(turn({ isError: true }), {}); }); @@ -156,7 +312,7 @@ describe('createFlueInstrumentation', () => { }; await withAgent(() => { - instrumentation.observe({ type: 'turn_start', turnId: 'turn_1' }, {}); + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', operationId: 'op_1' }, {}); instrumentation.observe(turn({ isError: true, response: { usage: empty } }), {}); }); @@ -181,12 +337,24 @@ describe('createFlueInstrumentation', () => { } satisfies FlueObservation; async function record(instr: FlueInstrumentation): Promise { - await instr.interceptor({ type: 'agent' }, AGENT_CTX, async () => { - instr.observe({ type: 'turn_start', turnId: 'turn_1' }, {}); + 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({ type: 'tool_start', toolCallId: 'c1', toolName: 'get_weather', args: { city: 'Berlin' } }, {}); - instr.observe({ type: 'tool', toolCallId: 'c1', toolName: 'get_weather', result: 'sunny' }, {}); + instr.observe( + { + type: 'tool_start', + toolCallId: 'c1', + toolName: 'get_weather', + args: { city: 'Berlin' }, + operationId: 'op_1', + }, + {}, + ); + instr.observe( + { type: 'tool', toolCallId: 'c1', toolName: 'get_weather', result: 'sunny', operationId: 'op_1' }, + {}, + ); }); } @@ -215,6 +383,33 @@ describe('createFlueInstrumentation', () => { expect(findSpan('execute_tool get_weather')?.data['gen_ai.tool.call.arguments']).toBeUndefined(); }); + // The client is replaced per request on Cloudflare, so options captured once at construction + // would be the wrong ones for every later request. + it('follows the current client when it is replaced', async () => { + const instr = createFlueInstrumentation(); + await record(instr); + expect(findSpan('chat claude-haiku-4.5')?.data['gen_ai.input.messages']).toBeDefined(); + + endedSpans.length = 0; + const strict = new TestClient( + getDefaultTestClientOptions({ + dsn: 'https://public@dsn.ingest.sentry.io/1337', + tracesSampleRate: 1, + traceLifecycle: 'stream', + dataCollection: { genAI: { inputs: false, outputs: false } }, + }), + ); + setCurrentClient(strict); + strict.init(); + strict.on('spanEnd', span => endedSpans.push(span)); + + await record(instr); + + const chat = findSpan('chat claude-haiku-4.5'); + expect(chat).toBeDefined(); + expect(chat?.data['gen_ai.input.messages']).toBeUndefined(); + }); + it('omits outputs when recordOutputs is false but keeps inputs', async () => { await record(createFlueInstrumentation({ recordOutputs: false })); @@ -227,8 +422,11 @@ describe('createFlueInstrumentation', () => { it('emits execute_tool spans keyed by tool call id', async () => { await withAgent(() => { - instrumentation.observe({ type: 'tool_start', toolCallId: 'call_1', toolName: 'get_weather' }, {}); - instrumentation.observe({ type: 'tool', toolCallId: 'call_1', toolName: 'get_weather' }, {}); + instrumentation.observe( + { type: 'tool_start', toolCallId: 'call_1', toolName: 'get_weather', operationId: 'op_1' }, + {}, + ); + instrumentation.observe({ type: 'tool', toolCallId: 'call_1', toolName: 'get_weather', operationId: 'op_1' }, {}); }); const json = findSpan('execute_tool get_weather'); @@ -240,13 +438,46 @@ describe('createFlueInstrumentation', () => { it('marks a failed tool call as errored', async () => { await withAgent(() => { - instrumentation.observe({ type: 'tool_start', toolCallId: 'call_1', toolName: 'boom' }, {}); - instrumentation.observe({ type: 'tool', toolCallId: 'call_1', toolName: 'boom', isError: true }, {}); + instrumentation.observe({ type: 'tool_start', toolCallId: 'call_1', toolName: 'boom', operationId: 'op_1' }, {}); + instrumentation.observe( + { type: 'tool', toolCallId: 'call_1', toolName: 'boom', isError: true, operationId: 'op_1' }, + {}, + ); }); expect(findSpan('execute_tool boom')?.status).toBe('internal_error'); }); + // Two agent runs overlap on a busy server. With shared closure state the second run is mistaken + // for a re-entry of the first: it gets no span, and its conversation id lands on the first's span. + it('keeps concurrent agent runs separate', async () => { + const runA = instrumentation.interceptor({ type: 'agent', operationId: 'op_a' }, { agentName: 'A' }, async () => { + instrumentation.observe( + { type: 'turn_start', turnId: 'turn_a', operationId: 'op_a', conversationId: 'conv_a' }, + {}, + ); + // B starts while A is still open. + await instrumentation.interceptor({ type: 'agent', operationId: 'op_b' }, { agentName: 'B' }, async () => { + instrumentation.observe( + { type: 'turn_start', turnId: 'turn_b', operationId: 'op_b', conversationId: 'conv_b' }, + {}, + ); + instrumentation.observe(turn({ turnId: 'turn_b', operationId: 'op_b' }), {}); + }); + instrumentation.observe(turn({ turnId: 'turn_a', operationId: 'op_a' }), {}); + }); + await runA; + + const agentA = findSpan('invoke_agent A'); + const agentB = findSpan('invoke_agent B'); + expect(agentA).toBeDefined(); + expect(agentB).toBeDefined(); + + // Each run keeps its own conversation; neither is overwritten by the other. + expect(agentA?.data['gen_ai.conversation.id']).toBe('conv_a'); + expect(agentB?.data['gen_ai.conversation.id']).toBe('conv_b'); + }); + it('ignores a settled turn or tool it never opened a span for', async () => { await withAgent(() => { instrumentation.observe(turn({ turnId: 'never_started' }), {}); @@ -258,8 +489,11 @@ describe('createFlueInstrumentation', () => { it('ends spans still open at dispose', async () => { await withAgent(() => { - instrumentation.observe({ type: 'turn_start', turnId: 'turn_1' }, {}); - instrumentation.observe({ type: 'tool_start', toolCallId: 'call_1', toolName: 'get_weather' }, {}); + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', operationId: 'op_1' }, {}); + instrumentation.observe( + { type: 'tool_start', toolCallId: 'call_1', toolName: 'get_weather', operationId: 'op_1' }, + {}, + ); }); // Never settled, so the span keeps the unqualified name it opened with. expect(findSpan('chat')).toBeUndefined(); From cd565eefa98cc641bdc50cef52ca4475d596f595 Mon Sep 17 00:00:00 2001 From: RulaKhaled Date: Fri, 11 Sep 2026 13:45:12 +0200 Subject: [PATCH 3/4] test(server-utils): Cover the Flue span cap and provider skip guard Both tests fail against the previous implementation: the first leaves `openai` registered before the run so the old first-entry-only guard short-circuits, and the second overflows the turn tracker to prove the evicted span is ended rather than dropped unsent. Co-Authored-By: Claude Opus 5 --- .../test/ai/lib/tracing/flue.test.ts | 26 +++++++++++++++++++ 1 file changed, 26 insertions(+) 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 066d94b16f76..5d29e7424759 100644 --- a/packages/server-utils/test/ai/lib/tracing/flue.test.ts +++ b/packages/server-utils/test/ai/lib/tracing/flue.test.ts @@ -3,6 +3,7 @@ import type { Span } from '@sentry/core'; import { _INTERNAL_clearAiProviderSkips, _INTERNAL_shouldSkipAiProviderWrapping, + _INTERNAL_skipAiProviderWrapping, getMainCarrier, setCurrentClient, spanToStaticSpanJSON, @@ -10,6 +11,8 @@ import { } from '@sentry/core'; import { ANTHROPIC_AI_INTEGRATION_NAME } from '../../../../src/ai/anthropic-ai/constants'; import { createFlueInstrumentation } from '../../../../src/ai/flue'; +import { MAX_TRACKED_FLUE_SPANS } from '../../../../src/ai/flue/constants'; +import { GOOGLE_GENAI_INTEGRATION_NAME } from '../../../../src/ai/google-genai/constants'; import type { FlueInstrumentation, FlueObservation } from '../../../../src/ai/flue/types'; import { OPENAI_INTEGRATION_NAME } from '../../../../src/ai/openai/constants'; import { getDefaultTestClientOptions, TestClient } from '../../../mocks/client'; @@ -118,6 +121,29 @@ describe('createFlueInstrumentation', () => { expect(_INTERNAL_shouldSkipAiProviderWrapping(OPENAI_INTEGRATION_NAME)).toBe(true); }); + // The guard has to hold for every provider, not just the first: another integration may have + // registered a skip for one of them already, which would otherwise short-circuit the rest. + it('applies the skip to every provider when only some are already registered', async () => { + _INTERNAL_skipAiProviderWrapping([OPENAI_INTEGRATION_NAME]); + + await withAgent(() => undefined); + + expect(_INTERNAL_shouldSkipAiProviderWrapping(ANTHROPIC_AI_INTEGRATION_NAME)).toBe(true); + expect(_INTERNAL_shouldSkipAiProviderWrapping(GOOGLE_GENAI_INTEGRATION_NAME)).toBe(true); + }); + + // A turn whose stream is abandoned never emits the settled `turn` that would remove it, so the + // tracker is capped. Eviction has to end the span it drops, or it is never sent. + it('ends the oldest chat span when the turn tracker overflows', async () => { + await withAgent(() => { + for (let i = 0; i <= MAX_TRACKED_FLUE_SPANS; i++) { + instrumentation.observe({ type: 'turn_start', turnId: `turn_${i}`, operationId: 'op_1' }, {}); + } + }); + + expect(endedSpans.filter(span => spanToStaticSpanJSON(span).data['sentry.op'] === 'gen_ai.chat')).toHaveLength(1); + }); + it('names agent spans `invoke_agent {name}` and sets the gen_ai op', async () => { await withAgent(() => undefined); From d69cb28e4478d733d394d1089e649684c0d003fc Mon Sep 17 00:00:00 2001 From: RulaKhaled Date: Fri, 11 Sep 2026 15:09:46 +0200 Subject: [PATCH 4/4] test(server-utils): Cover the active turn and tool spans The `model` and `tool` interceptor branches open no span; they make the span `observe` already opened active so the provider's HTTP call and the tool's own work nest inside it. Deleting both branches left all 28 tests green, and the e2e does not reach it either: its parent assertions come from the observation stream firing inside the agent operation, and nothing in that scenario opens a span inside a tool or model operation. Each case fails when its own branch is removed, and neither fails for the other's. Co-Authored-By: Claude Opus 5 --- .../test/ai/lib/tracing/flue.test.ts | 31 +++++++++++++++++++ 1 file changed, 31 insertions(+) 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 5d29e7424759..ce3d4e8e46ce 100644 --- a/packages/server-utils/test/ai/lib/tracing/flue.test.ts +++ b/packages/server-utils/test/ai/lib/tracing/flue.test.ts @@ -7,6 +7,7 @@ import { getMainCarrier, setCurrentClient, spanToStaticSpanJSON, + startInactiveSpan, startSpan, } from '@sentry/core'; import { ANTHROPIC_AI_INTEGRATION_NAME } from '../../../../src/ai/anthropic-ai/constants'; @@ -267,6 +268,36 @@ describe('createFlueInstrumentation', () => { expect(endedSpans).toHaveLength(0); }); + // The `model` and `tool` operations open no span of their own; they make the span `observe` + // already opened active, so the provider's HTTP call and the tool's own work nest inside it + // rather than landing beside it as siblings of the agent invocation. + it('makes the turn span active for the model operation it wraps', async () => { + await withAgent(async () => { + instrumentation.observe({ type: 'turn_start', turnId: 'turn_1', operationId: 'op_1' }, {}); + await instrumentation.interceptor({ type: 'model', turnId: 'turn_1' }, {}, async () => { + startInactiveSpan({ name: 'provider request' }).end(); + }); + instrumentation.observe(turn(), {}); + }); + + expect(findSpan('provider request')?.parent_span_id).toBe(findSpan('chat claude-haiku-4.5')?.span_id); + }); + + it('makes the tool span active for the tool operation it wraps', async () => { + await withAgent(async () => { + instrumentation.observe( + { type: 'tool_start', toolCallId: 'call_1', toolName: 'get_weather', operationId: 'op_1' }, + {}, + ); + await instrumentation.interceptor({ type: 'tool', toolCallId: 'call_1' }, {}, async () => { + startInactiveSpan({ name: 'tool work' }).end(); + }); + instrumentation.observe({ type: 'tool', toolCallId: 'call_1', toolName: 'get_weather', operationId: 'op_1' }, {}); + }); + + expect(findSpan('tool work')?.parent_span_id).toBe(findSpan('execute_tool get_weather')?.span_id); + }); + it('opens a chat span on turn_start and completes it from the settled turn', async () => { await withAgent(() => { instrumentation.observe(