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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 33 additions & 0 deletions dev-packages/bun-integration-tests/suites/streaming/index.ts
Original file line number Diff line number Diff line change
@@ -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!);
83 changes: 83 additions & 0 deletions dev-packages/bun-integration-tests/suites/streaming/test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
23 changes: 20 additions & 3 deletions packages/bun/src/integrations/bunserver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import type { Integration, IntegrationFn, MaxRequestBodySize, SpanAttributes } f
import {
captureBodyFromWinterCGRequest,
captureException,
classifyResponseStreaming,
continueTrace,
defineIntegration,
getClient,
Expand All @@ -14,7 +15,8 @@ import {
parseStringToURLObject,
SEMANTIC_ATTRIBUTE_SENTRY_ORIGIN,
setHttpStatus,
startSpan,
SPAN_STATUS_ERROR,
startSpanManual,
winterCGRequestToRequestData,
withIsolationScope,
filterCollectedUrl,
Expand All @@ -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;

Expand Down Expand Up @@ -294,7 +297,7 @@ function wrapRequestHandler<T extends RouteHandler = RouteHandler>(
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.
Expand All @@ -303,7 +306,7 @@ function wrapRequestHandler<T extends RouteHandler = RouteHandler>(
? `${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) {
Expand All @@ -319,14 +322,28 @@ function wrapRequestHandler<T extends RouteHandler = RouteHandler>(
);
}
}
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;
}
},
Expand Down
40 changes: 40 additions & 0 deletions packages/bun/src/utils/streaming.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
/** Monitor the source without reading ahead or treating cancellation as a source failure. */
export function monitorStream(
stream: ReadableStream<Uint8Array>,
onDone: () => void,
onError: (error: unknown) => void,
): ReadableStream<Uint8Array> {
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<Uint8Array>(
{
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 },
);
}
Loading