diff --git a/dev-packages/cloudflare-integration-tests/suites/flush-timeout/index.ts b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/index.ts new file mode 100644 index 000000000000..4e00dbb4109c --- /dev/null +++ b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/index.ts @@ -0,0 +1,52 @@ +import * as Sentry from '@sentry/cloudflare'; +import { WorkflowEntrypoint } from 'cloudflare:workers'; +import type { WorkflowEvent, WorkflowStep } from 'cloudflare:workers'; +import { lastSend } from './lastSend'; + +interface Env { + SERVER_URL: string; + ISSUE_WORKFLOW: Workflow; +} + +// The Workflow from https://github.com/getsentry/sentry-javascript/issues/24482. After its steps it flushes +// with the same timeout the SDK uses after each step and reports to SERVER_URL how the pending send ended. +export class IssueWorkflow extends WorkflowEntrypoint { + async run(_event: WorkflowEvent, step: WorkflowStep): Promise { + for (let index = 0; index < 100; index++) { + await step.do(`step-${index}`, async () => index); + } + + // Locally the step spans are not sent before `run()` returns, so flush here with the timeout the SDK uses + // after each step, while the Workflow can still observe whether the pending send gets aborted. + lastSend.aborted = false; + const flushed = await Sentry.flush(2000); + const send = lastSend.aborted ? 'aborted' : 'not aborted'; + await fetch(`${this.env.SERVER_URL}/result`, { method: 'POST', body: JSON.stringify({ flushed, send }) }); + } +} + +export default { + async fetch(request, env, ctx) { + const url = new URL(request.url); + + if (url.pathname === '/workflow/trigger') { + const instance = await env.ISSUE_WORKFLOW.create(); + return Response.json({ id: instance.id }); + } + + // The flush runs inside the invocation, so the send is still pending when its drain times out. + if (url.pathname === '/flush-with-timeout') { + Sentry.captureException(new Error('Captured on /flush-with-timeout')); + lastSend.aborted = false; + const flushed = await Sentry.flush(500); + return Response.json({ flushed, send: lastSend.aborted ? 'aborted' : 'not aborted' }); + } + + if (url.pathname === '/pending-wait-until') { + ctx.waitUntil(new Promise(resolve => setTimeout(resolve, 120_000))); + Sentry.captureException(new Error('Captured on /pending-wait-until')); + } + + return new Response('ok'); + }, +} satisfies ExportedHandler; diff --git a/dev-packages/cloudflare-integration-tests/suites/flush-timeout/instrument.server.ts b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/instrument.server.ts new file mode 100644 index 000000000000..13a9e52db068 --- /dev/null +++ b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/instrument.server.ts @@ -0,0 +1,25 @@ +import { defineCloudflareOptions } from '@sentry/cloudflare'; +import { lastSend } from './lastSend'; + +interface Env { + SENTRY_DSN: string; + SERVER_URL: string; + // "true" sends envelopes to SERVER_URL, a server that never answers + SLOW_INGEST?: string; + // "false" creates one client per invocation, which waits for the invocation's flush lock + CACHE_CLIENT?: string; + // "true" samples every trace, so the Workflow steps create spans to send + TRACING?: string; +} + +export default defineCloudflareOptions((env: Env) => ({ + dsn: env.SLOW_INGEST === 'true' ? `${env.SERVER_URL.replace('://', '://public@')}/1337` : env.SENTRY_DSN, + cacheClient: env.CACHE_CLIENT !== 'false', + tracesSampleRate: env.TRACING === 'true' ? 1 : undefined, + transportOptions: { + fetch: (input, init) => { + init?.signal?.addEventListener('abort', () => (lastSend.aborted = true)); + return fetch(input, init); + }, + }, +})); diff --git a/dev-packages/cloudflare-integration-tests/suites/flush-timeout/lastSend.ts b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/lastSend.ts new file mode 100644 index 000000000000..d9a7b995a130 --- /dev/null +++ b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/lastSend.ts @@ -0,0 +1,2 @@ +// Whether the transport aborted an envelope fetch since `aborted` was last reset. +export const lastSend: { aborted: boolean } = { aborted: false }; diff --git a/dev-packages/cloudflare-integration-tests/suites/flush-timeout/test.ts b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/test.ts new file mode 100644 index 000000000000..4eee8b06a45c --- /dev/null +++ b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/test.ts @@ -0,0 +1,76 @@ +import type { Envelope, Event } from '@sentry/core'; +import { createServer } from 'node:http'; +import type { AddressInfo } from 'node:net'; +import { expect, it, onTestFinished } from 'vitest'; +import { createRunner } from '../../runner'; + +// Starts an ingest server that never answers envelope requests, so every send stays pending. The Workflow +// posts its result to `/result`, which resolves the returned promise with the posted body. +async function startSilentIngest(): Promise<{ url: string; result: Promise }> { + let resolveResult!: (body: unknown) => void; + const result = new Promise(resolve => (resolveResult = resolve)); + + const server = createServer((req, res) => { + if (req.url !== '/result') { + return; + } + let body = ''; + req.on('data', chunk => (body += chunk)); + req.on('end', () => { + res.end(); + resolveResult(JSON.parse(body)); + }); + }); + await new Promise(resolve => server.listen(0, resolve)); + onTestFinished(() => { + server.closeAllConnections(); + server.close(); + }); + + return { url: `http://localhost:${(server.address() as AddressInfo).port}`, result }; +} + +it.for([true, false])( + 'cacheClient: %s - aborts a send that is still pending when the flush times out', + async (cacheClient, { signal }) => { + const ingest = await startSilentIngest(); + + const runner = createRunner(__dirname) + .withServerUrl(ingest.url) + .withWranglerArgs('--var', 'SLOW_INGEST:true', '--var', `CACHE_CLIENT:${cacheClient}`) + .start(signal); + + const result = await runner.makeRequest('get', '/flush-with-timeout'); + expect(result).toEqual({ flushed: false, send: 'aborted' }); + }, +); + +// A cached client sends the step spans in eager drains the Workflow cannot wait for, so this runs with one +// client per invocation. The transport abort itself is covered for both modes by the test above. +it('cacheClient: false - the Workflow from #24482 aborts its pending send when the flush times out', async ({ + signal, +}) => { + const ingest = await startSilentIngest(); + + const runner = createRunner(__dirname) + .withServerUrl(ingest.url) + .withWranglerArgs('--var', 'SLOW_INGEST:true', '--var', 'CACHE_CLIENT:false', '--var', 'TRACING:true') + .start(signal); + + await runner.makeRequest('get', '/workflow/trigger'); + expect(await ingest.result).toEqual({ flushed: false, send: 'aborted' }); +}); + +it('cacheClient: false - delivers events while a user waitUntil task is still running', async ({ signal }) => { + const runner = createRunner(__dirname) + .withWranglerArgs('--var', 'CACHE_CLIENT:false') + .expect((envelope: Envelope) => { + const event = envelope[1]?.[0]?.[1] as Event; + expect(event.exception?.values?.[0]?.value).toBe('Captured on /pending-wait-until'); + }) + .unordered() + .start(signal); + + await runner.makeRequest('get', '/pending-wait-until'); + await runner.completed(); +}); diff --git a/dev-packages/cloudflare-integration-tests/suites/flush-timeout/vite.config.mts b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/vite.config.mts new file mode 100644 index 000000000000..005f4448f6cb --- /dev/null +++ b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/vite.config.mts @@ -0,0 +1,7 @@ +import { cloudflare } from '@cloudflare/vite-plugin'; +import { sentryCloudflareVitePlugin } from '@sentry/cloudflare/vite'; +import { defineConfig } from 'vite'; + +export default defineConfig({ + plugins: [cloudflare(), sentryCloudflareVitePlugin()], +}); diff --git a/dev-packages/cloudflare-integration-tests/suites/flush-timeout/wrangler.jsonc b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/wrangler.jsonc new file mode 100644 index 000000000000..a74a1ff6059c --- /dev/null +++ b/dev-packages/cloudflare-integration-tests/suites/flush-timeout/wrangler.jsonc @@ -0,0 +1,14 @@ +{ + "$schema": "../../node_modules/wrangler/config-schema.json", + "name": "cloudflare-flush-timeout", + "main": "index.ts", + "compatibility_date": "2025-06-17", + "compatibility_flags": ["nodejs_compat"], + "workflows": [ + { + "name": "issue-workflow", + "binding": "ISSUE_WORKFLOW", + "class_name": "IssueWorkflow", + }, + ], +} diff --git a/packages/cloudflare/src/client.ts b/packages/cloudflare/src/client.ts index a71445b03f92..1a968cfe66b2 100644 --- a/packages/cloudflare/src/client.ts +++ b/packages/cloudflare/src/client.ts @@ -107,19 +107,30 @@ export class CloudflareClient extends ServerRuntimeClient { /** * Flushes pending operations and ensures all data is processed. - * If a timeout is provided, the operation will be completed within the specified time limit. * - * It will wait for all pending spans to complete before flushing. + * Each phase waits at most `timeout`: the flush lock of a per-invocation client, pending spans, event + * processing and the transport drain. So a flush can take a small multiple of `timeout`, which stays well + * below Cloudflare's 30 second `waitUntil` limit for the timeouts the SDK uses. Sends still pending when + * the drain times out are aborted. * - * @param {number} [timeout] - Optional timeout in milliseconds to force the completion of the flush operation. + * @param {number} [timeout] - Maximum time in milliseconds for each phase of the flush. * @return {Promise} A promise that resolves to a boolean indicating whether the flush operation was successful. */ public async flush(timeout?: number): Promise { // Wait for user waitUntil-registered work to settle before draining, so events // captured in that work are still in the buffer. Without this the final flush // can drain (and the client be disposed) before background captures land. + // + // Only per-invocation clients (`cacheClient: false`) have a flush lock; remove this with them in v12. + // The wait is bounded by `timeout` because a user `waitUntil` task that outlives the invocation + // would otherwise keep the flush from draining until the runtime cancels it. if (this._flushLock) { - await this._flushLock.finalize(); + let timer: ReturnType | undefined; + await Promise.race([ + this._flushLock.finalize(), + ...(timeout ? [new Promise(resolve => (timer = setTimeout(resolve, timeout)))] : []), + ]); + clearTimeout(timer); } if (this._pendingSpans.size > 0 && this._spanCompletionPromise) { diff --git a/packages/cloudflare/src/transport.ts b/packages/cloudflare/src/transport.ts index 25d9e05572b9..b7f5a3181680 100644 --- a/packages/cloudflare/src/transport.ts +++ b/packages/cloudflare/src/transport.ts @@ -29,6 +29,12 @@ export class IsolatedPromiseBuffer { // If we ever remove it from the interface we should also remove it here. public $: Array>; + /** + * Abort signal of the drain that is starting its requests. It is set only while `drain()` runs the task + * producers, so a request reads the signal of the drain that sends it. + */ + public drainSignal: AbortSignal | undefined; + private _taskProducers: (() => PromiseLike)[]; private readonly _bufferSize: number; @@ -58,9 +64,21 @@ export class IsolatedPromiseBuffer { const oldTaskProducers = [...this._taskProducers]; this._taskProducers = []; + const drainController = new AbortController(); + this.drainSignal = drainController.signal; + let tasks: PromiseLike[]; + try { + tasks = oldTaskProducers.map(taskProducer => taskProducer()); + } finally { + this.drainSignal = undefined; + } + return new Promise(resolve => { const timer = setTimeout(() => { if (timeout && timeout > 0) { + // Requests still pending when the drain times out are aborted. Otherwise Cloudflare keeps them + // until it cancels the invocation's `waitUntil` work and logs a warning. + drainController.abort(); resolve(false); } }, timeout); @@ -68,8 +86,8 @@ export class IsolatedPromiseBuffer { // This cannot reject // eslint-disable-next-line @typescript-eslint/no-floating-promises Promise.all( - oldTaskProducers.map(taskProducer => - taskProducer().then(null, () => { + tasks.map(task => + task.then(null, () => { // catch all failed requests }), ), @@ -86,12 +104,20 @@ export class IsolatedPromiseBuffer { * Creates a Transport that uses the native fetch API to send events to Sentry. */ export function makeCloudflareTransport(options: CloudflareTransportOptions): Transport { + const buffer = new IsolatedPromiseBuffer(options.bufferSize); + function makeRequest(request: TransportRequest): PromiseLike { + const drainSignal = buffer.drainSignal; + const callerSignal = options.fetchOptions?.signal ?? undefined; + const signal = + drainSignal && callerSignal ? AbortSignal.any([drainSignal, callerSignal]) : (drainSignal ?? callerSignal); + const requestOptions: RequestInit = { body: request.body as BodyInit, method: 'POST', headers: options.headers, ...options.fetchOptions, + ...(signal ? { signal } : {}), }; return suppressTracing(() => { @@ -118,5 +144,5 @@ export function makeCloudflareTransport(options: CloudflareTransportOptions): Tr }); } - return createTransport(options, makeRequest, new IsolatedPromiseBuffer(options.bufferSize)); + return createTransport(options, makeRequest, buffer); } diff --git a/packages/cloudflare/test/client.test.ts b/packages/cloudflare/test/client.test.ts index 09bff574e479..2d7828be238d 100644 --- a/packages/cloudflare/test/client.test.ts +++ b/packages/cloudflare/test/client.test.ts @@ -268,6 +268,31 @@ describe('CloudflareClient', () => { await flushPromise; expect(privateClient._transport.flush).toHaveBeenCalledWith(1000); }); + + it('drains the transport when the flush lock does not settle within the timeout', async () => { + vi.useFakeTimers(); + try { + const client = new CloudflareClient({ + ...MOCK_CLIENT_OPTIONS, + flushLock: { ready: Promise.resolve(), finalize: () => new Promise(() => undefined) }, + }); + + const privateClient = client as unknown as { + _transport: { flush: ReturnType }; + }; + + void client.flush(1000); + + await vi.advanceTimersByTimeAsync(999); + expect(privateClient._transport.flush).not.toHaveBeenCalled(); + + // The lock wait ends at 1000 ms; the client processing check after it also runs on timers. + await vi.advanceTimersByTimeAsync(100); + expect(privateClient._transport.flush).toHaveBeenCalledWith(1000); + } finally { + vi.useRealTimers(); + } + }); }); describe('span lifecycle tracking', () => { diff --git a/packages/cloudflare/test/request.test.ts b/packages/cloudflare/test/request.test.ts index 052147f84bc4..baf9282f9e15 100644 --- a/packages/cloudflare/test/request.test.ts +++ b/packages/cloudflare/test/request.test.ts @@ -1051,6 +1051,7 @@ describe('Durable Object (DO) context', () => { // Teardown is registered via waitUntil on error too expect(waitUntilSpy).toHaveBeenCalled(); + await Promise.all(waitUntilSpy.mock.calls.map(([promise]) => promise)); // And flush runs as part of that teardown expect(flushSpy).toHaveBeenCalled(); @@ -1072,6 +1073,7 @@ describe('Durable Object (DO) context', () => { ); expect(waitUntilSpy).toHaveBeenCalled(); + await Promise.all(waitUntilSpy.mock.calls.map(([promise]) => promise)); expect(flushSpy).toHaveBeenCalled(); flushSpy.mockRestore(); @@ -1092,6 +1094,7 @@ describe('Durable Object (DO) context', () => { ); expect(waitUntilSpy).toHaveBeenCalled(); + await Promise.all(waitUntilSpy.mock.calls.map(([promise]) => promise)); expect(flushSpy).toHaveBeenCalled(); flushSpy.mockRestore(); diff --git a/packages/cloudflare/test/transport.test.ts b/packages/cloudflare/test/transport.test.ts index fdb9fbc5e30f..1750176f5591 100644 --- a/packages/cloudflare/test/transport.test.ts +++ b/packages/cloudflare/test/transport.test.ts @@ -52,6 +52,7 @@ describe('Edge Transport', () => { expect(mockFetch).toHaveBeenLastCalledWith(DEFAULT_EDGE_TRANSPORT_OPTIONS.url, { body: serializeEnvelope(ERROR_ENVELOPE), method: 'POST', + signal: expect.any(AbortSignal), }); }); @@ -104,6 +105,7 @@ describe('Edge Transport', () => { body: serializeEnvelope(ERROR_ENVELOPE), method: 'POST', ...REQUEST_OPTIONS, + signal: expect.any(AbortSignal), }); }); @@ -249,4 +251,107 @@ describe('IsolatedPromiseBuffer', () => { await transport.flush(); expect(customFetch).toHaveBeenCalledTimes(1); }); + + it('aborts a request that is still pending when its drain times out', async () => { + vi.useFakeTimers(); + try { + let signal: AbortSignal | undefined; + const customFetch = vi.fn( + (_input: RequestInfo | URL, init?: RequestInit): Promise => + new Promise((_resolve, reject) => { + signal = init?.signal ?? undefined; + signal?.addEventListener('abort', () => reject(new DOMException('Aborted', 'AbortError')), { once: true }); + }), + ); + const transport = makeCloudflareTransport({ ...DEFAULT_EDGE_TRANSPORT_OPTIONS, fetch: customFetch }); + + await transport.send(ERROR_ENVELOPE); + const flush = transport.flush(1000); + + await vi.advanceTimersByTimeAsync(999); + expect(signal?.aborted).toBe(false); + + await vi.advanceTimersByTimeAsync(1); + await expect(flush).resolves.toBe(false); + expect(signal?.aborted).toBe(true); + } finally { + vi.useRealTimers(); + } + }); + + it('does not abort requests of a drain without a timeout', async () => { + vi.useFakeTimers(); + try { + let signal: AbortSignal | undefined; + const customFetch = vi.fn( + (_input: RequestInfo | URL, init?: RequestInit): Promise => + new Promise(() => { + signal = init?.signal ?? undefined; + }), + ); + const transport = makeCloudflareTransport({ ...DEFAULT_EDGE_TRANSPORT_OPTIONS, fetch: customFetch }); + + await transport.send(ERROR_ENVELOPE); + void transport.flush(); + + await vi.advanceTimersByTimeAsync(60_000); + expect(signal?.aborted).toBe(false); + } finally { + vi.useRealTimers(); + } + }); + + it('preserves a caller-provided abort signal', async () => { + let signal: AbortSignal | undefined; + const callerController = new AbortController(); + const customFetch = vi.fn( + (_input: RequestInfo | URL, init?: RequestInit): Promise => + new Promise((_resolve, reject) => { + signal = init?.signal ?? undefined; + signal?.addEventListener('abort', () => reject(new DOMException('Aborted', 'AbortError')), { once: true }); + }), + ); + const transport = makeCloudflareTransport({ + ...DEFAULT_EDGE_TRANSPORT_OPTIONS, + fetch: customFetch, + fetchOptions: { signal: callerController.signal }, + }); + + await transport.send(ERROR_ENVELOPE); + const flush = transport.flush(); + callerController.abort(); + + await expect(flush).resolves.toBe(true); + expect(signal?.aborted).toBe(true); + }); + + it('does not abort requests belonging to another drain', async () => { + vi.useFakeTimers(); + try { + const signals: AbortSignal[] = []; + const customFetch = vi.fn( + (_input: RequestInfo | URL, init?: RequestInit): Promise => + new Promise((_resolve, reject) => { + const signal = init?.signal as AbortSignal; + signals.push(signal); + signal.addEventListener('abort', () => reject(new DOMException('Aborted', 'AbortError')), { once: true }); + }), + ); + const transport = makeCloudflareTransport({ ...DEFAULT_EDGE_TRANSPORT_OPTIONS, fetch: customFetch }); + + await transport.send(ERROR_ENVELOPE); + void transport.flush(2000); + await transport.send(ERROR_ENVELOPE); + void transport.flush(1000); + + await vi.advanceTimersByTimeAsync(1000); + expect(signals[0]?.aborted).toBe(false); + expect(signals[1]?.aborted).toBe(true); + + await vi.advanceTimersByTimeAsync(1000); + expect(signals[0]?.aborted).toBe(true); + } finally { + vi.useRealTimers(); + } + }); });