diff --git a/dev-packages/bun-integration-tests/node-suites/excludes.ts b/dev-packages/bun-integration-tests/node-suites/excludes.ts index 20124be5a264..9705eff5ff1e 100644 --- a/dev-packages/bun-integration-tests/node-suites/excludes.ts +++ b/dev-packages/bun-integration-tests/node-suites/excludes.ts @@ -2,7 +2,7 @@ // fails on Bun is skipped with `test.skipIf` on `RUNTIME` in the Node suite, not listed here. // Node-only features: ANR and native thread watchdogs, child processes, the AWS Lambda Node runtime, -// and `node:sqlite`, which `flue` needs. +// `node:sqlite`, which `flue` needs, and the Vercel keep-alive, which needs `http.server.response.finish`. const NODE_ONLY = [ 'suites/anr/test.ts', 'suites/aws-serverless/**', @@ -10,6 +10,7 @@ const NODE_ONLY = [ 'suites/child-process/test.ts', 'suites/thread-blocked-native/test.ts', 'suites/tracing/flue/test.ts', + 'suites/vercel/keep-alive/test.ts', ]; // Bun does not publish `http.server.request.start`, so `@sentry/node` creates no `http.server` diff --git a/dev-packages/node-integration-tests/suites/vercel/keep-alive/scenario.ts b/dev-packages/node-integration-tests/suites/vercel/keep-alive/scenario.ts new file mode 100644 index 000000000000..15b47d11e9b6 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/vercel/keep-alive/scenario.ts @@ -0,0 +1,124 @@ +// Models a Vercel Node.js function: the platform exposes a per-request context with `waitUntil`, +// stops accepting `waitUntil` work once the response has closed, and freezes the instance right after the +// registered work settles. Nothing else (unref'd timers, SIGTERM) runs after the freeze. +import { AsyncLocalStorage } from 'node:async_hooks'; +import * as http from 'node:http'; +import * as net from 'node:net'; +import type { BaseTransportOptions, Envelope, Transport, TransportMakeRequestResponse } from '@sentry/core'; +import * as Sentry from '@sentry/node'; + +interface RequestStore { + waitUntil: (task: Promise) => void; + closed: boolean; +} + +const als = new AsyncLocalStorage(); +const registered: Promise[] = []; +let late = 0; + +(globalThis as Record)[Symbol.for('@vercel/request-context')] = { + get: () => als.getStore(), +}; + +function loggingTransport(_options: BaseTransportOptions): Transport { + return { + send(envelope: Envelope): Promise { + // eslint-disable-next-line no-console + console.log(JSON.stringify(envelope)); + return Promise.resolve({ statusCode: 200 }); + }, + flush(): PromiseLike { + return Promise.resolve(true); + }, + }; +} + +// Cases: `default`, `no-tracing` (no `tracesSampleRate`), `head` (HEAD requests never get a server span), +// `slow` (the handler responds late) and `late-root` (the root span ends after the response closes). +const testCase = process.env.TEST_CASE || 'default'; + +Sentry.init({ + dsn: 'https://public@dsn.ingest.sentry.io/1337', + tracesSampleRate: testCase === 'no-tracing' ? undefined : 1, + enableLogs: true, + traceLifecycle: process.env.TRACE_LIFECYCLE === 'stream' ? 'stream' : 'static', + transport: loggingTransport, + integrations: testCase === 'late-root' ? [Sentry.httpIntegration({ spans: false })] : [], +}); + +const server = http.createServer(async (_req, res) => { + if (testCase === 'slow') { + await new Promise(resolve => setTimeout(resolve, 300)); + } + + const respond = (): void => { + Sentry.logger.info('keep alive log'); + Sentry.metrics.count('keep_alive_metric', 1); + if (testCase === 'no-tracing') { + Sentry.captureException(new Error('keep alive error')); + } + // Without a server span (HEAD), a manual span would become a root span itself and mask the case. + if (testCase !== 'head') { + Sentry.startSpan({ name: 'keep-alive-child' }, () => undefined); + } + res.writeHead(200); + res.end('ok'); + }; + + if (testCase === 'late-root') { + // Some frameworks end the root span a moment after the response has closed. + Sentry.startSpanManual({ name: 'GET /', op: 'http.server' }, span => { + res.once('close', () => setTimeout(() => span.end(), 20)); + respond(); + }); + } else { + respond(); + } +}); + +// As on Vercel, `http.server.request.start` fires before the request context exists. The context is active +// for the handler and the response it writes. +const originalEmit = server.emit; +server.emit = function (this: http.Server, event: string, ...args: unknown[]): boolean { + if (event !== 'request') { + return originalEmit.call(this, event, ...args); + } + + const store: RequestStore = { + closed: false, + waitUntil: task => { + if (store.closed) { + late++; + } else { + registered.push(task); + } + }, + }; + // Registered before the original emit, so it runs before Sentry's own 'close' listeners. + (args[1] as http.ServerResponse).once('close', () => { + store.closed = true; + }); + return als.run(store, () => originalEmit.call(this, event, ...args)); +} as typeof server.emit; + +server.listen(0, async () => { + const { port } = server.address() as net.AddressInfo; + + // A raw socket keeps the request free of trace propagation headers, which would otherwise carry a + // sampling decision into the server span. + await new Promise(resolve => { + const socket = net.connect(port, 'localhost', () => { + socket.write( + `${testCase === 'head' ? 'HEAD' : 'GET'} / HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n`, + ); + }); + socket.resume(); + socket.on('close', () => resolve()); + }); + + // The freeze: only what was registered in time keeps the instance alive. + await Promise.all(registered); + // eslint-disable-next-line no-console + console.log(`WAITUNTIL registered=${registered.length} late=${late}`); + process.exit(0); +}); diff --git a/dev-packages/node-integration-tests/suites/vercel/keep-alive/test.ts b/dev-packages/node-integration-tests/suites/vercel/keep-alive/test.ts new file mode 100644 index 000000000000..e104ddbfd051 --- /dev/null +++ b/dev-packages/node-integration-tests/suites/vercel/keep-alive/test.ts @@ -0,0 +1,110 @@ +import { afterAll, describe, expect, test, vi } from 'vitest'; +import { cleanupChildProcesses, createRunner } from '../../../utils/runner'; + +afterAll(() => { + cleanupChildProcesses(); +}); + +describe.each(['stream', 'static'] as const)('keeps the instance alive for telemetry (%s)', lifecycle => { + const logExpectation = { + log: { + items: [expect.objectContaining({ body: 'keep alive log' })], + }, + }; + + const rootSpanExpectation = + lifecycle === 'stream' + ? { + span: { + items: expect.arrayContaining([ + expect.objectContaining({ name: 'keep-alive-child', is_segment: false }), + expect.objectContaining({ is_segment: true }), + ]), + }, + } + : { + transaction: { + contexts: { trace: { op: 'http.server' } }, + spans: [expect.objectContaining({ description: 'keep-alive-child' })], + }, + }; + + const metricExpectation = { + trace_metric: { + items: [expect.objectContaining({ name: 'keep_alive_metric', value: 1 })], + }, + }; + + test('sends the root span, log and metric before the instance freezes', async () => { + const runner = createRunner(__dirname, 'scenario.ts') + .withEnv({ VERCEL: '1', TRACE_LIFECYCLE: lifecycle }) + .unordered() + .expect(rootSpanExpectation) + .expect(logExpectation) + .expect(metricExpectation) + .start(); + + await runner.completed(); + + await vi.waitFor(() => expect(runner.getLogs().join('\n')).toContain('WAITUNTIL registered=1 late=0')); + }); + + test('sends the root span of a request that responds late', async () => { + const runner = createRunner(__dirname, 'scenario.ts') + .withEnv({ VERCEL: '1', TRACE_LIFECYCLE: lifecycle, TEST_CASE: 'slow' }) + .unordered() + .expect(rootSpanExpectation) + .expect(logExpectation) + .expect(metricExpectation) + .start(); + + await runner.completed(); + + await vi.waitFor(() => expect(runner.getLogs().join('\n')).toContain('WAITUNTIL registered=1 late=0')); + }); + + test('sends a root span that ends after the response closes', async () => { + const runner = createRunner(__dirname, 'scenario.ts') + .withEnv({ VERCEL: '1', TRACE_LIFECYCLE: lifecycle, TEST_CASE: 'late-root' }) + .unordered() + .expect(rootSpanExpectation) + .expect(logExpectation) + .expect(metricExpectation) + .start(); + + await runner.completed(); + + await vi.waitFor(() => expect(runner.getLogs().join('\n')).toContain('WAITUNTIL registered=1 late=0')); + }); + + test('sends errors and logs when tracing is disabled', async () => { + const runner = createRunner(__dirname, 'scenario.ts') + .withEnv({ VERCEL: '1', TRACE_LIFECYCLE: lifecycle, TEST_CASE: 'no-tracing' }) + .unordered() + .expect({ + event: { + exception: { values: [expect.objectContaining({ value: 'keep alive error' })] }, + }, + }) + .expect(logExpectation) + .expect(metricExpectation) + .start(); + + await runner.completed(); + + await vi.waitFor(() => expect(runner.getLogs().join('\n')).toContain('WAITUNTIL registered=1 late=0')); + }); + + test('sends logs for requests that get no server span', async () => { + const runner = createRunner(__dirname, 'scenario.ts') + .withEnv({ VERCEL: '1', TRACE_LIFECYCLE: lifecycle, TEST_CASE: 'head' }) + .unordered() + .expect(logExpectation) + .expect(metricExpectation) + .start(); + + await runner.completed(); + + await vi.waitFor(() => expect(runner.getLogs().join('\n')).toContain('WAITUNTIL registered=1 late=0')); + }); +}); diff --git a/packages/node/src/sdk/index.ts b/packages/node/src/sdk/index.ts index b740fa3e6612..e4a9e64c6185 100644 --- a/packages/node/src/sdk/index.ts +++ b/packages/node/src/sdk/index.ts @@ -40,6 +40,7 @@ import { getSpotlightConfig } from '../utils/spotlight'; import { defaultStackParser, getSentryRelease } from './api'; import { NodeClient } from './client'; import { initOpenTelemetry } from './initOtel'; +import { setupVercelKeepAlive } from './vercel'; /** * Get the base default integrations shared by all Node SDK default-integration sets. @@ -189,13 +190,8 @@ function _init( updateScopeFromEnvVariables(); - // Ensure we flush events when vercel functions are ended - // See: https://vercel.com/docs/functions/functions-api-reference#sigterm-signal if (process.env.VERCEL) { - process.on('SIGTERM', async () => { - // We have 500ms for processing here, so we try to make sure to have enough time to send the events - await client.flush(200); - }); + setupVercelKeepAlive(client); } // Add Node SDK specific OpenTelemetry setup. `setupEventContextTrace` reads the active span from the diff --git a/packages/node/src/sdk/vercel.ts b/packages/node/src/sdk/vercel.ts new file mode 100644 index 000000000000..9a18c5e1c9a2 --- /dev/null +++ b/packages/node/src/sdk/vercel.ts @@ -0,0 +1,86 @@ +import type { Client, Span } from '@sentry/core'; +import { getActiveSpan, getRootSpan, GLOBAL_OBJ } from '@sentry/core'; +import type { HttpServerResponse } from '@sentry/core/server'; +import { _INTERNAL_safeUnref } from '@sentry/core/server'; +import { subscribeDiagnosticsChannel } from '@sentry/server-utils'; + +// On Vercel, `http.server.request.start` fires without the request context, while this channel fires with it. +const HTTP_ON_SERVER_RESPONSE_FINISH = 'http.server.response.finish'; + +// Some frameworks end the request's root span shortly after the response closes. +const REQUEST_END_TIMEOUT_MS = 2000; + +interface VercelRequestContextGlobal { + get?(): { waitUntil?: (task: Promise) => void } | undefined; +} + +/** + * Keeps Vercel Node.js functions alive until the SDK has sent the telemetry of each request. + * + * Vercel can suspend a function as soon as the response is sent, so buffered events, spans, logs and metrics + * (and the timers that flush them) arrive late or never. When a response + * finishes, we register one `waitUntil` that waits for the response to close and the request's root span to + * end (2 seconds at most), and then flushes the client. + */ +export function setupVercelKeepAlive(client: Client): void { + // Ensure we flush events when vercel functions are ended + // See: https://vercel.com/docs/functions/functions-api-reference#sigterm-signal + process.on('SIGTERM', async () => { + // We have 500ms for processing here, so we try to make sure to have enough time to send the events + await client.flush(200); + }); + + subscribeDiagnosticsChannel(HTTP_ON_SERVER_RESPONSE_FINISH, message => { + const { response } = message as { response?: HttpServerResponse }; + const requestContextGlobal: VercelRequestContextGlobal | undefined = + // @ts-expect-error Vercel sets this global, so `GLOBAL_OBJ` does not type it + GLOBAL_OBJ[Symbol.for('@vercel/request-context')]; + const requestContext = requestContextGlobal?.get?.(); + if (!response || !requestContext?.waitUntil) { + return; + } + + const activeSpan = getActiveSpan(); + const rootSpan = activeSpan && getRootSpan(activeSpan); + + requestContext.waitUntil( + waitForRequestEnd(client, response, rootSpan, REQUEST_END_TIMEOUT_MS).then(() => client.flush(2000)), + ); + }); +} + +function waitForRequestEnd( + client: Client, + response: HttpServerResponse, + rootSpan: Span | undefined, + timeout: number, +): Promise { + return new Promise(resolve => { + let responseClosed = false; + // Ended and unsampled spans are not recording. + let rootSpanEnded = !rootSpan?.isRecording(); + + const done = (): void => { + clearTimeout(timer); + unsubscribe(); + resolve(); + }; + const doneIfRequestEnded = (): void => { + if (responseClosed && rootSpanEnded) { + done(); + } + }; + + const timer = _INTERNAL_safeUnref(setTimeout(done, timeout)); + const unsubscribe = client.on('spanEnd', span => { + if (span === rootSpan) { + rootSpanEnded = true; + doneIfRequestEnded(); + } + }); + response.once('close', () => { + responseClosed = true; + doneIfRequestEnded(); + }); + }); +} diff --git a/packages/node/test/sdk/vercel.test.ts b/packages/node/test/sdk/vercel.test.ts new file mode 100644 index 000000000000..211feef8bc60 --- /dev/null +++ b/packages/node/test/sdk/vercel.test.ts @@ -0,0 +1,170 @@ +import { EventEmitter } from 'node:events'; +import type { Client, Span } from '@sentry/core'; +import type * as SentryCore from '@sentry/core'; +import { getActiveSpan } from '@sentry/core'; +import type * as SentryServerUtils from '@sentry/server-utils'; +import { subscribeDiagnosticsChannel } from '@sentry/server-utils'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { setupVercelKeepAlive } from '../../src/sdk/vercel'; + +vi.mock('@sentry/core', async importOriginal => ({ + ...(await importOriginal()), + getActiveSpan: vi.fn(), +})); + +vi.mock('@sentry/server-utils', async importOriginal => ({ + ...(await importOriginal()), + subscribeDiagnosticsChannel: vi.fn(), +})); + +const REQUEST_CONTEXT = Symbol.for('@vercel/request-context'); + +function createClient() { + const spanEndCallbacks = new Set<(span: Span) => void>(); + return { + on: vi.fn((_hook: 'spanEnd', callback: (span: Span) => void) => { + spanEndCallbacks.add(callback); + return () => spanEndCallbacks.delete(callback); + }), + flush: vi.fn().mockResolvedValue(true), + endSpan: (span: Span) => spanEndCallbacks.forEach(callback => callback(span)), + }; +} + +function setActiveRootSpan(): Span { + const span = { isRecording: () => true } as unknown as Span; + vi.mocked(getActiveSpan).mockReturnValue(span); + return span; +} + +function finishResponse() { + const response = new EventEmitter(); + const [, onMessage] = vi.mocked(subscribeDiagnosticsChannel).mock.calls.at(-1)!; + onMessage({ request: new EventEmitter(), response }, 'http.server.response.finish'); + return response; +} + +describe('setupVercelKeepAlive', () => { + let waitUntil: ReturnType; + + beforeEach(() => { + vi.useFakeTimers(); + waitUntil = vi.fn(); + (globalThis as any)[REQUEST_CONTEXT] = { get: () => ({ waitUntil }) }; + vi.spyOn(process, 'on').mockImplementation(() => process); + }); + + afterEach(() => { + (globalThis as any)[REQUEST_CONTEXT] = undefined; + vi.mocked(getActiveSpan).mockReturnValue(undefined); + vi.useRealTimers(); + vi.restoreAllMocks(); + }); + + it('registers exactly one waitUntil per finished response', () => { + setupVercelKeepAlive(createClient() as unknown as Client); + + // Vercel does not publish `http.server.request.start` for function requests. + expect(subscribeDiagnosticsChannel).toHaveBeenCalledWith('http.server.response.finish', expect.any(Function)); + + finishResponse(); + expect(waitUntil).toHaveBeenCalledTimes(1); + + finishResponse(); + expect(waitUntil).toHaveBeenCalledTimes(2); + }); + + it('keeps the task pending until the response closes, then flushes the client once', async () => { + const client = createClient(); + setupVercelKeepAlive(client as unknown as Client); + + const response = finishResponse(); + const task = waitUntil.mock.calls[0]![0] as Promise; + const settled = vi.fn(); + void task.then(settled); + + await vi.advanceTimersByTimeAsync(0); + expect(settled).not.toHaveBeenCalled(); + expect(client.flush).not.toHaveBeenCalled(); + + response.emit('close'); + await task; + + expect(client.flush).toHaveBeenCalledTimes(1); + expect(client.flush).toHaveBeenCalledWith(2000); + }); + + it('waits for the root span to end after the response closes, then flushes', async () => { + const client = createClient(); + setupVercelKeepAlive(client as unknown as Client); + const rootSpan = setActiveRootSpan(); + + const response = finishResponse(); + const task = waitUntil.mock.calls[0]![0] as Promise; + + response.emit('close'); + await vi.advanceTimersByTimeAsync(0); + expect(client.flush).not.toHaveBeenCalled(); + + client.endSpan(rootSpan); + await task; + + expect(client.flush).toHaveBeenCalledTimes(1); + }); + + it('flushes after 2 seconds when the root span does not end', async () => { + const client = createClient(); + setupVercelKeepAlive(client as unknown as Client); + setActiveRootSpan(); + + const response = finishResponse(); + response.emit('close'); + + await vi.advanceTimersByTimeAsync(1999); + expect(client.flush).not.toHaveBeenCalled(); + + await vi.advanceTimersByTimeAsync(1); + expect(client.flush).toHaveBeenCalledTimes(1); + }); + + it('flushes after 2 seconds when the response does not close', async () => { + const client = createClient(); + setupVercelKeepAlive(client as unknown as Client); + + finishResponse(); + + await vi.advanceTimersByTimeAsync(1999); + expect(client.flush).not.toHaveBeenCalled(); + + await vi.advanceTimersByTimeAsync(1); + expect(client.flush).toHaveBeenCalledTimes(1); + }); + + it('does nothing without a request context', () => { + (globalThis as any)[REQUEST_CONTEXT] = undefined; + setupVercelKeepAlive(createClient() as unknown as Client); + + finishResponse(); + + expect(waitUntil).not.toHaveBeenCalled(); + }); + + it('does nothing when the request context has no waitUntil', () => { + (globalThis as any)[REQUEST_CONTEXT] = { get: () => ({}) }; + setupVercelKeepAlive(createClient() as unknown as Client); + + expect(() => finishResponse()).not.toThrow(); + expect(waitUntil).not.toHaveBeenCalled(); + }); + + it('does not depend on a span', async () => { + const client = createClient(); + setupVercelKeepAlive(client as unknown as Client); + + const response = finishResponse(); + response.emit('close'); + await (waitUntil.mock.calls[0]![0] as Promise); + + expect(client.flush).toHaveBeenCalledTimes(1); + }); +});