Skip to content

Commit 7ad0606

Browse files
matthewbjonescodex
andcommitted
fix(cloudflare): Enforce flush timeout across workflow lifecycle (#24482)
Cloudflare's flush deadline could expire while transport fetches remained active, Workflow steps could inherit the previous step's flush point and schedule redundant eager drains, and flush-lock finalization could wait outside the configured timeout. Abort requests when a transport drain expires, reset the flush point at each Workflow step boundary, and share one deadline across flush-lock, pending-span, and transport phases. Add focused regression coverage for all three paths and update timer-based tests to await the user task before asserting teardown. Fixes #24482 Co-authored-by: OpenAI Codex <codex@openai.com>
1 parent c94fd11 commit 7ad0606

12 files changed

Lines changed: 233 additions & 20 deletions

‎packages/cloudflare/src/client.ts‎

Lines changed: 39 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,22 @@ import type { makeFlushLock } from './flush';
1515
import type { CloudflareTransportOptions } from './transport';
1616
import { getInvocationState, getInvocationWaitUntil } from './utils/invocationContext';
1717

18+
async function waitForPromise(promise: PromiseLike<unknown>, timeout: number): Promise<boolean> {
19+
let timer: ReturnType<typeof setTimeout> | undefined;
20+
try {
21+
return await Promise.race([
22+
Promise.resolve(promise).then(() => true),
23+
new Promise<false>(resolve => {
24+
timer = setTimeout(() => resolve(false), timeout);
25+
}),
26+
]);
27+
} finally {
28+
if (timer) {
29+
clearTimeout(timer);
30+
}
31+
}
32+
}
33+
1834
/**
1935
* The Sentry Cloudflare SDK Client.
2036
*
@@ -115,18 +131,34 @@ export class CloudflareClient extends ServerRuntimeClient {
115131
* @return {Promise<boolean>} A promise that resolves to a boolean indicating whether the flush operation was successful.
116132
*/
117133
public async flush(timeout?: number): Promise<boolean> {
134+
const deadline = timeout && timeout > 0 ? Date.now() + timeout : undefined;
135+
const remainingTimeout = (): number | undefined =>
136+
deadline === undefined ? timeout : Math.max(0, deadline - Date.now());
137+
118138
// Wait for user waitUntil-registered work to settle before draining, so events
119139
// captured in that work are still in the buffer. Without this the final flush
120140
// can drain (and the client be disposed) before background captures land.
121141
if (this._flushLock) {
122-
await this._flushLock.finalize();
142+
const lockTimeout = remainingTimeout();
143+
if (lockTimeout && lockTimeout > 0) {
144+
if (!(await waitForPromise(this._flushLock.finalize(), lockTimeout))) {
145+
return false;
146+
}
147+
} else if (deadline !== undefined) {
148+
return false;
149+
} else {
150+
await this._flushLock.finalize();
151+
}
123152
}
124153

125154
if (this._pendingSpans.size > 0 && this._spanCompletionPromise) {
126155
DEBUG_BUILD &&
127156
debug.log('[CloudflareClient] Waiting for', this._pendingSpans.size, 'pending spans to complete...');
128157

129-
const timeoutMs = timeout ?? 5000;
158+
const timeoutMs = remainingTimeout() ?? 5000;
159+
if (deadline !== undefined && timeoutMs <= 0) {
160+
return false;
161+
}
130162
const spanCompletionRace = Promise.race([
131163
this._spanCompletionPromise,
132164
new Promise(resolve =>
@@ -146,7 +178,11 @@ export class CloudflareClient extends ServerRuntimeClient {
146178
// of them would also start an eager drain and a `waitUntil` registration.
147179
this._inBoundaryFlush = true;
148180
try {
149-
return await super.flush(timeout);
181+
const transportTimeout = remainingTimeout();
182+
if (deadline !== undefined && (!transportTimeout || transportTimeout <= 0)) {
183+
return false;
184+
}
185+
return await super.flush(transportTimeout);
150186
} finally {
151187
this._inBoundaryFlush = false;
152188
}

‎packages/cloudflare/src/transport.ts‎

Lines changed: 54 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,9 @@ export interface CloudflareTransportOptions extends BaseTransportOptions {
1616
*/
1717
const DEFAULT_TRANSPORT_BUFFER_SIZE = 256;
1818

19+
type TaskProducer = () => PromiseLike<TransportMakeRequestResponse>;
20+
type RunTask = (taskProducer: TaskProducer, signal: AbortSignal) => PromiseLike<TransportMakeRequestResponse>;
21+
1922
/**
2023
* This is a modified promise buffer that collects tasks until drain is called.
2124
* We need this in the edge runtime because edge function invocations may not share I/O objects, like fetch requests
@@ -29,20 +32,23 @@ export class IsolatedPromiseBuffer {
2932
// If we ever remove it from the interface we should also remove it here.
3033
public $: Array<PromiseLike<TransportMakeRequestResponse>>;
3134

32-
private _taskProducers: (() => PromiseLike<TransportMakeRequestResponse>)[];
35+
private _taskProducers: TaskProducer[];
3336

3437
private readonly _bufferSize: number;
3538

36-
public constructor(_bufferSize = DEFAULT_TRANSPORT_BUFFER_SIZE) {
39+
private readonly _runTask: RunTask;
40+
41+
public constructor(_bufferSize = DEFAULT_TRANSPORT_BUFFER_SIZE, _runTask: RunTask = taskProducer => taskProducer()) {
3742
this.$ = [];
3843
this._taskProducers = [];
3944
this._bufferSize = _bufferSize;
45+
this._runTask = _runTask;
4046
}
4147

4248
/**
4349
* @inheritdoc
4450
*/
45-
public add(taskProducer: () => PromiseLike<TransportMakeRequestResponse>): PromiseLike<TransportMakeRequestResponse> {
51+
public add(taskProducer: TaskProducer): PromiseLike<TransportMakeRequestResponse> {
4652
if (this._taskProducers.length >= this._bufferSize) {
4753
return Promise.reject(SENTRY_BUFFER_FULL_ERROR);
4854
}
@@ -57,19 +63,22 @@ export class IsolatedPromiseBuffer {
5763
public drain(timeout?: number): PromiseLike<boolean> {
5864
const oldTaskProducers = [...this._taskProducers];
5965
this._taskProducers = [];
66+
const drainController = new AbortController();
67+
const tasks = oldTaskProducers.map(taskProducer => this._runTask(taskProducer, drainController.signal));
6068

6169
return new Promise(resolve => {
6270
const timer = setTimeout(() => {
6371
if (timeout && timeout > 0) {
72+
drainController.abort();
6473
resolve(false);
6574
}
6675
}, timeout);
6776

6877
// This cannot reject
6978
// eslint-disable-next-line @typescript-eslint/no-floating-promises
7079
Promise.all(
71-
oldTaskProducers.map(taskProducer =>
72-
taskProducer().then(null, () => {
80+
tasks.map(task =>
81+
task.then(null, () => {
7382
// catch all failed requests
7483
}),
7584
),
@@ -86,15 +95,36 @@ export class IsolatedPromiseBuffer {
8695
* Creates a Transport that uses the native fetch API to send events to Sentry.
8796
*/
8897
export function makeCloudflareTransport(options: CloudflareTransportOptions): Transport {
98+
let activeDrainSignal: AbortSignal | undefined;
99+
89100
function makeRequest(request: TransportRequest): PromiseLike<TransportMakeRequestResponse> {
101+
const controller = new AbortController();
102+
const callerSignal = options.fetchOptions?.signal;
103+
const drainSignal = activeDrainSignal;
104+
const abortFromCallerSignal = (): void => controller.abort();
105+
const abortFromDrainSignal = (): void => controller.abort();
106+
107+
if (callerSignal?.aborted) {
108+
controller.abort();
109+
} else {
110+
callerSignal?.addEventListener('abort', abortFromCallerSignal, { once: true });
111+
}
112+
113+
if (drainSignal?.aborted) {
114+
controller.abort();
115+
} else {
116+
drainSignal?.addEventListener('abort', abortFromDrainSignal, { once: true });
117+
}
118+
90119
const requestOptions: RequestInit = {
91120
body: request.body as BodyInit,
92121
method: 'POST',
93122
headers: options.headers,
94123
...options.fetchOptions,
124+
signal: controller.signal,
95125
};
96126

97-
return suppressTracing(() => {
127+
const requestPromise = suppressTracing(() => {
98128
return (options.fetch ?? fetch)(options.url, requestOptions).then(async response => {
99129
// Consume the response body to satisfy Cloudflare Workers' fetch requirements.
100130
// The runtime requires all fetch response bodies to be read or explicitly canceled
@@ -116,7 +146,24 @@ export function makeCloudflareTransport(options: CloudflareTransportOptions): Tr
116146
};
117147
});
118148
});
149+
150+
return Promise.resolve(requestPromise).finally(() => {
151+
callerSignal?.removeEventListener('abort', abortFromCallerSignal);
152+
drainSignal?.removeEventListener('abort', abortFromDrainSignal);
153+
});
154+
}
155+
156+
function runTaskWithinDrain(
157+
taskProducer: TaskProducer,
158+
signal: AbortSignal,
159+
): PromiseLike<TransportMakeRequestResponse> {
160+
activeDrainSignal = signal;
161+
try {
162+
return taskProducer();
163+
} finally {
164+
activeDrainSignal = undefined;
165+
}
119166
}
120167

121-
return createTransport(options, makeRequest, new IsolatedPromiseBuffer(options.bufferSize));
168+
return createTransport(options, makeRequest, new IsolatedPromiseBuffer(options.bufferSize, runTaskWithinDrain));
122169
}

‎packages/cloudflare/src/workflows.ts‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ import { instrumentEnv } from './instrumentations/worker/instrumentEnv';
3030
import { addCloudResourceContext } from './scope-utils';
3131
import { init } from './sdk';
3232
import { instrumentContext } from './utils/instrumentContext';
33+
import { getInvocationState } from './utils/invocationContext';
3334
import type { DefaultEnv, ResolveEnv, StrictCloudflareOptions } from './types';
3435
import { withInvocationIsolationScope } from './utils/invocationScope';
3536

@@ -124,6 +125,14 @@ class WrappedWorkflowStep implements WorkflowStep {
124125
// run's isolation scope (and with it the invocation state that ties eager sends
125126
// to this invocation's `waitUntil`) has to be restored explicitly.
126127
return withIsolationScope(this._isolationScope, () => {
128+
// Each Workflow step is its own RPC invocation with its own boundary flush.
129+
// The isolation scope is shared across steps, so clear the previous step's
130+
// flush point before capturing anything for this one.
131+
const invocationState = getInvocationState();
132+
if (invocationState) {
133+
invocationState.flushPointReached = false;
134+
}
135+
127136
const stepResult = startSpan(
128137
{
129138
name,

‎packages/cloudflare/test/client.test.ts‎

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -266,7 +266,29 @@ describe('CloudflareClient', () => {
266266

267267
releaseLock();
268268
await flushPromise;
269-
expect(privateClient._transport.flush).toHaveBeenCalledWith(1000);
269+
expect(privateClient._transport.flush).toHaveBeenCalledWith(expect.any(Number));
270+
const transportTimeout = privateClient._transport.flush.mock.calls[0]?.[0];
271+
expect(transportTimeout).toBeGreaterThan(0);
272+
expect(transportTimeout).toBeLessThanOrEqual(1000);
273+
});
274+
275+
it('includes the flush lock in the timeout', async () => {
276+
const finalize = vi.fn(() => new Promise<void>(() => undefined));
277+
const client = new CloudflareClient({
278+
...MOCK_CLIENT_OPTIONS,
279+
flushLock: { ready: Promise.resolve(), finalize },
280+
});
281+
282+
const privateClient = client as unknown as {
283+
_transport: { flush: ReturnType<typeof vi.fn> };
284+
};
285+
const result = await Promise.race([
286+
client.flush(10),
287+
new Promise<'did-not-settle'>(resolve => setTimeout(() => resolve('did-not-settle'), 30)),
288+
]);
289+
290+
expect(result).toBe(false);
291+
expect(privateClient._transport.flush).not.toHaveBeenCalled();
270292
});
271293
});
272294

‎packages/cloudflare/test/instrumentations/worker/instrumentEmail.test.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -302,7 +302,7 @@ describe('instrumentEmail', () => {
302302
} as unknown as ExecutionContext);
303303
expect(flush).not.toBeCalled();
304304
expect(waitUntil).toBeCalled();
305-
vi.advanceTimersToNextTimer().runAllTimers();
305+
await vi.advanceTimersToNextTimerAsync();
306306
await Promise.all(waits);
307307
expect(flush).toHaveBeenCalledOnce();
308308
});

‎packages/cloudflare/test/instrumentations/worker/instrumentFetch.test.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -172,7 +172,7 @@ describe('instrumentFetch', () => {
172172
.then(response => response.text());
173173
expect(flush).not.toBeCalled();
174174
expect(waitUntil).toBeCalled();
175-
vi.advanceTimersToNextTimer().runAllTimers();
175+
await vi.advanceTimersToNextTimerAsync();
176176
await Promise.all(waits);
177177
expect(flush).toHaveBeenCalledOnce();
178178
});

‎packages/cloudflare/test/instrumentations/worker/instrumentQueue.test.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -318,7 +318,7 @@ describe('instrumentQueue', () => {
318318
} as unknown as ExecutionContext);
319319
expect(flush).not.toBeCalled();
320320
expect(waitUntil).toBeCalled();
321-
vi.advanceTimersToNextTimer().runAllTimers();
321+
await vi.advanceTimersToNextTimerAsync();
322322
await Promise.all(waits);
323323
expect(flush).toHaveBeenCalledOnce();
324324
});

‎packages/cloudflare/test/instrumentations/worker/instrumentScheduled.test.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -323,7 +323,7 @@ describe('instrumentScheduled', () => {
323323
} as unknown as ExecutionContext);
324324
expect(flush).not.toBeCalled();
325325
expect(waitUntil).toBeCalled();
326-
vi.advanceTimersToNextTimer().runAllTimers();
326+
await vi.advanceTimersToNextTimerAsync();
327327
await Promise.all(waits);
328328
expect(flush).toHaveBeenCalledOnce();
329329
});

‎packages/cloudflare/test/instrumentations/worker/instrumentTail.test.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -270,7 +270,7 @@ describe('instrumentTail', () => {
270270
} as unknown as ExecutionContext);
271271
expect(flush).not.toBeCalled();
272272
expect(waitUntil).toBeCalled();
273-
vi.advanceTimersToNextTimer().runAllTimers();
273+
await vi.advanceTimersToNextTimerAsync();
274274
await Promise.all(waits);
275275
expect(flush).toHaveBeenCalledOnce();
276276
});

‎packages/cloudflare/test/request.test.ts‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -150,7 +150,7 @@ describe('withSentry', () => {
150150
return new Response('test');
151151
}).then(response => response.text());
152152
expect(waitUntil).toBeCalled();
153-
vi.advanceTimersToNextTimer().runAllTimers();
153+
await vi.advanceTimersToNextTimerAsync();
154154
await Promise.all(waits);
155155

156156
const after = flushSpy.mock.calls.length;
@@ -1051,6 +1051,7 @@ describe('Durable Object (DO) context', () => {
10511051

10521052
// Teardown is registered via waitUntil on error too
10531053
expect(waitUntilSpy).toHaveBeenCalled();
1054+
await Promise.all(waitUntilSpy.mock.calls.map(([promise]) => promise));
10541055
// And flush runs as part of that teardown
10551056
expect(flushSpy).toHaveBeenCalled();
10561057

@@ -1072,6 +1073,7 @@ describe('Durable Object (DO) context', () => {
10721073
);
10731074

10741075
expect(waitUntilSpy).toHaveBeenCalled();
1076+
await Promise.all(waitUntilSpy.mock.calls.map(([promise]) => promise));
10751077
expect(flushSpy).toHaveBeenCalled();
10761078

10771079
flushSpy.mockRestore();
@@ -1092,6 +1094,7 @@ describe('Durable Object (DO) context', () => {
10921094
);
10931095

10941096
expect(waitUntilSpy).toHaveBeenCalled();
1097+
await Promise.all(waitUntilSpy.mock.calls.map(([promise]) => promise));
10951098
expect(flushSpy).toHaveBeenCalled();
10961099

10971100
flushSpy.mockRestore();

0 commit comments

Comments
 (0)