Skip to content
Merged
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
2 changes: 1 addition & 1 deletion .size-limit.js
Original file line number Diff line number Diff line change
Expand Up @@ -480,7 +480,7 @@ module.exports = [
ignore: [...builtinModules, ...nodePrefixedBuiltinModules],
gzip: false,
brotli: false,
limit: '492 KiB',
limit: '495 KiB',
disablePlugins: ['@size-limit/webpack'],
webpack: false,
modifyEsbuildConfig: function (config) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
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. Each step flushes its span to an
// ingest that never answers. The run reports to SERVER_URL once the SDK has aborted one of those pending sends.
export class IssueWorkflow extends WorkflowEntrypoint<Env> {
async run(_event: WorkflowEvent<unknown>, step: WorkflowStep): Promise<void> {
const stepSendAborted = new Promise<void>(resolve => (lastSend.onAbort = resolve));

for (let index = 0; index < 100; index++) {
await step.do(`step-${index}`, async () => index);
}

await stepSendAborted;
await fetch(`${this.env.SERVER_URL}/result`, { method: 'POST', body: JSON.stringify({ send: 'aborted' }) });
}
}

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;
// "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,
tracesSampleRate: env.TRACING === 'true' ? 1 : undefined,
transportOptions: {
fetch: (input, init) => {
init?.signal?.addEventListener('abort', () => {
lastSend.aborted = true;
lastSend.onAbort?.();
});
return fetch(input, init);
},
},
}));
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
// `aborted`: whether the transport aborted an envelope fetch since it was last reset.
// `onAbort`: called each time the transport aborts an envelope fetch.
export const lastSend: { aborted: boolean; onAbort?: () => void } = { aborted: false };
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
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('aborts a send that is still pending when the flush times out', async ({ signal }) => {
const ingest = await startSilentIngest();

const runner = createRunner(__dirname)
.withServerUrl(ingest.url)
.withWranglerArgs('--var', 'SLOW_INGEST:true')
.start(signal);

const result = await runner.makeRequest('get', '/flush-with-timeout');
expect(result).toEqual({ flushed: false, send: 'aborted' });
});

it('the Workflow from #24482 aborts the pending send of a step when its flush times out', async ({ signal }) => {
const ingest = await startSilentIngest();

const runner = createRunner(__dirname)
.withServerUrl(ingest.url)
.withWranglerArgs('--var', 'SLOW_INGEST:true', '--var', 'TRACING:true')
.start(signal);

await runner.makeRequest('get', '/workflow/trigger');
expect(await ingest.result).toEqual({ send: 'aborted' });
});

it('delivers events while a user waitUntil task is still running', async ({ signal }) => {
const runner = createRunner(__dirname)
.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,14 @@
import { cloudflare } from '@cloudflare/vite-plugin';
import { sentryCloudflareVitePlugin } from '@sentry/cloudflare/vite';
import { defineConfig } from 'vite';

export default defineConfig({
plugins: [
cloudflare(),
sentryCloudflareVitePlugin({
_experimental: {
autoInstrumentation: true,
},
}),
],
});
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",
},
],
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
import { WorkflowEntrypoint } from 'cloudflare:workers';
import type { WorkflowEvent, WorkflowStep } from 'cloudflare:workers';

interface Env {
SENTRY_DSN: string;
SLEEP_WORKFLOW: Workflow;
}

export class SleepWorkflow extends WorkflowEntrypoint<Env> {
async run(_event: WorkflowEvent<unknown>, step: WorkflowStep): Promise<void> {
await step.do('before-sleep', async () => 'done');
await step.sleep('pause', '1 hour');
await step.do('after-sleep', async () => 'done');
}
}

export default {
async fetch(request: Request, env: Env): Promise<Response> {
const url = new URL(request.url);

if (url.pathname === '/workflow/trigger') {
const instance = await env.SLEEP_WORKFLOW.create();
return Response.json({ id: instance.id });
}

return new Response('Not found', { status: 404 });
},
} satisfies ExportedHandler<Env>;
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
import { defineCloudflareOptions } from '@sentry/cloudflare';

export default defineCloudflareOptions((env: { SENTRY_DSN: string }) => ({
dsn: env.SENTRY_DSN,
tracesSampleRate: 1.0,
traceLifecycle: 'stream',
}));
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
import type { SerializedStreamedSpanContainer } from '@sentry/core';
import { expect, it } from 'vitest';
import { createRunner } from '../../../runner';

it('sends the span of a step before the Workflow goes to sleep', async ({ signal }) => {
const runner = createRunner(__dirname)
.expect(envelope => {
const container = envelope[1].find(item => item[0].type === 'span')?.[1] as SerializedStreamedSpanContainer;

expect(container.items).toHaveLength(1);
expect(container.items[0]!.name).toBe('before-sleep');
expect(envelope[0].trace).toEqual({
environment: 'production',
public_key: 'public',
trace_id: container.items[0]!.trace_id,
transaction: 'before-sleep',
sampled: 'true',
sample_rand: expect.any(String),
sample_rate: '1',
});
})
.unordered()
.start(signal);

await runner.makeRequest('get', '/workflow/trigger');
await runner.completed();
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
import { cloudflare } from '@cloudflare/vite-plugin';
import { sentryCloudflareVitePlugin } from '@sentry/cloudflare/vite';
import { defineConfig } from 'vite';

export default defineConfig({
plugins: [
cloudflare(),
sentryCloudflareVitePlugin({
_experimental: {
autoInstrumentation: true,
},
}),
],
});
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
{
"$schema": "../../../node_modules/wrangler/config-schema.json",
"name": "cloudflare-workflow-step-flush",
"main": "index.ts",
"compatibility_date": "2025-06-17",
"compatibility_flags": ["nodejs_compat"],
"workflows": [
{
"name": "sleep-workflow",
"binding": "SLEEP_WORKFLOW",
"class_name": "SleepWorkflow",
},
],
}
17 changes: 13 additions & 4 deletions packages/cloudflare/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -85,16 +85,25 @@ 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> {
// 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);
}

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 @@ -23,6 +23,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 @@ -52,18 +58,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 @@ -80,12 +98,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,
method: 'POST',
headers: options.headers,
...options.fetchOptions,
...(signal ? { signal } : {}),
};

return suppressTracing(() => {
Expand All @@ -112,5 +138,5 @@ export function makeCloudflareTransport(options: CloudflareTransportOptions): Tr
});
}

return createTransport(options, makeRequest, new IsolatedPromiseBuffer(options.bufferSize));
return createTransport(options, makeRequest, buffer);
}
Loading
Loading