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
29 changes: 29 additions & 0 deletions packages/worker-utils/src/do-retry-scope.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
import { afterEach, describe, expect, it, vi } from 'vitest';
import { DEFAULT_DO_RETRY_CONFIG, withDORetry } from './do-retry.js';

afterEach(() => vi.useRealTimers());

describe('withDORetry scopes', () => {
it('does not start another attempt after its deadline expires', async () => {
vi.useFakeTimers();
const operation = vi.fn(() => new Promise<never>(() => undefined));
const pending = withDORetry(
() => ({}),
operation,
'scoped_operation',
{
...DEFAULT_DO_RETRY_CONFIG,
scope: { deadlineAt: Date.now() + 100 },
},
{ warn: () => undefined, error: () => undefined }
);
const outcome = pending.then(
() => undefined,
error => error
);

await vi.advanceTimersByTimeAsync(100);
await expect(outcome).resolves.toMatchObject({ name: 'TimeoutError' });
expect(operation).toHaveBeenCalledOnce();
});
});
172 changes: 133 additions & 39 deletions packages/worker-utils/src/do-retry.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,20 @@
// Cloudflare Workers provides scheduler.wait() for cooperative delays.
// Not in standard webworker lib types.
declare const scheduler: undefined | { wait(ms: number): Promise<void> };
declare const scheduler:
| undefined
| { wait(ms: number, options?: { signal?: AbortSignal }): Promise<void> };

export type DORetryScope = {
deadlineAt: number;
signal?: AbortSignal;
assertCurrent?: () => void;
};

export type DORetryConfig = {
maxAttempts: number;
baseBackoffMs: number;
maxBackoffMs: number;
scope?: DORetryScope;
};

export const DEFAULT_DO_RETRY_CONFIG: DORetryConfig = {
Expand Down Expand Up @@ -44,11 +53,76 @@ function calculateBackoff(attempt: number, config: DORetryConfig): number {
return Math.min(config.maxBackoffMs, jitteredBackoff);
}

function waitMs(ms: number): Promise<void> {
function waitMs(ms: number, signal?: AbortSignal): Promise<void> {
signal?.throwIfAborted();
if (typeof scheduler !== 'undefined' && 'wait' in scheduler) {
return scheduler.wait(ms);
return signal ? scheduler.wait(ms, { signal }) : scheduler.wait(ms);
}
if (!signal) return new Promise(resolve => setTimeout(resolve, ms));
return new Promise((resolve, reject) => {
const onAbort = () => {
clearTimeout(timeoutId);
signal.removeEventListener('abort', onAbort);
reject(signal.reason);
};
const timeoutId = setTimeout(() => {
signal.removeEventListener('abort', onAbort);
resolve();
}, ms);
signal.addEventListener('abort', onAbort, { once: true });
if (signal.aborted) onAbort();
});
}

function createRetryScope({ deadlineAt, signal, assertCurrent }: DORetryScope) {
if (!Number.isFinite(deadlineAt))
throw new RangeError('Durable Object retry deadlineAt must be finite');

const controller = new AbortController();
const deadlineError = new DOMException('Durable Object retry deadline exceeded', 'TimeoutError');
const onAbort = () => controller.abort(signal?.reason);
signal?.addEventListener('abort', onAbort, { once: true });
if (signal?.aborted) onAbort();
const timeoutId = setTimeout(
() => controller.abort(deadlineError),
Math.max(0, deadlineAt - Date.now())
);

return {
deadlineAt,
signal: controller.signal,
assertActive() {
if (Date.now() >= deadlineAt) controller.abort(deadlineError);
controller.signal.throwIfAborted();
try {
assertCurrent?.();
} catch (error) {
controller.abort(error);
throw error;
}
if (Date.now() >= deadlineAt) controller.abort(deadlineError);
controller.signal.throwIfAborted();
},
dispose() {
clearTimeout(timeoutId);
signal?.removeEventListener('abort', onAbort);
},
};
}

async function waitWithSignal<T>(pending: Promise<T>, signal: AbortSignal): Promise<T> {
let onAbort: (() => void) | undefined;
const cancelled = new Promise<never>((_, reject) => {
onAbort = () => reject(signal.reason);
signal.addEventListener('abort', onAbort, { once: true });
if (signal.aborted) onAbort();
});

try {
return await Promise.race([pending, cancelled]);
} finally {
if (onAbort) signal.removeEventListener('abort', onAbort);
}
return new Promise(resolve => setTimeout(resolve, ms));
}

type DORetryLogger = {
Expand Down Expand Up @@ -87,48 +161,68 @@ export async function withDORetry<TStub, TResult>(
logger: DORetryLogger = console
): Promise<TResult> {
let lastError: Error | undefined;
const scope = config.scope ? createRetryScope(config.scope) : undefined;

for (let attempt = 0; attempt < config.maxAttempts; attempt++) {
try {
// Create fresh stub for each attempt
const stub = getStub();
return await operation(stub);
} catch (error) {
lastError = error instanceof Error ? error : new Error(String(error));

// Check if we should retry
if (!isRetryableError(error)) {
logger.warn('[do-retry] Non-retryable error', {
try {
for (let attempt = 0; attempt < config.maxAttempts; attempt++) {
scope?.assertActive();
try {
// Create fresh stub for each attempt
const stub = getStub();
scope?.assertActive();
const result = scope
? await waitWithSignal(operation(stub), scope.signal)
: await operation(stub);
scope?.assertActive();
return result;
} catch (error) {
scope?.assertActive();
lastError = error instanceof Error ? error : new Error(String(error));

// Check if we should retry
if (!isRetryableError(error)) {
logger.warn('[do-retry] Non-retryable error', {
operation: operationName,
attempt: attempt + 1,
error: lastError.message,
retryable: false,
});
throw lastError;
}

// Check if we have retries left
if (attempt + 1 >= config.maxAttempts) {
logger.error('[do-retry] All retry attempts exhausted', {
operation: operationName,
attempts: attempt + 1,
error: lastError.message,
});
throw lastError;
}

// Calculate backoff and wait
const requestedBackoffMs = calculateBackoff(attempt, config);
const backoffMs = scope
? Math.min(requestedBackoffMs, Math.max(0, scope.deadlineAt - Date.now()))
: requestedBackoffMs;
logger.warn('[do-retry] Retrying', {
operation: operationName,
attempt: attempt + 1,
backoffMs: Math.round(backoffMs),
error: lastError.message,
retryable: false,
});
throw lastError;
}

// Check if we have retries left
if (attempt + 1 >= config.maxAttempts) {
logger.error('[do-retry] All retry attempts exhausted', {
operation: operationName,
attempts: attempt + 1,
error: lastError.message,
});
throw lastError;
scope?.assertActive();
try {
await waitMs(backoffMs, scope?.signal);
} finally {
scope?.assertActive();
}
}

// Calculate backoff and wait
const backoffMs = calculateBackoff(attempt, config);
logger.warn('[do-retry] Retrying', {
operation: operationName,
attempt: attempt + 1,
backoffMs: Math.round(backoffMs),
error: lastError.message,
});

await waitMs(backoffMs);
}
}

throw lastError ?? new Error('Unexpected retry loop exit');
throw lastError ?? new Error('Unexpected retry loop exit');
} finally {
scope?.dispose();
}
}
2 changes: 1 addition & 1 deletion packages/worker-utils/src/index.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
export { getCachedSecret, clearSecretCacheForTest } from './cached-secret.js';

export { withDORetry, DEFAULT_DO_RETRY_CONFIG } from './do-retry.js';
export type { DORetryConfig } from './do-retry.js';
export type { DORetryConfig, DORetryScope } from './do-retry.js';

export { backendAuthMiddleware } from './backend-auth-middleware.js';

Expand Down
Loading