Skip to content

Commit 8d117cd

Browse files
pranaygpclaude
andauthored
Retry 5xx errors from workflow-server in step handler (#1011)
* Retry 5xx errors from workflow-server in step handler Add `withServerErrorRetry` helper that retries world calls on 5xx errors with exponential backoff (500ms, 1s, 2s ≈ 3.5s total). Applied to all `world.events.create` calls in the step handler so transient workflow-server errors don't consume step attempts. If retries are exhausted, the error is thrown to the queue for higher-level retry without burning a step attempt. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com> * Address PR review comments on 5xx retry handling - Fix misleading `maxAttempts` log field to `maxRetries` in withServerErrorRetry - Update step-handler comment to accurately note that queue retries may still consume step attempts since step_started has already incremented - Add unit tests for withServerErrorRetry (7 tests covering success, retry/backoff, exhaustion, and non-5xx passthrough) Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com> * Add tests for 429/5xx retry handling Unit tests for withThrottleRetry and withServerErrorRetry helpers, plus an e2e test that exercises the 5xx retry codepath during step execution via run-scoped fault injection. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com> * Wrap VQS errors in WorkflowAPIError and add retry to queueMessage VQS throws its own error types (InternalServerError, ConsumerDiscoveryError, ConsumerRegistryNotConfiguredError) that don't match WorkflowAPIError.is(). Wrapping them at the world-vercel boundary enables withServerErrorRetry in queueMessage() to automatically retry transient queue failures. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com> * Update changeset to include @workflow/world-vercel Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com> * Remove queue retry logic per review feedback Queue retrying will be handled natively by the @vercel/queue client instead. Reverts VQS error wrapping and withServerErrorRetry in queueMessage(). Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
1 parent 94760b4 commit 8d117cd

6 files changed

Lines changed: 499 additions & 45 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"@workflow/core": patch
3+
---
4+
5+
Retry 5xx errors from workflow-server in step handler to avoid consuming step attempts on transient infrastructure errors

‎packages/core/e2e/e2e.test.ts‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -889,6 +889,33 @@ describe('e2e', () => {
889889
expect(result.failed).toBe(true);
890890
expect(result.attempt).toBe(1);
891891
});
892+
893+
test(
894+
'workflow completes despite transient 5xx on step_completed',
895+
{ timeout: 120_000 },
896+
async () => {
897+
const run = await start(
898+
await e2e('serverError5xxRetryWorkflow'),
899+
[42]
900+
);
901+
const result = await run.returnValue;
902+
903+
// Correct result proves workflow completed successfully
904+
expect(result.result).toBe(84); // 42 * 2
905+
906+
// retryCount > 0 proves the fault injection actually triggered
907+
expect(result.retryCount).toBe(2);
908+
909+
// attempt === 1 proves no step attempt was consumed by the 5xx retries
910+
const { json: steps } = await cliInspectJson(
911+
`steps --runId ${run.runId}`
912+
);
913+
const doWorkStep = steps.find((s: any) =>
914+
s.stepName.includes('doWork')
915+
);
916+
expect(doWorkStep.attempt).toBe(1);
917+
}
918+
);
892919
});
893920

894921
describe('catchability', () => {

‎packages/core/src/runtime/helpers.test.ts‎

Lines changed: 241 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,20 @@
1-
import { describe, expect, it } from 'vitest';
2-
import { getWorkflowQueueName } from './helpers';
1+
import { WorkflowAPIError } from '@workflow/errors';
2+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
3+
import {
4+
getWorkflowQueueName,
5+
withServerErrorRetry,
6+
withThrottleRetry,
7+
} from './helpers.js';
8+
9+
// Mock the logger to suppress output during tests
10+
vi.mock('../logger.js', () => ({
11+
runtimeLogger: {
12+
warn: vi.fn(),
13+
debug: vi.fn(),
14+
info: vi.fn(),
15+
error: vi.fn(),
16+
},
17+
}));
318

419
describe('getWorkflowQueueName', () => {
520
it('should return a valid queue name for a simple workflow name', () => {
@@ -70,3 +85,227 @@ describe('getWorkflowQueueName', () => {
7085
expect(() => getWorkflowQueueName('')).toThrow('Invalid workflow name');
7186
});
7287
});
88+
89+
describe('withServerErrorRetry', () => {
90+
beforeEach(() => {
91+
vi.useFakeTimers();
92+
});
93+
94+
afterEach(() => {
95+
vi.useRealTimers();
96+
vi.clearAllMocks();
97+
});
98+
99+
it('should return the result on success', async () => {
100+
const fn = vi.fn().mockResolvedValue('ok');
101+
const result = await withServerErrorRetry(fn);
102+
expect(result).toBe('ok');
103+
expect(fn).toHaveBeenCalledTimes(1);
104+
});
105+
106+
it('should retry on 5xx WorkflowAPIError and succeed', async () => {
107+
const fn = vi
108+
.fn()
109+
.mockRejectedValueOnce(
110+
new WorkflowAPIError('Internal Server Error', { status: 500 })
111+
)
112+
.mockResolvedValueOnce('recovered');
113+
114+
const promise = withServerErrorRetry(fn);
115+
await vi.advanceTimersByTimeAsync(500);
116+
const result = await promise;
117+
118+
expect(result).toBe('recovered');
119+
expect(fn).toHaveBeenCalledTimes(2);
120+
});
121+
122+
it('should retry up to 3 times with exponential backoff', async () => {
123+
const fn = vi
124+
.fn()
125+
.mockRejectedValueOnce(new WorkflowAPIError('error', { status: 502 }))
126+
.mockRejectedValueOnce(new WorkflowAPIError('error', { status: 503 }))
127+
.mockRejectedValueOnce(new WorkflowAPIError('error', { status: 500 }))
128+
.mockResolvedValueOnce('finally');
129+
130+
const promise = withServerErrorRetry(fn);
131+
132+
// First retry after 500ms
133+
await vi.advanceTimersByTimeAsync(500);
134+
// Second retry after 1000ms
135+
await vi.advanceTimersByTimeAsync(1000);
136+
// Third retry after 2000ms
137+
await vi.advanceTimersByTimeAsync(2000);
138+
139+
const result = await promise;
140+
expect(result).toBe('finally');
141+
expect(fn).toHaveBeenCalledTimes(4);
142+
});
143+
144+
it('should throw after exhausting all retries', async () => {
145+
const error = new WorkflowAPIError('server down', { status: 500 });
146+
const fn = vi.fn().mockRejectedValue(error);
147+
148+
const promise = withServerErrorRetry(fn).catch((e) => e);
149+
150+
// Advance through all 3 retry delays
151+
await vi.advanceTimersByTimeAsync(500);
152+
await vi.advanceTimersByTimeAsync(1000);
153+
await vi.advanceTimersByTimeAsync(2000);
154+
155+
const result = await promise;
156+
expect(result).toBeInstanceOf(WorkflowAPIError);
157+
expect(result.message).toBe('server down');
158+
// 1 initial + 3 retries = 4 total calls
159+
expect(fn).toHaveBeenCalledTimes(4);
160+
});
161+
162+
it('should not retry non-5xx WorkflowAPIErrors', async () => {
163+
const error = new WorkflowAPIError('Not Found', { status: 404 });
164+
const fn = vi.fn().mockRejectedValue(error);
165+
166+
await expect(withServerErrorRetry(fn)).rejects.toThrow('Not Found');
167+
expect(fn).toHaveBeenCalledTimes(1);
168+
});
169+
170+
it('should not retry non-WorkflowAPIError errors', async () => {
171+
const error = new Error('some other error');
172+
const fn = vi.fn().mockRejectedValue(error);
173+
174+
await expect(withServerErrorRetry(fn)).rejects.toThrow('some other error');
175+
expect(fn).toHaveBeenCalledTimes(1);
176+
});
177+
178+
it('should not retry 429 errors (handled by withThrottleRetry)', async () => {
179+
const error = new WorkflowAPIError('Too Many Requests', {
180+
status: 429,
181+
retryAfter: 5,
182+
});
183+
const fn = vi.fn().mockRejectedValue(error);
184+
185+
await expect(withServerErrorRetry(fn)).rejects.toThrow('Too Many Requests');
186+
expect(fn).toHaveBeenCalledTimes(1);
187+
});
188+
});
189+
190+
describe('withThrottleRetry', () => {
191+
beforeEach(() => {
192+
vi.useFakeTimers();
193+
});
194+
195+
afterEach(() => {
196+
vi.useRealTimers();
197+
vi.clearAllMocks();
198+
});
199+
200+
it('should pass through the result on success', async () => {
201+
const fn = vi.fn().mockResolvedValue(undefined);
202+
const result = await withThrottleRetry(fn);
203+
expect(result).toBeUndefined();
204+
expect(fn).toHaveBeenCalledTimes(1);
205+
});
206+
207+
it('should pass through { timeoutSeconds } returned by fn', async () => {
208+
const fn = vi.fn().mockResolvedValue({ timeoutSeconds: 42 });
209+
const result = await withThrottleRetry(fn);
210+
expect(result).toEqual({ timeoutSeconds: 42 });
211+
expect(fn).toHaveBeenCalledTimes(1);
212+
});
213+
214+
it('should re-throw non-429 errors including 5xx', async () => {
215+
const error = new WorkflowAPIError('Internal Server Error', {
216+
status: 500,
217+
});
218+
const fn = vi.fn().mockRejectedValue(error);
219+
220+
await expect(withThrottleRetry(fn)).rejects.toThrow(
221+
'Internal Server Error'
222+
);
223+
expect(fn).toHaveBeenCalledTimes(1);
224+
});
225+
226+
it('should re-throw non-WorkflowAPIError errors', async () => {
227+
const error = new Error('random failure');
228+
const fn = vi.fn().mockRejectedValue(error);
229+
230+
await expect(withThrottleRetry(fn)).rejects.toThrow('random failure');
231+
expect(fn).toHaveBeenCalledTimes(1);
232+
});
233+
234+
it('should wait in-process and retry once for short retryAfter (<10s)', async () => {
235+
const fn = vi
236+
.fn()
237+
.mockRejectedValueOnce(
238+
new WorkflowAPIError('Throttled', { status: 429, retryAfter: 5 })
239+
)
240+
.mockResolvedValueOnce(undefined);
241+
242+
const promise = withThrottleRetry(fn);
243+
// Advance past the 5s wait
244+
await vi.advanceTimersByTimeAsync(5000);
245+
const result = await promise;
246+
247+
expect(result).toBeUndefined();
248+
expect(fn).toHaveBeenCalledTimes(2);
249+
});
250+
251+
it('should defer to queue when both attempts are throttled (double 429)', async () => {
252+
const fn = vi
253+
.fn()
254+
.mockRejectedValueOnce(
255+
new WorkflowAPIError('Throttled', { status: 429, retryAfter: 3 })
256+
)
257+
.mockRejectedValueOnce(
258+
new WorkflowAPIError('Throttled again', { status: 429, retryAfter: 7 })
259+
);
260+
261+
const promise = withThrottleRetry(fn);
262+
// Advance past the 3s in-process wait
263+
await vi.advanceTimersByTimeAsync(3000);
264+
const result = await promise;
265+
266+
expect(result).toEqual({ timeoutSeconds: 7 });
267+
expect(fn).toHaveBeenCalledTimes(2);
268+
});
269+
270+
it('should re-throw non-429 error on retry failure', async () => {
271+
const fn = vi
272+
.fn()
273+
.mockRejectedValueOnce(
274+
new WorkflowAPIError('Throttled', { status: 429, retryAfter: 2 })
275+
)
276+
.mockRejectedValueOnce(new Error('connection lost'));
277+
278+
// Capture the rejection early to prevent unhandled rejection warning
279+
const promise = withThrottleRetry(fn).catch((e) => e);
280+
await vi.advanceTimersByTimeAsync(2000);
281+
282+
const result = await promise;
283+
expect(result).toBeInstanceOf(Error);
284+
expect(result.message).toBe('connection lost');
285+
expect(fn).toHaveBeenCalledTimes(2);
286+
});
287+
288+
it('should defer to queue immediately for long retryAfter (>=10s)', async () => {
289+
const fn = vi
290+
.fn()
291+
.mockRejectedValueOnce(
292+
new WorkflowAPIError('Throttled', { status: 429, retryAfter: 15 })
293+
);
294+
295+
const result = await withThrottleRetry(fn);
296+
297+
expect(result).toEqual({ timeoutSeconds: 15 });
298+
expect(fn).toHaveBeenCalledTimes(1);
299+
});
300+
301+
it('should default to 30s (defer to queue) when no retryAfter is provided', async () => {
302+
const error = new WorkflowAPIError('Throttled', { status: 429 });
303+
// retryAfter is undefined, so it defaults to 30 (>=10 → defer)
304+
const fn = vi.fn().mockRejectedValue(error);
305+
306+
const result = await withThrottleRetry(fn);
307+
308+
expect(result).toEqual({ timeoutSeconds: 30 });
309+
expect(fn).toHaveBeenCalledTimes(1);
310+
});
311+
});

‎packages/core/src/runtime/helpers.ts‎

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -468,3 +468,43 @@ export async function withThrottleRetry(
468468
throw err;
469469
}
470470
}
471+
472+
/**
473+
* Retries a function when it throws a 5xx WorkflowAPIError.
474+
* Used to handle transient workflow-server errors without consuming step attempts.
475+
*
476+
* Retries up to 3 times with exponential backoff (500ms, 1s, 2s ≈ 3.5s total).
477+
* If all retries fail, the original error is re-thrown.
478+
*/
479+
export async function withServerErrorRetry<T>(
480+
fn: () => Promise<T>
481+
): Promise<T> {
482+
const delays = [500, 1000, 2000];
483+
for (let attempt = 0; attempt <= delays.length; attempt++) {
484+
try {
485+
return await fn();
486+
} catch (err) {
487+
if (
488+
WorkflowAPIError.is(err) &&
489+
err.status !== undefined &&
490+
err.status >= 500 &&
491+
attempt < delays.length
492+
) {
493+
runtimeLogger.warn(
494+
'Server error (5xx) from workflow-server, retrying in-process',
495+
{
496+
status: err.status,
497+
attempt: attempt + 1,
498+
maxRetries: delays.length,
499+
nextDelayMs: delays[attempt],
500+
url: err.url,
501+
}
502+
);
503+
await new Promise((resolve) => setTimeout(resolve, delays[attempt]));
504+
continue;
505+
}
506+
throw err;
507+
}
508+
}
509+
throw new Error('withServerErrorRetry: unreachable');
510+
}

0 commit comments

Comments
 (0)