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
Original file line number Diff line number Diff line change
@@ -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<Env> {
async run(_event: WorkflowEvent<unknown>, step: WorkflowStep): Promise<void> {
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<Env>;
Original file line number Diff line number Diff line change
@@ -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);
},
},
}));
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
// Whether the transport aborted an envelope fetch since `aborted` was last reset.
export const lastSend: { aborted: boolean } = { aborted: false };
Original file line number Diff line number Diff line change
@@ -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<unknown> }> {
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<void>(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();
});
Original file line number Diff line number Diff line change
@@ -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()],
});
Original file line number Diff line number Diff line change
@@ -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",
},
],
}
19 changes: 15 additions & 4 deletions packages/cloudflare/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<boolean>} A promise that resolves to a boolean indicating whether the flush operation was successful.
*/
public async flush(timeout?: number): Promise<boolean> {
// 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<typeof setTimeout> | undefined;
await Promise.race([
this._flushLock.finalize(),
...(timeout ? [new Promise<void>(resolve => (timer = setTimeout(resolve, timeout)))] : []),
]);
clearTimeout(timer);
Comment thread
JPeer264 marked this conversation as resolved.
}
Comment thread
JPeer264 marked this conversation as resolved.

if (this._pendingSpans.size > 0 && this._spanCompletionPromise) {
Expand Down
32 changes: 29 additions & 3 deletions packages/cloudflare/src/transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,12 @@ export class IsolatedPromiseBuffer {
// If we ever remove it from the interface we should also remove it here.
public $: Array<PromiseLike<TransportMakeRequestResponse>>;

/**
* 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<TransportMakeRequestResponse>)[];

private readonly _bufferSize: number;
Expand Down Expand Up @@ -58,18 +64,30 @@ export class IsolatedPromiseBuffer {
const oldTaskProducers = [...this._taskProducers];
this._taskProducers = [];

const drainController = new AbortController();
this.drainSignal = drainController.signal;
let tasks: PromiseLike<TransportMakeRequestResponse>[];
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);

// 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
}),
),
Expand All @@ -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<TransportMakeRequestResponse> {
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(() => {
Expand All @@ -118,5 +144,5 @@ export function makeCloudflareTransport(options: CloudflareTransportOptions): Tr
});
}

return createTransport(options, makeRequest, new IsolatedPromiseBuffer(options.bufferSize));
return createTransport(options, makeRequest, buffer);
}
25 changes: 25 additions & 0 deletions packages/cloudflare/test/client.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>(() => undefined) },
});

const privateClient = client as unknown as {
_transport: { flush: ReturnType<typeof vi.fn> };
};

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', () => {
Expand Down
3 changes: 3 additions & 0 deletions packages/cloudflare/test/request.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand All @@ -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();
Expand All @@ -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();
Expand Down
Loading
Loading