From 6aa7185abb7d2015215eb212c48089d39c72fdf8 Mon Sep 17 00:00:00 2001 From: Minwook Shin <163576506+minwookshin@users.noreply.github.com> Date: Sat, 3 Oct 2026 18:08:21 -0400 Subject: [PATCH 1/2] feat(bun): end server spans when streaming responses finish Co-Authored-By: GPT-6 --- .../suites/streaming/index.ts | 28 +++++ .../suites/streaming/test.ts | 42 +++++++ packages/bun/src/integrations/bunserver.ts | 22 +++- .../bun/test/integrations/bunserver.test.ts | 103 +++++++++++++++++- packages/cloudflare/src/utils/streaming.ts | 43 +------- packages/core/src/index.ts | 3 + packages/core/src/utils/responseStreaming.ts | 41 +++++++ .../test/lib/utils/responseStreaming.test.ts | 34 ++++++ packages/deno/src/utils/streaming.ts | 44 +------- 9 files changed, 273 insertions(+), 87 deletions(-) create mode 100644 dev-packages/bun-integration-tests/suites/streaming/index.ts create mode 100644 dev-packages/bun-integration-tests/suites/streaming/test.ts create mode 100644 packages/core/src/utils/responseStreaming.ts create mode 100644 packages/core/test/lib/utils/responseStreaming.test.ts 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..c717365d3b48 --- /dev/null +++ b/dev-packages/bun-integration-tests/suites/streaming/index.ts @@ -0,0 +1,28 @@ +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) { + if (new URL(request.url).pathname === '/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); + 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..88e5a285fb96 --- /dev/null +++ b/dev-packages/bun-integration-tests/suites/streaming/test.ts @@ -0,0 +1,42 @@ +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(); +}); diff --git a/packages/bun/src/integrations/bunserver.ts b/packages/bun/src/integrations/bunserver.ts index 080251c7dba3..09f1584cde4f 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, @@ -294,7 +296,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 +305,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 +321,28 @@ function wrapRequestHandler( ); } } + if (response?.body && !response.body.locked && classifyResponseStreaming(response).isStreaming) { + const { readable, writable } = new TransformStream(); + const streamedResponse = new Response(readable, { + status: response.status, + statusText: response.statusText, + headers: response.headers, + }); + // pipeTo settles on completion, source errors, and downstream cancellation. + void response.body.pipeTo(writable).then(endSpan, endSpan); + return streamedResponse; + } + 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/test/integrations/bunserver.test.ts b/packages/bun/test/integrations/bunserver.test.ts index 391b3c981c59..781c30480298 100644 --- a/packages/bun/test/integrations/bunserver.test.ts +++ b/packages/bun/test/integrations/bunserver.test.ts @@ -9,8 +9,9 @@ describe('Bun Serve Integration', () => { const mockSpan = SentryCore.startInactiveSpan({ name: 'test span' }); const setAttributesSpy = spyOn(mockSpan, 'setAttributes'); 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,6 +30,7 @@ describe('Bun Serve Integration', () => { beforeEach(() => { startSpanSpy.mockClear(); + endSpanSpy.mockReset(); continueTraceSpy.mockClear(); setAttributesSpy.mockClear(); // Header attributes are only collected while a client is active, so every test sets up its own instead of @@ -46,6 +48,102 @@ 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('ends the span and cancels the source when a streaming response is cancelled', async () => { + 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`)); + await response.body!.cancel('client disconnected'); + await ended; + expect(cancelled).toBe('client disconnected'); + expect(endSpanSpy).toHaveBeenCalledTimes(1); + } finally { + await server.stop(true); + } + }); + + test('ends the span and preserves the error when its response stream errors', 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 expect(response.text()).rejects.toBe(error); + await ended; + 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 +469,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. From e39a1ce6c00ea6380f5b535c6c0d1d41a37f074c Mon Sep 17 00:00:00 2001 From: Minwook Shin <163576506+minwookshin@users.noreply.github.com> Date: Mon, 5 Oct 2026 02:18:45 -0400 Subject: [PATCH 2/2] fix(bun): capture response stream failures without reporting cancellation Co-Authored-By: GPT-6 --- .../suites/streaming/index.ts | 7 +- .../suites/streaming/test.ts | 41 ++++++++++ packages/bun/src/integrations/bunserver.ts | 11 +-- packages/bun/src/utils/streaming.ts | 40 ++++++++++ .../bun/test/integrations/bunserver.test.ts | 79 +++++++++++++++---- 5 files changed, 158 insertions(+), 20 deletions(-) create mode 100644 packages/bun/src/utils/streaming.ts diff --git a/dev-packages/bun-integration-tests/suites/streaming/index.ts b/dev-packages/bun-integration-tests/suites/streaming/index.ts index c717365d3b48..aeb5dc2e6460 100644 --- a/dev-packages/bun-integration-tests/suites/streaming/index.ts +++ b/dev-packages/bun-integration-tests/suites/streaming/index.ts @@ -6,7 +6,8 @@ Sentry.init({ dsn: process.env.SENTRY_DSN, tracesSampleRate: 1 }); const server = Bun.serve({ port: 0, fetch(request) { - if (new URL(request.url).pathname === '/error') { + const path = new URL(request.url).pathname; + if (path === '/error') { throw new Error('handler failed'); } const span = Sentry.getActiveSpan(); @@ -14,6 +15,10 @@ const server = Bun.serve({ 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(); diff --git a/dev-packages/bun-integration-tests/suites/streaming/test.ts b/dev-packages/bun-integration-tests/suites/streaming/test.ts index 88e5a285fb96..25a0f138a04b 100644 --- a/dev-packages/bun-integration-tests/suites/streaming/test.ts +++ b/dev-packages/bun-integration-tests/suites/streaming/test.ts @@ -40,3 +40,44 @@ test('records a thrown handler as an error span', async () => { 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 09f1584cde4f..6265ee49f734 100644 --- a/packages/bun/src/integrations/bunserver.ts +++ b/packages/bun/src/integrations/bunserver.ts @@ -40,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; @@ -322,15 +323,15 @@ function wrapRequestHandler( } } if (response?.body && !response.body.locked && classifyResponseStreaming(response).isStreaming) { - const { readable, writable } = new TransformStream(); - const streamedResponse = new Response(readable, { + 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, }); - // pipeTo settles on completion, source errors, and downstream cancellation. - void response.body.pipeTo(writable).then(endSpan, endSpan); - return streamedResponse; } endSpan(); return response; 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 781c30480298..e21657a68c5f 100644 --- a/packages/bun/test/integrations/bunserver.test.ts +++ b/packages/bun/test/integrations/bunserver.test.ts @@ -8,6 +8,8 @@ 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 endSpanSpy = spyOn(mockSpan, 'end'); const startSpanSpy = spyOn(SentryCore, 'startSpanManual').mockImplementation((_opts, cb) => { @@ -33,6 +35,8 @@ describe('Bun Serve Integration', () => { 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(); @@ -57,7 +61,7 @@ describe('Bun Serve Integration', () => { controller = value; }, }); - const ended = new Promise(resolve => endSpanSpy.mockImplementation(resolve)); + 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({ @@ -85,32 +89,74 @@ describe('Bun Serve Integration', () => { }, ); - test('ends the span and cancels the source when a streaming response is cancelled', async () => { - let cancelled: unknown; - const ended = new Promise(resolve => endSpanSpy.mockImplementation(resolve)); - const source = new ReadableStream({ - cancel(reason) { - cancelled = reason; + 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`)); - await response.body!.cancel('client disconnected'); - await ended; - expect(cancelled).toBe('client disconnected'); + 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('ends the span and preserves the error when its response stream errors', async () => { + 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 ended = new Promise(resolve => endSpanSpy.mockImplementation(() => resolve())); const source = new ReadableStream({ start(value) { controller = value; @@ -124,8 +170,13 @@ describe('Bun Serve Integration', () => { const response = await server.fetch(new Request(`http://localhost:${port}/events`)); const error = new Error('stream interrupted'); controller!.error(error); - await expect(response.text()).rejects.toBe(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);