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 packages/cloudflare/src/flush.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -35,7 +35,7 @@ const flushLockRegistries = new WeakMap<ExecutionContext['waitUntil'], FlushLock
*
* By using the original waitUntil for flush operations, we bypass this issue.
*/
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] | undefined {
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] {
// eslint-disable-next-line @typescript-eslint/unbound-method
const currentWaitUntil = context.waitUntil;
const original = flushLockRegistries.get(currentWaitUntil)?.originalWaitUntil;
Expand Down
2 changes: 1 addition & 1 deletion packages/cloudflare/src/request.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,7 +58,7 @@ export function wrapRequestHandler(
// to track pending tasks. If we use the instrumented version for flushAndDispose,
// it acquires the lock, then flushAndDispose tries to wait for the same lock,
// creating a deadlock.
const waitUntil = context ? getOriginalWaitUntil(context)?.bind(context) : undefined;
const waitUntil = context ? getOriginalWaitUntil(context).bind(context) : undefined;
const errorMechanismType = getRequestErrorMechanismType(context);

const client = init({ ...options, ctx: context });
Expand Down
4 changes: 2 additions & 2 deletions packages/cloudflare/src/workflows.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,7 +23,7 @@ import type {
} from 'cloudflare:workers';
import { setAsyncLocalStorageAsyncContextStrategy } from './async';
import type { CloudflareOptions } from './client';
import { flushAndDispose } from './flush';
import { flushAndDispose, getOriginalWaitUntil } from './flush';
import { instrumentEnv } from './instrumentations/worker/instrumentEnv';
import { addCloudResourceContext } from './scope-utils';
import { init } from './sdk';
Expand DownExpand Up@@ -214,7 +214,7 @@ export function instrumentWorkflowWithSentry<
setAsyncLocalStorageAsyncContextStrategy();

return withIsolationScope(async isolationScope => {
const waitUntil = context.waitUntil.bind(context);
const waitUntil = getOriginalWaitUntil(context).bind(context);
const client = init({ ...options, ctx: context, enableDedupe: false });
isolationScope.setClient(client);

Expand Down
6 changes: 3 additions & 3 deletions packages/cloudflare/test/flush.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -165,7 +165,7 @@ describe('getOriginalWaitUntil', () => {

expect(result).not.toBe(context.waitUntil);
expect(result).toBeDefined();
result!(Promise.resolve());
result(Promise.resolve());
expect(originalWaitUntil).toHaveBeenCalled();
});

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

Expand All@@ -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();
Expand Down
129 changes: 129 additions & 0 deletions packages/cloudflare/test/workflow.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -208,6 +208,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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
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<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
runCount += 1;
await step.do('reused step', async () => {
if (runCount === 2) {
this._ctx.waitUntil(
new Promise<void>(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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
waitUntilPromises.push(promise);
}),
passThroughOnException: vi.fn(),
props: {},
};

let runCount = 0;

class ReusedErrorWorkflow {
public constructor(private _ctx: ExecutionContext) {}

public async run(_event: Readonly<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
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<any>),
maybeCallback?: (...args: unknown[]) => Promise<any>,
) => {
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) {}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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 packages/cloudflare/src/flush.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -35,7 +35,7 @@ const flushLockRegistries = new WeakMap<ExecutionContext['waitUntil'], FlushLock
*
* By using the original waitUntil for flush operations, we bypass this issue.
*/
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] | undefined {
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] {
// eslint-disable-next-line @typescript-eslint/unbound-method
const currentWaitUntil = context.waitUntil;
const original = flushLockRegistries.get(currentWaitUntil)?.originalWaitUntil;
Expand Down
2 changes: 1 addition & 1 deletion packages/cloudflare/src/request.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,7 +58,7 @@ export function wrapRequestHandler(
// to track pending tasks. If we use the instrumented version for flushAndDispose,
// it acquires the lock, then flushAndDispose tries to wait for the same lock,
// creating a deadlock.
const waitUntil = context ? getOriginalWaitUntil(context)?.bind(context) : undefined;
const waitUntil = context ? getOriginalWaitUntil(context).bind(context) : undefined;
const errorMechanismType = getRequestErrorMechanismType(context);

const client = init({ ...options, ctx: context });
Expand Down
4 changes: 2 additions & 2 deletions packages/cloudflare/src/workflows.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,7 +23,7 @@ import type {
} from 'cloudflare:workers';
import { setAsyncLocalStorageAsyncContextStrategy } from './async';
import type { CloudflareOptions } from './client';
import { flushAndDispose } from './flush';
import { flushAndDispose, getOriginalWaitUntil } from './flush';
import { instrumentEnv } from './instrumentations/worker/instrumentEnv';
import { addCloudResourceContext } from './scope-utils';
import { init } from './sdk';
Expand DownExpand Up@@ -214,7 +214,7 @@ export function instrumentWorkflowWithSentry<
setAsyncLocalStorageAsyncContextStrategy();

return withIsolationScope(async isolationScope => {
const waitUntil = context.waitUntil.bind(context);
const waitUntil = getOriginalWaitUntil(context).bind(context);
const client = init({ ...options, ctx: context, enableDedupe: false });
isolationScope.setClient(client);

Expand Down
6 changes: 3 additions & 3 deletions packages/cloudflare/test/flush.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -165,7 +165,7 @@ describe('getOriginalWaitUntil', () => {

expect(result).not.toBe(context.waitUntil);
expect(result).toBeDefined();
result!(Promise.resolve());
result(Promise.resolve());
expect(originalWaitUntil).toHaveBeenCalled();
});

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

Expand All@@ -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();
Expand Down
129 changes: 129 additions & 0 deletions packages/cloudflare/test/workflow.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -208,6 +208,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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
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<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
runCount += 1;
await step.do('reused step', async () => {
if (runCount === 2) {
this._ctx.waitUntil(
new Promise<void>(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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
waitUntilPromises.push(promise);
}),
passThroughOnException: vi.fn(),
props: {},
};

let runCount = 0;

class ReusedErrorWorkflow {
public constructor(private _ctx: ExecutionContext) {}

public async run(_event: Readonly<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
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<any>),
maybeCallback?: (...args: unknown[]) => Promise<any>,
) => {
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) {}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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 packages/cloudflare/src/flush.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -35,7 +35,7 @@ const flushLockRegistries = new WeakMap<ExecutionContext['waitUntil'], FlushLock
*
* By using the original waitUntil for flush operations, we bypass this issue.
*/
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] | undefined {
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] {
// eslint-disable-next-line @typescript-eslint/unbound-method
const currentWaitUntil = context.waitUntil;
const original = flushLockRegistries.get(currentWaitUntil)?.originalWaitUntil;
Expand Down
2 changes: 1 addition & 1 deletion packages/cloudflare/src/request.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,7 +58,7 @@ export function wrapRequestHandler(
// to track pending tasks. If we use the instrumented version for flushAndDispose,
// it acquires the lock, then flushAndDispose tries to wait for the same lock,
// creating a deadlock.
const waitUntil = context ? getOriginalWaitUntil(context)?.bind(context) : undefined;
const waitUntil = context ? getOriginalWaitUntil(context).bind(context) : undefined;
const errorMechanismType = getRequestErrorMechanismType(context);

const client = init({ ...options, ctx: context });
Expand Down
4 changes: 2 additions & 2 deletions packages/cloudflare/src/workflows.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,7 +23,7 @@ import type {
} from 'cloudflare:workers';
import { setAsyncLocalStorageAsyncContextStrategy } from './async';
import type { CloudflareOptions } from './client';
import { flushAndDispose } from './flush';
import { flushAndDispose, getOriginalWaitUntil } from './flush';
import { instrumentEnv } from './instrumentations/worker/instrumentEnv';
import { addCloudResourceContext } from './scope-utils';
import { init } from './sdk';
Expand DownExpand Up@@ -214,7 +214,7 @@ export function instrumentWorkflowWithSentry<
setAsyncLocalStorageAsyncContextStrategy();

return withIsolationScope(async isolationScope => {
const waitUntil = context.waitUntil.bind(context);
const waitUntil = getOriginalWaitUntil(context).bind(context);
const client = init({ ...options, ctx: context, enableDedupe: false });
isolationScope.setClient(client);

Expand Down
6 changes: 3 additions & 3 deletions packages/cloudflare/test/flush.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -165,7 +165,7 @@ describe('getOriginalWaitUntil', () => {

expect(result).not.toBe(context.waitUntil);
expect(result).toBeDefined();
result!(Promise.resolve());
result(Promise.resolve());
expect(originalWaitUntil).toHaveBeenCalled();
});

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

Expand All@@ -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();
Expand Down
129 changes: 129 additions & 0 deletions packages/cloudflare/test/workflow.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -208,6 +208,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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
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<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
runCount += 1;
await step.do('reused step', async () => {
if (runCount === 2) {
this._ctx.waitUntil(
new Promise<void>(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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
waitUntilPromises.push(promise);
}),
passThroughOnException: vi.fn(),
props: {},
};

let runCount = 0;

class ReusedErrorWorkflow {
public constructor(private _ctx: ExecutionContext) {}

public async run(_event: Readonly<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
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<any>),
maybeCallback?: (...args: unknown[]) => Promise<any>,
) => {
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) {}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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 packages/cloudflare/src/flush.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -35,7 +35,7 @@ const flushLockRegistries = new WeakMap<ExecutionContext['waitUntil'], FlushLock
*
* By using the original waitUntil for flush operations, we bypass this issue.
*/
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] | undefined {
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] {
// eslint-disable-next-line @typescript-eslint/unbound-method
const currentWaitUntil = context.waitUntil;
const original = flushLockRegistries.get(currentWaitUntil)?.originalWaitUntil;
Expand Down
2 changes: 1 addition & 1 deletion packages/cloudflare/src/request.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,7 +58,7 @@ export function wrapRequestHandler(
// to track pending tasks. If we use the instrumented version for flushAndDispose,
// it acquires the lock, then flushAndDispose tries to wait for the same lock,
// creating a deadlock.
const waitUntil = context ? getOriginalWaitUntil(context)?.bind(context) : undefined;
const waitUntil = context ? getOriginalWaitUntil(context).bind(context) : undefined;
const errorMechanismType = getRequestErrorMechanismType(context);

const client = init({ ...options, ctx: context });
Expand Down
4 changes: 2 additions & 2 deletions packages/cloudflare/src/workflows.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,7 +23,7 @@ import type {
} from 'cloudflare:workers';
import { setAsyncLocalStorageAsyncContextStrategy } from './async';
import type { CloudflareOptions } from './client';
import { flushAndDispose } from './flush';
import { flushAndDispose, getOriginalWaitUntil } from './flush';
import { instrumentEnv } from './instrumentations/worker/instrumentEnv';
import { addCloudResourceContext } from './scope-utils';
import { init } from './sdk';
Expand DownExpand Up@@ -214,7 +214,7 @@ export function instrumentWorkflowWithSentry<
setAsyncLocalStorageAsyncContextStrategy();

return withIsolationScope(async isolationScope => {
const waitUntil = context.waitUntil.bind(context);
const waitUntil = getOriginalWaitUntil(context).bind(context);
const client = init({ ...options, ctx: context, enableDedupe: false });
isolationScope.setClient(client);

Expand Down
6 changes: 3 additions & 3 deletions packages/cloudflare/test/flush.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -165,7 +165,7 @@ describe('getOriginalWaitUntil', () => {

expect(result).not.toBe(context.waitUntil);
expect(result).toBeDefined();
result!(Promise.resolve());
result(Promise.resolve());
expect(originalWaitUntil).toHaveBeenCalled();
});

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

Expand All@@ -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();
Expand Down
129 changes: 129 additions & 0 deletions packages/cloudflare/test/workflow.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -208,6 +208,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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
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<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
runCount += 1;
await step.do('reused step', async () => {
if (runCount === 2) {
this._ctx.waitUntil(
new Promise<void>(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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
waitUntilPromises.push(promise);
}),
passThroughOnException: vi.fn(),
props: {},
};

let runCount = 0;

class ReusedErrorWorkflow {
public constructor(private _ctx: ExecutionContext) {}

public async run(_event: Readonly<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
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<any>),
maybeCallback?: (...args: unknown[]) => Promise<any>,
) => {
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) {}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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 packages/cloudflare/src/flush.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -35,7 +35,7 @@ const flushLockRegistries = new WeakMap<ExecutionContext['waitUntil'], FlushLock
*
* By using the original waitUntil for flush operations, we bypass this issue.
*/
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] | undefined {
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] {
// eslint-disable-next-line @typescript-eslint/unbound-method
const currentWaitUntil = context.waitUntil;
const original = flushLockRegistries.get(currentWaitUntil)?.originalWaitUntil;
Expand Down
2 changes: 1 addition & 1 deletion packages/cloudflare/src/request.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,7 +58,7 @@ export function wrapRequestHandler(
// to track pending tasks. If we use the instrumented version for flushAndDispose,
// it acquires the lock, then flushAndDispose tries to wait for the same lock,
// creating a deadlock.
const waitUntil = context ? getOriginalWaitUntil(context)?.bind(context) : undefined;
const waitUntil = context ? getOriginalWaitUntil(context).bind(context) : undefined;
const errorMechanismType = getRequestErrorMechanismType(context);

const client = init({ ...options, ctx: context });
Expand Down
4 changes: 2 additions & 2 deletions packages/cloudflare/src/workflows.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,7 +23,7 @@ import type {
} from 'cloudflare:workers';
import { setAsyncLocalStorageAsyncContextStrategy } from './async';
import type { CloudflareOptions } from './client';
import { flushAndDispose } from './flush';
import { flushAndDispose, getOriginalWaitUntil } from './flush';
import { instrumentEnv } from './instrumentations/worker/instrumentEnv';
import { addCloudResourceContext } from './scope-utils';
import { init } from './sdk';
Expand DownExpand Up@@ -214,7 +214,7 @@ export function instrumentWorkflowWithSentry<
setAsyncLocalStorageAsyncContextStrategy();

return withIsolationScope(async isolationScope => {
const waitUntil = context.waitUntil.bind(context);
const waitUntil = getOriginalWaitUntil(context).bind(context);
const client = init({ ...options, ctx: context, enableDedupe: false });
isolationScope.setClient(client);

Expand Down
6 changes: 3 additions & 3 deletions packages/cloudflare/test/flush.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -165,7 +165,7 @@ describe('getOriginalWaitUntil', () => {

expect(result).not.toBe(context.waitUntil);
expect(result).toBeDefined();
result!(Promise.resolve());
result(Promise.resolve());
expect(originalWaitUntil).toHaveBeenCalled();
});

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

Expand All@@ -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();
Expand Down
129 changes: 129 additions & 0 deletions packages/cloudflare/test/workflow.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -208,6 +208,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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
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<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
runCount += 1;
await step.do('reused step', async () => {
if (runCount === 2) {
this._ctx.waitUntil(
new Promise<void>(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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
waitUntilPromises.push(promise);
}),
passThroughOnException: vi.fn(),
props: {},
};

let runCount = 0;

class ReusedErrorWorkflow {
public constructor(private _ctx: ExecutionContext) {}

public async run(_event: Readonly<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
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<any>),
maybeCallback?: (...args: unknown[]) => Promise<any>,
) => {
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) {}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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 packages/cloudflare/src/flush.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -35,7 +35,7 @@ const flushLockRegistries = new WeakMap<ExecutionContext['waitUntil'], FlushLock
*
* By using the original waitUntil for flush operations, we bypass this issue.
*/
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] | undefined {
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] {
// eslint-disable-next-line @typescript-eslint/unbound-method
const currentWaitUntil = context.waitUntil;
const original = flushLockRegistries.get(currentWaitUntil)?.originalWaitUntil;
Expand Down
2 changes: 1 addition & 1 deletion packages/cloudflare/src/request.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,7 +58,7 @@ export function wrapRequestHandler(
// to track pending tasks. If we use the instrumented version for flushAndDispose,
// it acquires the lock, then flushAndDispose tries to wait for the same lock,
// creating a deadlock.
const waitUntil = context ? getOriginalWaitUntil(context)?.bind(context) : undefined;
const waitUntil = context ? getOriginalWaitUntil(context).bind(context) : undefined;
const errorMechanismType = getRequestErrorMechanismType(context);

const client = init({ ...options, ctx: context });
Expand Down
4 changes: 2 additions & 2 deletions packages/cloudflare/src/workflows.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,7 +23,7 @@ import type {
} from 'cloudflare:workers';
import { setAsyncLocalStorageAsyncContextStrategy } from './async';
import type { CloudflareOptions } from './client';
import { flushAndDispose } from './flush';
import { flushAndDispose, getOriginalWaitUntil } from './flush';
import { instrumentEnv } from './instrumentations/worker/instrumentEnv';
import { addCloudResourceContext } from './scope-utils';
import { init } from './sdk';
Expand DownExpand Up@@ -214,7 +214,7 @@ export function instrumentWorkflowWithSentry<
setAsyncLocalStorageAsyncContextStrategy();

return withIsolationScope(async isolationScope => {
const waitUntil = context.waitUntil.bind(context);
const waitUntil = getOriginalWaitUntil(context).bind(context);
const client = init({ ...options, ctx: context, enableDedupe: false });
isolationScope.setClient(client);

Expand Down
6 changes: 3 additions & 3 deletions packages/cloudflare/test/flush.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -165,7 +165,7 @@ describe('getOriginalWaitUntil', () => {

expect(result).not.toBe(context.waitUntil);
expect(result).toBeDefined();
result!(Promise.resolve());
result(Promise.resolve());
expect(originalWaitUntil).toHaveBeenCalled();
});

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

Expand All@@ -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();
Expand Down
129 changes: 129 additions & 0 deletions packages/cloudflare/test/workflow.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -208,6 +208,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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
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<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
runCount += 1;
await step.do('reused step', async () => {
if (runCount === 2) {
this._ctx.waitUntil(
new Promise<void>(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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
waitUntilPromises.push(promise);
}),
passThroughOnException: vi.fn(),
props: {},
};

let runCount = 0;

class ReusedErrorWorkflow {
public constructor(private _ctx: ExecutionContext) {}

public async run(_event: Readonly<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
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<any>),
maybeCallback?: (...args: unknown[]) => Promise<any>,
) => {
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) {}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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 packages/cloudflare/src/flush.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -35,7 +35,7 @@ const flushLockRegistries = new WeakMap<ExecutionContext['waitUntil'], FlushLock
*
* By using the original waitUntil for flush operations, we bypass this issue.
*/
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] | undefined {
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] {
// eslint-disable-next-line @typescript-eslint/unbound-method
const currentWaitUntil = context.waitUntil;
const original = flushLockRegistries.get(currentWaitUntil)?.originalWaitUntil;
Expand Down
2 changes: 1 addition & 1 deletion packages/cloudflare/src/request.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,7 +58,7 @@ export function wrapRequestHandler(
// to track pending tasks. If we use the instrumented version for flushAndDispose,
// it acquires the lock, then flushAndDispose tries to wait for the same lock,
// creating a deadlock.
const waitUntil = context ? getOriginalWaitUntil(context)?.bind(context) : undefined;
const waitUntil = context ? getOriginalWaitUntil(context).bind(context) : undefined;
const errorMechanismType = getRequestErrorMechanismType(context);

const client = init({ ...options, ctx: context });
Expand Down
4 changes: 2 additions & 2 deletions packages/cloudflare/src/workflows.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,7 +23,7 @@ import type {
} from 'cloudflare:workers';
import { setAsyncLocalStorageAsyncContextStrategy } from './async';
import type { CloudflareOptions } from './client';
import { flushAndDispose } from './flush';
import { flushAndDispose, getOriginalWaitUntil } from './flush';
import { instrumentEnv } from './instrumentations/worker/instrumentEnv';
import { addCloudResourceContext } from './scope-utils';
import { init } from './sdk';
Expand DownExpand Up@@ -214,7 +214,7 @@ export function instrumentWorkflowWithSentry<
setAsyncLocalStorageAsyncContextStrategy();

return withIsolationScope(async isolationScope => {
const waitUntil = context.waitUntil.bind(context);
const waitUntil = getOriginalWaitUntil(context).bind(context);
const client = init({ ...options, ctx: context, enableDedupe: false });
isolationScope.setClient(client);

Expand Down
6 changes: 3 additions & 3 deletions packages/cloudflare/test/flush.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -165,7 +165,7 @@ describe('getOriginalWaitUntil', () => {

expect(result).not.toBe(context.waitUntil);
expect(result).toBeDefined();
result!(Promise.resolve());
result(Promise.resolve());
expect(originalWaitUntil).toHaveBeenCalled();
});

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

Expand All@@ -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();
Expand Down
129 changes: 129 additions & 0 deletions packages/cloudflare/test/workflow.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -208,6 +208,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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
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<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
runCount += 1;
await step.do('reused step', async () => {
if (runCount === 2) {
this._ctx.waitUntil(
new Promise<void>(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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
waitUntilPromises.push(promise);
}),
passThroughOnException: vi.fn(),
props: {},
};

let runCount = 0;

class ReusedErrorWorkflow {
public constructor(private _ctx: ExecutionContext) {}

public async run(_event: Readonly<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
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<any>),
maybeCallback?: (...args: unknown[]) => Promise<any>,
) => {
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) {}
Expand Down
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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 packages/cloudflare/src/flush.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -35,7 +35,7 @@ const flushLockRegistries = new WeakMap<ExecutionContext['waitUntil'], FlushLock
*
* By using the original waitUntil for flush operations, we bypass this issue.
*/
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] | undefined {
export function getOriginalWaitUntil(context: ExecutionContextCompat): ExecutionContext['waitUntil'] {
// eslint-disable-next-line @typescript-eslint/unbound-method
const currentWaitUntil = context.waitUntil;
const original = flushLockRegistries.get(currentWaitUntil)?.originalWaitUntil;
Expand Down
2 changes: 1 addition & 1 deletion packages/cloudflare/src/request.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -58,7 +58,7 @@ export function wrapRequestHandler(
// to track pending tasks. If we use the instrumented version for flushAndDispose,
// it acquires the lock, then flushAndDispose tries to wait for the same lock,
// creating a deadlock.
const waitUntil = context ? getOriginalWaitUntil(context)?.bind(context) : undefined;
const waitUntil = context ? getOriginalWaitUntil(context).bind(context) : undefined;
const errorMechanismType = getRequestErrorMechanismType(context);

const client = init({ ...options, ctx: context });
Expand Down
4 changes: 2 additions & 2 deletions packages/cloudflare/src/workflows.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,7 +23,7 @@ import type {
} from 'cloudflare:workers';
import { setAsyncLocalStorageAsyncContextStrategy } from './async';
import type { CloudflareOptions } from './client';
import { flushAndDispose } from './flush';
import { flushAndDispose, getOriginalWaitUntil } from './flush';
import { instrumentEnv } from './instrumentations/worker/instrumentEnv';
import { addCloudResourceContext } from './scope-utils';
import { init } from './sdk';
Expand DownExpand Up@@ -214,7 +214,7 @@ export function instrumentWorkflowWithSentry<
setAsyncLocalStorageAsyncContextStrategy();

return withIsolationScope(async isolationScope => {
const waitUntil = context.waitUntil.bind(context);
const waitUntil = getOriginalWaitUntil(context).bind(context);
const client = init({ ...options, ctx: context, enableDedupe: false });
isolationScope.setClient(client);

Expand Down
6 changes: 3 additions & 3 deletions packages/cloudflare/test/flush.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -165,7 +165,7 @@ describe('getOriginalWaitUntil', () => {

expect(result).not.toBe(context.waitUntil);
expect(result).toBeDefined();
result!(Promise.resolve());
result(Promise.resolve());
expect(originalWaitUntil).toHaveBeenCalled();
});

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

Expand All@@ -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();
Expand Down
129 changes: 129 additions & 0 deletions packages/cloudflare/test/workflow.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -208,6 +208,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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
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<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
runCount += 1;
await step.do('reused step', async () => {
if (runCount === 2) {
this._ctx.waitUntil(
new Promise<void>(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<unknown>[] = [];
const context: ExecutionContext = {
waitUntil: vi.fn((promise: Promise<unknown>) => {
waitUntilPromises.push(promise);
}),
passThroughOnException: vi.fn(),
props: {},
};

let runCount = 0;

class ReusedErrorWorkflow {
public constructor(private _ctx: ExecutionContext) {}

public async run(_event: Readonly<WorkflowEvent<Params>>, step: WorkflowStep): Promise<void> {
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<any>),
maybeCallback?: (...args: unknown[]) => Promise<any>,
) => {
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) {}
Expand Down
Loading