diff --git a/dev-packages/node-integration-tests/suites/tracing/vercelai/v6_v7/instrument-abort.mjs b/dev-packages/node-integration-tests/suites/tracing/vercelai/v6_v7/instrument-abort.mjs new file mode 100644 index 000000000000..46a27dd03b74 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/vercelai/v6_v7/instrument-abort.mjs @@ -0,0 +1,9 @@ +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, + transport: loggingTransport, +}); diff --git a/dev-packages/node-integration-tests/suites/tracing/vercelai/v6_v7/scenario-aborted-stream-text.mjs b/dev-packages/node-integration-tests/suites/tracing/vercelai/v6_v7/scenario-aborted-stream-text.mjs new file mode 100644 index 000000000000..d7f4e7755f32 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/tracing/vercelai/v6_v7/scenario-aborted-stream-text.mjs @@ -0,0 +1,34 @@ +import * as Sentry from '@sentry/node'; +import { streamText } from 'ai'; +import { MockLanguageModelV3 } from 'ai/test'; + +async function run() { + await Sentry.startSpan({ op: 'function', name: 'main' }, async () => { + const controller = new AbortController(); + + // Abort the moment the model is asked for a stream, so the operation fails before its first + // chunk. `fetch()` rejects with the signal's reason on abort, and `@hono/node-server` aborts + // with a plain string — so there is no `AbortError` name to suppress by. + const model = new MockLanguageModelV3({ + doStream: ({ abortSignal }) => + new Promise((_, reject) => { + abortSignal.addEventListener('abort', () => reject(abortSignal.reason), { once: true }); + controller.abort('Client connection prematurely closed.'); + }), + }); + + const result = streamText({ + experimental_telemetry: { isEnabled: true, recordInputs: true, recordOutputs: true }, + maxRetries: 0, + model, + prompt: 'Stream me a response', + abortSignal: controller.signal, + }); + + for await (const _part of result.textStream) { + void _part; + } + }); +} + +run(); diff --git a/dev-packages/node-integration-tests/suites/tracing/vercelai/v6_v7/test.ts b/dev-packages/node-integration-tests/suites/tracing/vercelai/v6_v7/test.ts index d014e9308862..2468f974af55 100644 --- a/dev-packages/node-integration-tests/suites/tracing/vercelai/v6_v7/test.ts +++ b/dev-packages/node-integration-tests/suites/tracing/vercelai/v6_v7/test.ts @@ -955,4 +955,45 @@ describe.each(matrix)('Vercel AI integration (version %s)', (version, vercelAiVe }, }, ); + + createEsmTests( + __dirname, + 'scenario-aborted-stream-text.mjs', + 'instrument-abort.mjs', + (createRunner, test) => { + test('aborting a stream with a non-AbortError reason leaves no unhandled rejection', async () => { + await createRunner().ensureNoErrorOutput().start().completed(); + }); + + test('an aborted stream finishes its spans with an error status and no result attributes', async () => { + await createRunner() + .expect({ + span: container => { + const invokeAgent = container.items.find( + span => span.attributes['sentry.op']?.value === 'gen_ai.invoke_agent', + )!; + expect(invokeAgent).toBeDefined(); + expect(invokeAgent.status).toBe('error'); + expect(invokeAgent.attributes[GEN_AI_REQUEST_MODEL]?.value).toBe('mock-model-id'); + expect(invokeAgent.attributes[GEN_AI_RESPONSE_MODEL]).toBeUndefined(); + expect(invokeAgent.attributes[GEN_AI_USAGE_TOTAL_TOKENS]).toBeUndefined(); + expect(invokeAgent.attributes[GEN_AI_OUTPUT_MESSAGES]).toBeUndefined(); + + const generateContent = container.items.find( + span => span.attributes['sentry.op']?.value === 'gen_ai.generate_content', + )!; + expect(generateContent).toBeDefined(); + expect(generateContent.status).toBe('error'); + }, + }) + .start() + .completed(); + }); + }, + { + additionalDependencies: { + ai: vercelAiVersion, + }, + }, + ); }); diff --git a/packages/server-utils/src/tracing-channel.ts b/packages/server-utils/src/tracing-channel.ts index fb3d602ce1eb..075057f47c31 100644 --- a/packages/server-utils/src/tracing-channel.ts +++ b/packages/server-utils/src/tracing-channel.ts @@ -26,6 +26,9 @@ export type TracingChannelPayloadWithSpan = TData & { * The context's active store value, used to restore the context for asyncStart continuations for callback-based tracing. */ _sentryCallerStore?: unknown; + + /** Set by Node's tracing channel when the traced operation failed. */ + error?: unknown; }; /* @@ -66,7 +69,7 @@ export interface TracingChannelLifeCycleOptions { deferSpanEnd?: (args: { span: Span; data: TracingChannelPayloadWithSpan; - /** Ends the span: `end()` on success, `end(error)` on failure. Idempotent. */ + /** Ends the span: `end()` on success, `end(error)` on failure (which marks `data` as errored). Idempotent. */ end: (error?: unknown) => void; }) => boolean; @@ -158,6 +161,9 @@ export function bindTracingChannelToSpan( ended = true; if (error !== undefined) { annotateSpanError(span, error); + // Without this the payload still looks successful, so `beforeSpanEnd` enriches the span from + // a `result` the operation never produced. + data.error = error; } endBoundSpan(data, beforeSpanEnd); diff --git a/packages/server-utils/test/tracing-channel.test.ts b/packages/server-utils/test/tracing-channel.test.ts index 9e1f15484134..693be7de75df 100644 --- a/packages/server-utils/test/tracing-channel.test.ts +++ b/packages/server-utils/test/tracing-channel.test.ts @@ -837,6 +837,32 @@ describe('bindTracingChannelToSpan', () => { expect(endSpy).toHaveBeenCalledTimes(1); }); + it('`end(error)` marks the payload as failed for `beforeSpanEnd`', () => { + installTestAsyncContextStrategy(); + initTestClient(); + const span = startInactiveSpan({ name: 'channel-span' }); + const beforeSpanEnd = vi.fn(); + let captured: (error?: unknown) => void = () => undefined; + const { channel } = bindTracingChannelToSpan( + tracingChannel<{ operation: string }>('test:defer:payload-error'), + () => span, + { + beforeSpanEnd, + deferSpanEnd({ end }) { + captured = end; + return true; + }, + }, + ); + + channel.traceSync(() => 'stream', { operation: 'read' }); + const error = new Error('stream aborted'); + captured(error); + + expect(beforeSpanEnd).toHaveBeenCalledTimes(1); + expect(beforeSpanEnd).toHaveBeenCalledWith(span, expect.objectContaining({ error })); + }); + it('captures the error via `end(error)` when `captureError` is set', () => { const captureExceptionSpy = vi.spyOn(SentryCore, 'captureException').mockReturnValue('event-id'); const { end } = setupDeferred('test:defer:capture', { captureError: true });