From 564c7f3b776aaa88dcaa1ba74ce4bb6bf3b99b0c Mon Sep 17 00:00:00 2001 From: JPeer264 Date: Tue, 4 Aug 2026 15:04:16 +0200 Subject: [PATCH] fix(cloudflare): Get original waituntil in workflows --- packages/cloudflare/src/flush.ts | 2 +- packages/cloudflare/src/request.ts | 2 +- packages/cloudflare/src/workflows.ts | 4 +- packages/cloudflare/test/flush.test.ts | 6 +- packages/cloudflare/test/workflow.test.ts | 129 ++++++++++++++++++++++ 5 files changed, 136 insertions(+), 7 deletions(-) diff --git a/packages/cloudflare/src/flush.ts b/packages/cloudflare/src/flush.ts index 77911cd8752c..fe86e21dbd62 100644 --- a/packages/cloudflare/src/flush.ts +++ b/packages/cloudflare/src/flush.ts @@ -35,7 +35,7 @@ const flushLockRegistries = new WeakMap { - const waitUntil = context.waitUntil.bind(context); + const waitUntil = getOriginalWaitUntil(context).bind(context); const client = init({ ...options, ctx: context, enableDedupe: false }); isolationScope.setClient(client); diff --git a/packages/cloudflare/test/flush.test.ts b/packages/cloudflare/test/flush.test.ts index 49ce15dc5153..bcef56a8c101 100644 --- a/packages/cloudflare/test/flush.test.ts +++ b/packages/cloudflare/test/flush.test.ts @@ -165,7 +165,7 @@ describe('getOriginalWaitUntil', () => { expect(result).not.toBe(context.waitUntil); expect(result).toBeDefined(); - result!(Promise.resolve()); + result(Promise.resolve()); expect(originalWaitUntil).toHaveBeenCalled(); }); @@ -183,7 +183,7 @@ describe('getOriginalWaitUntil', () => { const result = getOriginalWaitUntil(context); expect(result).not.toBe(context.waitUntil); - result!(Promise.resolve()); + result(Promise.resolve()); expect(originalWaitUntil).toHaveBeenCalled(); }); @@ -207,7 +207,7 @@ describe('getOriginalWaitUntil', () => { } as unknown as Client; const originalWaitUntil = getOriginalWaitUntil(context); - originalWaitUntil!.call(context, flushAndDispose(mockClient)); + originalWaitUntil.call(context, flushAndDispose(mockClient)); await vi.waitFor(() => Promise.all(waitUntilPromises)); expect(mockClient.flush).toHaveBeenCalled(); diff --git a/packages/cloudflare/test/workflow.test.ts b/packages/cloudflare/test/workflow.test.ts index 5ea96ba8449d..2578c1e0343d 100644 --- a/packages/cloudflare/test/workflow.test.ts +++ b/packages/cloudflare/test/workflow.test.ts @@ -209,6 +209,135 @@ describe.skipIf(NODE_MAJOR_VERSION < 20)('workflows', () => { await expect(drainWaitUntilLikeCloudflareVitestPool(waitUntilPromises)).resolves.toBeUndefined(); }); + test('teardown does not deadlock when a workflow instance is reused across runs', async () => { + const waitUntilPromises: Promise[] = []; + const context: ExecutionContext = { + waitUntil: vi.fn((promise: Promise) => { + waitUntilPromises.push(promise); + }), + passThroughOnException: vi.fn(), + props: {}, + }; + + let runCount = 0; + let releaseAppWork: () => void = () => undefined; + + class ReusedWorkflow { + public constructor(private _ctx: ExecutionContext) {} + + public async run(_event: Readonly>, step: WorkflowStep): Promise { + runCount += 1; + await step.do('reused step', async () => { + if (runCount === 2) { + this._ctx.waitUntil( + new Promise(resolve => { + releaseAppWork = resolve; + }), + ); + } + }); + } + } + + const TestWorkflowInstrumented = instrumentWorkflowWithSentry(getSentryOptions, ReusedWorkflow as any); + // Cloudflare reuses a Workflow instance across runs, so the context + // captured at construction is instrumented by the first run's init() + const workflow = new TestWorkflowInstrumented(context, {}) as ReusedWorkflow; + const event = { payload: {}, timestamp: new Date(), instanceId: INSTANCE_ID }; + + await workflow.run(event, mockStep); + await drainWaitUntilLikeCloudflareVitestPool(waitUntilPromises); + + await workflow.run(event, mockStep); + + releaseAppWork(); + + // Both the application work and the teardown promise must settle + await expect(drainWaitUntilLikeCloudflareVitestPool(waitUntilPromises)).resolves.toBeUndefined(); + }); + + test('step errors are still captured when a workflow instance is reused across runs', async () => { + const waitUntilPromises: Promise[] = []; + const context: ExecutionContext = { + waitUntil: vi.fn((promise: Promise) => { + waitUntilPromises.push(promise); + }), + passThroughOnException: vi.fn(), + props: {}, + }; + + let runCount = 0; + + class ReusedErrorWorkflow { + public constructor(private _ctx: ExecutionContext) {} + + public async run(_event: Readonly>, step: WorkflowStep): Promise { + runCount += 1; + await step.do('flaky step', async () => { + if (runCount === 2) { + throw new Error('second run error'); + } + }); + } + } + + // Fails the step through every retry without backoff, so the error is + // captured on the final attempt and surfaces from run() + const alwaysFailStep: WorkflowStep = { + do: vi + .fn() + .mockImplementation( + async ( + _name: string, + configOrCallback: WorkflowStepConfig | ((...args: unknown[]) => Promise), + maybeCallback?: (...args: unknown[]) => Promise, + ) => { + const retryLimit = 2; + const callback = (typeof configOrCallback === 'function' ? configOrCallback : maybeCallback)!; + let lastError: unknown; + for (let attempt = 1; attempt <= retryLimit + 1; attempt++) { + try { + return await callback({ attempt, config: { retries: { limit: retryLimit }, timeout: 60000 } }); + } catch (err) { + lastError = err; + } + } + throw lastError; + }, + ), + sleep: vi.fn(), + sleepUntil: vi.fn(), + waitForEvent: vi.fn(), + }; + + const TestWorkflowInstrumented = instrumentWorkflowWithSentry(getSentryOptions, ReusedErrorWorkflow as any); + const workflow = new TestWorkflowInstrumented(context, {}) as ReusedErrorWorkflow; + const event = { payload: {}, timestamp: new Date(), instanceId: INSTANCE_ID }; + + await workflow.run(event, mockStep); + await drainWaitUntilLikeCloudflareVitestPool(waitUntilPromises); + + await expect(workflow.run(event, alwaysFailStep)).rejects.toThrow('second run error'); + await expect(drainWaitUntilLikeCloudflareVitestPool(waitUntilPromises)).resolves.toBeUndefined(); + + const errorEnvelopes = mockTransport.send.mock.calls.filter(call => { + const items = (call[0] as any)[1] as any[]; + return items.some(i => i[0].type === 'event'); + }); + expect(errorEnvelopes).toHaveLength(1); + expect(errorEnvelopes[0]![0][1][0][1]).toMatchObject({ + exception: { + values: [ + expect.objectContaining({ + type: 'Error', + value: 'second run error', + mechanism: { type: 'auto.faas.cloudflare.workflow', handled: true }, + }), + ], + }, + }); + }); + test('Wraps env with instrumentEnv', async () => { class EnvTestWorkflow { constructor(_ctx: ExecutionContext, _env: unknown) {}