diff --git a/dev-packages/bun-integration-tests/suites/streaming/index.ts b/dev-packages/bun-integration-tests/suites/streaming/index.ts new file mode 100644 index 000000000000..aeb5dc2e6460 --- /dev/null +++ b/dev-packages/bun-integration-tests/suites/streaming/index.ts @@ -0,0 +1,33 @@ +import { sendPortToRunner } from '@sentry-internal/node-integration-tests'; +import * as Sentry from '@sentry/bun'; + +Sentry.init({ dsn: process.env.SENTRY_DSN, tracesSampleRate: 1 }); + +const server = Bun.serve({ + port: 0, + fetch(request) { + const path = new URL(request.url).pathname; + if (path === '/error') { + throw new Error('handler failed'); + } + const span = Sentry.getActiveSpan(); + const body = new ReadableStream({ + async start(controller) { + controller.enqueue(new TextEncoder().encode('data: first\n\n')); + await Bun.sleep(50); + if (path === '/stream-error') { + controller.error(new Error('stream failed')); + return; + } + controller.enqueue(new TextEncoder().encode(`data: recording=${span?.isRecording()}\n\n`)); + span?.setAttribute('test.stream.completed', true); + controller.close(); + }, + }); + return new Response(body, { status: 201, headers: { 'content-type': 'text/event-stream' } }); + }, + error() { + return new Response('failed', { status: 500 }); + }, +}); +sendPortToRunner(server.port!); diff --git a/dev-packages/bun-integration-tests/suites/streaming/test.ts b/dev-packages/bun-integration-tests/suites/streaming/test.ts new file mode 100644 index 000000000000..25a0f138a04b --- /dev/null +++ b/dev-packages/bun-integration-tests/suites/streaming/test.ts @@ -0,0 +1,83 @@ +import { afterAll, expect, test } from 'vitest'; +import { cleanupChildProcesses, createRunner } from '../../../node-integration-tests/utils/runner'; + +afterAll(cleanupChildProcesses); + +test('records streaming work after the response headers have been returned', async () => { + const runner = createRunner(__dirname, 'index.ts') + .withMockSentryServer() + .unordered() + .expect({ + span: container => { + const root = container.items.find(item => item.is_segment); + expect(root).toEqual( + expect.objectContaining({ + attributes: expect.objectContaining({ + 'url.path': { value: '/events', type: 'string' }, + 'http.response.status_code': { value: 201, type: 'integer' }, + 'test.stream.completed': { value: true, type: 'boolean' }, + }), + }), + ); + }, + }) + .start(); + expect(await runner.makeRequest('get', '/events')).toBe('data: first\n\ndata: recording=true\n\n'); + await runner.completed(); +}); + +test('records a thrown handler as an error span', async () => { + const runner = createRunner(__dirname, 'index.ts') + .withMockSentryServer() + .unordered() + .expect({ + span: container => { + const root = container.items.find(item => item.is_segment); + expect(root).toEqual(expect.objectContaining({ status: 'error' })); + }, + }) + .start(); + await runner.makeRequest('get', '/error', { expectError: true }); + await runner.completed(); +}); + +test('captures a response stream failure with its request context and error span', async () => { + let spanTraceId: string | undefined; + let eventTraceId: string | undefined; + const runner = createRunner(__dirname, 'index.ts') + .withMockSentryServer() + .unordered() + .expect({ + span: container => { + const roots = container.items.filter(item => item.is_segment); + expect(roots).toHaveLength(1); + const root = roots[0]!; + expect(root.status).toBe('error'); + expect(root.attributes['url.path']?.value).toBe('/stream-error'); + expect(root.attributes['http.response.status_code']?.value).toBe(201); + spanTraceId = root.trace_id; + }, + }) + .expect({ + event: event => { + expect(event.exception?.values).toHaveLength(1); + const exception = event.exception!.values![0]!; + expect(exception.type).toBe('Error'); + expect(exception.value).toBe('stream failed'); + expect(exception.mechanism).toEqual({ type: 'auto.http.bun.serve', handled: false }); + expect(event.request?.method).toBe('GET'); + expect(new URL(event.request!.url!).pathname).toBe('/stream-error'); + eventTraceId = event.contexts?.trace?.trace_id; + }, + }) + .start(); + + await expect.poll(() => runner.getPort(), { timeout: 30_000 }).toBeTypeOf('number'); + const response = await fetch(`http://localhost:${runner.getPort()}/stream-error`); + expect(response.status).toBe(201); + // Bun versions differ in whether a failed stream rejects or closes the HTTP response body. + await Promise.allSettled([response.text()]); + await runner.completed(); + expect(spanTraceId).toMatch(/^[a-f\d]{32}$/); + expect(eventTraceId).toBe(spanTraceId); +}); diff --git a/packages/bun/src/integrations/bunserver.ts b/packages/bun/src/integrations/bunserver.ts index 080251c7dba3..6265ee49f734 100644 --- a/packages/bun/src/integrations/bunserver.ts +++ b/packages/bun/src/integrations/bunserver.ts @@ -2,6 +2,7 @@ import type { Integration, IntegrationFn, MaxRequestBodySize, SpanAttributes } f import { captureBodyFromWinterCGRequest, captureException, + classifyResponseStreaming, continueTrace, defineIntegration, getClient, @@ -14,7 +15,8 @@ import { parseStringToURLObject, SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN, setHttpStatus, - startSpan, + SPAN_STATUS_ERROR, + startSpanManual, winterCGRequestToRequestData, withIsolationScope, filterCollectedUrl, @@ -38,6 +40,7 @@ import { URL_SCHEME, } from '@sentry/conventions/attributes'; import { HTTP_SERVER } from '@sentry/conventions/op'; +import { monitorStream } from '../utils/streaming'; const INTEGRATION_NAME = 'BunServer' as const; @@ -294,7 +297,7 @@ function wrapRequestHandler( baggage: request.headers.get('baggage'), }, () => - startSpan( + startSpanManual( { attributes: { ...attributes, [SENTRY_OP]: HTTP_SERVER }, // With span streaming, span names have to be low cardinality, so we can't fall back to the URL path. @@ -303,7 +306,7 @@ function wrapRequestHandler( ? `${request.method} ${routeName}` : request.method?.toUpperCase() || HTTP_SPAN_NAME_FALLBACK, }, - async span => { + async (span, endSpan) => { try { const response = (await target.apply(thisArg, args)) as Response | undefined; if (response?.status) { @@ -319,14 +322,28 @@ function wrapRequestHandler( ); } } + if (response?.body && !response.body.locked && classifyResponseStreaming(response).isStreaming) { + const body = monitorStream(response.body, endSpan, error => { + span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' }); + captureException(error, { mechanism: { type: 'auto.http.bun.serve', handled: false } }); + }); + return new Response(body, { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }); + } + endSpan(); return response; } catch (e) { + span.setStatus({ code: SPAN_STATUS_ERROR, message: 'internal_error' }); captureException(e, { mechanism: { type: 'auto.http.bun.serve', handled: false, }, }); + endSpan(); throw e; } }, diff --git a/packages/bun/src/utils/streaming.ts b/packages/bun/src/utils/streaming.ts new file mode 100644 index 000000000000..d19466a6a035 --- /dev/null +++ b/packages/bun/src/utils/streaming.ts @@ -0,0 +1,40 @@ +/** Monitor the source without reading ahead or treating cancellation as a source failure. */ +export function monitorStream( + stream: ReadableStream, + onDone: () => void, + onError: (error: unknown) => void, +): ReadableStream { + const reader = stream.getReader(); + // A cancelled reader closes normally; only a source failure should be captured as an error. + void reader.closed.then(onDone, error => { + onError(error); + onDone(); + }); + + return new ReadableStream( + { + async pull(controller) { + try { + const { done, value } = await reader.read(); + if (done) { + controller.close(); + reader.releaseLock(); + } else { + controller.enqueue(value); + } + } catch (error) { + controller.error(error); + reader.releaseLock(); + } + }, + async cancel(reason) { + try { + await reader.cancel(reason); + } finally { + reader.releaseLock(); + } + }, + }, + { highWaterMark: 0 }, + ); +} diff --git a/packages/bun/test/integrations/bunserver.test.ts b/packages/bun/test/integrations/bunserver.test.ts index 391b3c981c59..e21657a68c5f 100644 --- a/packages/bun/test/integrations/bunserver.test.ts +++ b/packages/bun/test/integrations/bunserver.test.ts @@ -8,9 +8,12 @@ import { instrumentBunServe } from '../../src/integrations/bunserver'; describe('Bun Serve Integration', () => { const mockSpan = SentryCore.startInactiveSpan({ name: 'test span' }); const setAttributesSpy = spyOn(mockSpan, 'setAttributes'); + const setStatusSpy = spyOn(mockSpan, 'setStatus'); + const captureExceptionSpy = spyOn(SentryCore, 'captureException'); const continueTraceSpy = spyOn(SentryCore, 'continueTrace'); - const startSpanSpy = spyOn(SentryCore, 'startSpan').mockImplementation((_opts, cb) => { - return cb(mockSpan as unknown as SentryCore.Span); + const endSpanSpy = spyOn(mockSpan, 'end'); + const startSpanSpy = spyOn(SentryCore, 'startSpanManual').mockImplementation((_opts, cb) => { + return cb(mockSpan as unknown as SentryCore.Span, () => mockSpan.end()); }); const setupClient = (options?: BunOptions): void => { @@ -29,8 +32,11 @@ describe('Bun Serve Integration', () => { beforeEach(() => { startSpanSpy.mockClear(); + endSpanSpy.mockReset(); continueTraceSpy.mockClear(); setAttributesSpy.mockClear(); + setStatusSpy.mockClear(); + captureExceptionSpy.mockClear(); // Header attributes are only collected while a client is active, so every test sets up its own instead of // relying on one leaking in from whichever test file `bun test` happened to run first. setupClient(); @@ -46,6 +52,149 @@ describe('Bun Serve Integration', () => { port += 1; }); + test.each(['fetch', 'route'])( + 'keeps a streaming %s response span open until the last chunk is consumed', + async mode => { + let controller: ReadableStreamDefaultController; + const source = new ReadableStream({ + start(value) { + controller = value; + }, + }); + const ended = new Promise(resolve => endSpanSpy.mockImplementation(() => resolve())); + const handler = () => + new Response(source, { status: 201, headers: { 'content-type': 'text/event-stream', 'x-stream': 'events' } }); + const server = Bun.serve({ + port, + ...(mode === 'fetch' ? { fetch: handler } : { routes: { '/events': { GET: handler } } }), + }); + try { + controller!.enqueue(new TextEncoder().encode('data: first\n\n')); + const response = await fetch(`http://localhost:${port}/events`); + const reader = response.body!.getReader(); + expect(response.status).toBe(201); + expect(response.headers.get('x-stream')).toBe('events'); + expect(new TextDecoder().decode((await reader.read()).value)).toBe('data: first\n\n'); + expect(endSpanSpy).toHaveBeenCalledTimes(0); + + controller!.enqueue(new TextEncoder().encode('data: last\n\n')); + controller!.close(); + expect(new TextDecoder().decode((await reader.read()).value)).toBe('data: last\n\n'); + expect((await reader.read()).done).toBe(true); + await ended; + expect(endSpanSpy).toHaveBeenCalledTimes(1); + } finally { + await server.stop(true); + } + }, + ); + + test.each(['client disconnected', new Error('client disconnected')])( + 'ends the span without capturing a cancelled stream: %s', + async reason => { + let cancelled: unknown; + const ended = new Promise(resolve => endSpanSpy.mockImplementation(() => resolve())); + const source = new ReadableStream({ + cancel(reason) { + cancelled = reason; + }, + }); + const server = Bun.serve({ + port, + fetch: () => new Response(source, { headers: { 'content-type': 'text/event-stream' } }), + }); + try { + const response = await server.fetch(new Request(`http://localhost:${port}/events`)); + const reader = response.body!.getReader(); + const pendingRead = reader.read(); + await reader.cancel(reason); + expect(await pendingRead).toEqual({ done: true, value: undefined }); + await ended; + expect(cancelled).toBe(reason); + expect(endSpanSpy).toHaveBeenCalledTimes(1); + expect(captureExceptionSpy).not.toHaveBeenCalled(); + expect(setStatusSpy).not.toHaveBeenCalledWith({ + code: SentryCore.SPAN_STATUS_ERROR, + message: 'internal_error', + }); + } finally { + await server.stop(true); + } + }, + ); + + test('preserves backpressure instead of reading ahead of the response consumer', async () => { + let pulls = 0; + const source = new ReadableStream( + { + pull(controller) { + pulls += 1; + controller.enqueue(new TextEncoder().encode(`data: ${pulls}\n\n`)); + }, + }, + { highWaterMark: 0 }, + ); + const server = Bun.serve({ + port, + fetch: () => new Response(source, { headers: { 'content-type': 'text/event-stream' } }), + }); + try { + const response = await server.fetch(new Request(`http://localhost:${port}/events`)); + expect(pulls).toBe(0); + const reader = response.body!.getReader(); + expect(new TextDecoder().decode((await reader.read()).value)).toBe('data: 1\n\n'); + expect(pulls).toBe(1); + expect(endSpanSpy).not.toHaveBeenCalled(); + + await reader.cancel(); + expect(endSpanSpy).toHaveBeenCalledTimes(1); + expect(captureExceptionSpy).not.toHaveBeenCalled(); + } finally { + await server.stop(true); + } + }); + + test('captures a response stream error and marks its span as failed', async () => { + let controller: ReadableStreamDefaultController; + const ended = new Promise(resolve => endSpanSpy.mockImplementation(() => resolve())); + const source = new ReadableStream({ + start(value) { + controller = value; + }, + }); + const server = Bun.serve({ + port, + fetch: () => new Response(source, { headers: { 'content-type': 'application/x-ndjson' } }), + }); + try { + const response = await server.fetch(new Request(`http://localhost:${port}/events`)); + const error = new Error('stream interrupted'); + controller!.error(error); + await ended; + await expect(response.text()).rejects.toBe(error); + expect(captureExceptionSpy).toHaveBeenCalledTimes(1); + expect(captureExceptionSpy).toHaveBeenCalledWith(error, { + mechanism: { type: 'auto.http.bun.serve', handled: false }, + }); + expect(setStatusSpy).toHaveBeenCalledWith({ code: SentryCore.SPAN_STATUS_ERROR, message: 'internal_error' }); + expect(endSpanSpy).toHaveBeenCalledTimes(1); + } finally { + await server.stop(true); + } + }); + + test('ends a non-streaming response span before the body is consumed', async () => { + const server = Bun.serve({ port, fetch: () => Response.json({ ok: true }) }); + try { + const response = await server.fetch(new Request(`http://localhost:${port}/status`)); + expect(endSpanSpy).toHaveBeenCalledTimes(1); + expect(await response.json()).toEqual({ ok: true }); + expect(endSpanSpy).toHaveBeenCalledTimes(1); + } finally { + await server.stop(true); + } + }); + test('generates a transaction around a request', async () => { const server = Bun.serve({ async fetch(_req) { @@ -371,6 +520,7 @@ describe('Bun Serve Integration', () => { expect(await initialResponse.text()).toBe('Initial handler'); expect(startSpanSpy).toHaveBeenCalledTimes(1); startSpanSpy.mockClear(); + endSpanSpy.mockReset(); // Reload server with new handler server.reload({ diff --git a/packages/cloudflare/src/utils/streaming.ts b/packages/cloudflare/src/utils/streaming.ts index fee67cbb9f2a..c057385724aa 100644 --- a/packages/cloudflare/src/utils/streaming.ts +++ b/packages/cloudflare/src/utils/streaming.ts @@ -1,41 +1,2 @@ -export type StreamingGuess = { - isStreaming: boolean; -}; - -/** - * Classifies a Response as streaming or non-streaming. - * - * Heuristics: - * - No body → not streaming - * - Known streaming Content-Types → streaming (SSE, NDJSON, JSON streaming) - * - text/plain without Content-Length → streaming (some AI APIs) - * - Otherwise → not streaming (conservative default, including HTML/SSR) - * - * We avoid probing the stream to prevent blocking on transform streams (like injectTraceMetaTags) - * or SSR streams that may not have data ready immediately. - */ -export function classifyResponseStreaming(res: Response): StreamingGuess { - if (!res.body) { - return { isStreaming: false }; - } - - const contentType = res.headers.get('content-type') ?? ''; - const contentLength = res.headers.get('content-length'); - - // Streaming: Known streaming content types - // - text/event-stream: Server-Sent Events (Vercel AI SDK, real-time APIs) - // - application/x-ndjson, application/ndjson: Newline-delimited JSON - // - application/stream+json: JSON streaming - // - text/plain (without Content-Length): Some AI APIs use this for streaming text - if ( - /^text\/event-stream\b/i.test(contentType) || - /^application\/(x-)?ndjson\b/i.test(contentType) || - /^application\/stream\+json\b/i.test(contentType) || - (/^text\/plain\b/i.test(contentType) && !contentLength) - ) { - return { isStreaming: true }; - } - - // Default: treat as non-streaming - return { isStreaming: false }; -} +export { classifyResponseStreaming } from '@sentry/core'; +export type { StreamingGuess } from '@sentry/core'; diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 20b52be30603..42018b33bb08 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -541,3 +541,6 @@ export { } from './utils/randomSafeContext'; export { warnOnRemovedBuildOptions } from './build-time-plugins/warnOnRemovedBuildOptions'; export { startSpan, startInactiveSpan, startSpanManual } from './tracing/trace'; + +export { classifyResponseStreaming } from './utils/responseStreaming'; +export type { StreamingGuess } from './utils/responseStreaming'; diff --git a/packages/core/src/utils/responseStreaming.ts b/packages/core/src/utils/responseStreaming.ts new file mode 100644 index 000000000000..fee67cbb9f2a --- /dev/null +++ b/packages/core/src/utils/responseStreaming.ts @@ -0,0 +1,41 @@ +export type StreamingGuess = { + isStreaming: boolean; +}; + +/** + * Classifies a Response as streaming or non-streaming. + * + * Heuristics: + * - No body → not streaming + * - Known streaming Content-Types → streaming (SSE, NDJSON, JSON streaming) + * - text/plain without Content-Length → streaming (some AI APIs) + * - Otherwise → not streaming (conservative default, including HTML/SSR) + * + * We avoid probing the stream to prevent blocking on transform streams (like injectTraceMetaTags) + * or SSR streams that may not have data ready immediately. + */ +export function classifyResponseStreaming(res: Response): StreamingGuess { + if (!res.body) { + return { isStreaming: false }; + } + + const contentType = res.headers.get('content-type') ?? ''; + const contentLength = res.headers.get('content-length'); + + // Streaming: Known streaming content types + // - text/event-stream: Server-Sent Events (Vercel AI SDK, real-time APIs) + // - application/x-ndjson, application/ndjson: Newline-delimited JSON + // - application/stream+json: JSON streaming + // - text/plain (without Content-Length): Some AI APIs use this for streaming text + if ( + /^text\/event-stream\b/i.test(contentType) || + /^application\/(x-)?ndjson\b/i.test(contentType) || + /^application\/stream\+json\b/i.test(contentType) || + (/^text\/plain\b/i.test(contentType) && !contentLength) + ) { + return { isStreaming: true }; + } + + // Default: treat as non-streaming + return { isStreaming: false }; +} diff --git a/packages/core/test/lib/utils/responseStreaming.test.ts b/packages/core/test/lib/utils/responseStreaming.test.ts new file mode 100644 index 000000000000..5ec7779bd7e8 --- /dev/null +++ b/packages/core/test/lib/utils/responseStreaming.test.ts @@ -0,0 +1,34 @@ +import { describe, expect, test } from 'vitest'; +import { classifyResponseStreaming } from '../../../src/utils/responseStreaming'; + +describe('classifyResponseStreaming', () => { + test.each([ + ['text/event-stream; charset=utf-8', true], + ['application/x-ndjson', true], + ['application/ndjson', true], + ['application/stream+json', true], + ['TEXT/EVENT-STREAM', true], + ['text/plain', true], + ['text/html', false], + ['application/json', false], + ['application/octet-stream', false], + ])('classifies %s without reading the body', (contentType, isStreaming) => { + const response = new Response(new ReadableStream(), { headers: { 'content-type': contentType } }); + + expect(classifyResponseStreaming(response)).toEqual({ isStreaming }); + expect(response.bodyUsed).toBe(false); + expect(response.body?.locked).toBe(false); + }); + + test('treats plain text with a known length as non-streaming', () => { + const response = new Response('ready', { headers: { 'content-type': 'text/plain', 'content-length': '5' } }); + + expect(classifyResponseStreaming(response)).toEqual({ isStreaming: false }); + }); + + test('treats a bodyless response as non-streaming even with streaming headers', () => { + const response = new Response(null, { status: 204, headers: { 'content-type': 'text/event-stream' } }); + + expect(classifyResponseStreaming(response)).toEqual({ isStreaming: false }); + }); +}); diff --git a/packages/deno/src/utils/streaming.ts b/packages/deno/src/utils/streaming.ts index 491131e6e093..1385d2022765 100644 --- a/packages/deno/src/utils/streaming.ts +++ b/packages/deno/src/utils/streaming.ts @@ -1,46 +1,8 @@ import type { Span } from '@sentry/core'; +import { classifyResponseStreaming } from '@sentry/core'; -export type StreamingGuess = { - isStreaming: boolean; -}; - -/** - * Classifies a Response as streaming or non-streaming. - * - * Heuristics: - * - No body → not streaming - * - Known streaming Content-Types → streaming (SSE, NDJSON, JSON streaming) - * - text/plain without Content-Length → streaming (some AI APIs) - * - Otherwise → not streaming (conservative default, including HTML/SSR) - * - * We avoid probing the stream to prevent blocking on transform streams (like injectTraceMetaTags) - * or SSR streams that may not have data ready immediately. - */ -export function classifyResponseStreaming(res: Response): StreamingGuess { - if (!res.body) { - return { isStreaming: false }; - } - - const contentType = res.headers.get('content-type') ?? ''; - const contentLength = res.headers.get('content-length'); - - // Streaming: Known streaming content types - // - text/event-stream: Server-Sent Events (Vercel AI SDK, real-time APIs) - // - application/x-ndjson, application/ndjson: Newline-delimited JSON - // - application/stream+json: JSON streaming - // - text/plain (without Content-Length): Some AI APIs use this for streaming text - if ( - /^text\/event-stream\b/i.test(contentType) || - /^application\/(x-)?ndjson\b/i.test(contentType) || - /^application\/stream\+json\b/i.test(contentType) || - (/^text\/plain\b/i.test(contentType) && !contentLength) - ) { - return { isStreaming: true }; - } - - // Default: treat as non-streaming - return { isStreaming: false }; -} +export { classifyResponseStreaming } from '@sentry/core'; +export type { StreamingGuess } from '@sentry/core'; /** * Tee a stream, and end the provided span when the stream ends.