Skip to content

Commit 897aac9

Browse files
pranaygpclaudeVaguelySerious
authored
fix(world-vercel): retry idempotent event POSTs in-process to avoid step re-execution (#2675)
* fix(world-vercel): retry idempotent event POSTs in-process to avoid step re-execution undici's RetryAgent never retries a POST, so a transient transport blip (UND_ERR_REQ_RETRY, ECONNRESET, socket/headers timeout, transient 5xx) when committing a step's terminal event bubbles out, the queue redelivers, and the step's user code re-executes with attempt++ even though it already ran to completion. workflow-server makes these writes idempotent in outcome: entity handlers run before the event-log row is inserted and state transitions are conditional writes excluding terminal states, so a retry whose original landed throws before any row is written and surfaces as a 409 the SDK already handles. This adds a bounded in-process retry (new event-retry.ts) gated by a validated per-event EVENT_RETRY_ELIGIBILITY map, excluding step_started (double-increments attempt), step_retrying (appends a duplicate row), and hook_received (no server guard). Complements #2666, which routes completion-persistence failures to queue redelivery instead of recording them as user failures; this avoids the redelivery (and re-execution) for transient blips. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> * fix(world-vercel): address review on event-POST retry - Don't retry external cancellation: drop AbortError from the transient set (keep self-timeout TimeoutError), so a caller-requested abort isn't re-issued or stalled by the backoff budget. - Add DEBUG-gated logging on each retry and on retry exhaustion so an in-process retry vs. a fall-through to queue redelivery is distinguishable in logs. - Clarify docs: a landed retry surfaces as 409 for most types, but run_started/attr_set return 200 success (not a 409). - Add createWorkflowRunEvent integration tests (events-retry.test.ts): eventType is threaded into the retry wrapper, and the 404->HookNotFoundError mapping still fires after the retry loop. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: Peter Wielander <mittgfu@gmail.com>
1 parent 1dcdafd commit 897aac9

5 files changed

Lines changed: 689 additions & 1 deletion

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
'@workflow/world-vercel': patch
3+
---
4+
5+
Retry transient transport failures (e.g. `UND_ERR_REQ_RETRY`, `ECONNRESET`, socket timeouts, 5xx) in-process for idempotent-on-retry event POSTs, so a brief network blip after a step completes no longer re-executes the step. `step_started`, `step_retrying`, and `hook_received` are excluded as they are not safe to blindly retry.
Lines changed: 228 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,228 @@
1+
import {
2+
EntityConflictError,
3+
RunExpiredError,
4+
ThrottleError,
5+
TooEarlyError,
6+
WorkflowWorldError,
7+
} from '@workflow/errors';
8+
import { EventTypeSchema } from '@workflow/world';
9+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
10+
import {
11+
EVENT_RETRY_ELIGIBILITY,
12+
isRetryableEventPostError,
13+
MAX_EVENT_POST_RETRIES,
14+
withEventPostRetry,
15+
} from './event-retry.js';
16+
17+
const transportErr = (code: string) =>
18+
Object.assign(new Error(`transport ${code}`), { code });
19+
20+
// undici `fetch` wraps low-level failures in a TypeError whose `cause` carries
21+
// the real code — the classifier must walk the cause chain.
22+
const fetchFailed = (code: string) =>
23+
Object.assign(new TypeError('fetch failed'), {
24+
cause: Object.assign(new Error('underlying'), { code }),
25+
});
26+
27+
describe('EVENT_RETRY_ELIGIBILITY', () => {
28+
it('classifies every world event type (no gaps)', () => {
29+
for (const type of EventTypeSchema.options) {
30+
const policy = EVENT_RETRY_ELIGIBILITY[type];
31+
expect(policy, `missing policy for ${type}`).toBeDefined();
32+
expect(typeof policy.retryable).toBe('boolean');
33+
expect(policy.reason.length).toBeGreaterThan(0);
34+
}
35+
});
36+
37+
it('marks idempotent-on-retry events retryable', () => {
38+
for (const type of [
39+
'run_created',
40+
'run_started',
41+
'run_completed',
42+
'run_failed',
43+
'run_cancelled',
44+
'attr_set',
45+
'step_created',
46+
'step_completed',
47+
'step_failed',
48+
'wait_created',
49+
'wait_completed',
50+
'hook_created',
51+
'hook_disposed',
52+
] as const) {
53+
expect(EVENT_RETRY_ELIGIBILITY[type].retryable, type).toBe(true);
54+
}
55+
});
56+
57+
it('excludes events that are unsafe to blindly retry', () => {
58+
// step_started double-increments attempt; step_retrying / hook_received
59+
// append a duplicate event row; hook_conflict is server-originated.
60+
for (const type of [
61+
'step_started',
62+
'step_retrying',
63+
'hook_received',
64+
'hook_conflict',
65+
] as const) {
66+
expect(EVENT_RETRY_ELIGIBILITY[type].retryable, type).toBe(false);
67+
}
68+
});
69+
});
70+
71+
describe('isRetryableEventPostError', () => {
72+
it('retries transient 5xx', () => {
73+
for (const status of [500, 502, 503, 504]) {
74+
expect(
75+
isRetryableEventPostError(new WorkflowWorldError('boom', { status }))
76+
).toBe(true);
77+
}
78+
});
79+
80+
it('does not retry 4xx definitive responses', () => {
81+
for (const status of [400, 404]) {
82+
expect(
83+
isRetryableEventPostError(new WorkflowWorldError('nope', { status }))
84+
).toBe(false);
85+
}
86+
});
87+
88+
it('does not retry server-considered conflicts/terminal/too-early/throttle', () => {
89+
expect(isRetryableEventPostError(new EntityConflictError('409'))).toBe(
90+
false
91+
);
92+
expect(isRetryableEventPostError(new RunExpiredError('410'))).toBe(false);
93+
expect(isRetryableEventPostError(new TooEarlyError('425'))).toBe(false);
94+
expect(isRetryableEventPostError(new ThrottleError('429'))).toBe(false);
95+
});
96+
97+
it('retries a body-parse failure (write may have landed)', () => {
98+
expect(
99+
isRetryableEventPostError(
100+
new WorkflowWorldError('parse', { code: 'PARSE_ERROR' })
101+
)
102+
).toBe(true);
103+
});
104+
105+
it('retries raw transport errors by code', () => {
106+
for (const code of [
107+
'ECONNRESET',
108+
'ETIMEDOUT',
109+
'UND_ERR_SOCKET',
110+
'UND_ERR_REQ_RETRY',
111+
'UND_ERR_HEADERS_TIMEOUT',
112+
]) {
113+
expect(isRetryableEventPostError(transportErr(code)), code).toBe(true);
114+
}
115+
});
116+
117+
it('retries transport errors hidden under a `fetch failed` cause', () => {
118+
expect(isRetryableEventPostError(fetchFailed('UND_ERR_SOCKET'))).toBe(true);
119+
});
120+
121+
it('retries our own timeout (TimeoutError) but not an external abort (AbortError)', () => {
122+
// Self-deadline via AbortSignal.timeout → TimeoutError → ambiguous, retry.
123+
expect(
124+
isRetryableEventPostError(
125+
Object.assign(new Error('timed out'), { name: 'TimeoutError' })
126+
)
127+
).toBe(true);
128+
// Caller-requested cancellation surfaces as AbortError → must NOT be retried.
129+
expect(
130+
isRetryableEventPostError(
131+
Object.assign(new Error('aborted'), { name: 'AbortError' })
132+
)
133+
).toBe(false);
134+
// Also when wrapped by makeRequest as a WorkflowWorldError(cause).
135+
expect(
136+
isRetryableEventPostError(
137+
new WorkflowWorldError('request aborted', {
138+
cause: Object.assign(new Error('aborted'), { name: 'AbortError' }),
139+
})
140+
)
141+
).toBe(false);
142+
});
143+
144+
it('does not retry an unclassified error', () => {
145+
expect(isRetryableEventPostError(new Error('something else'))).toBe(false);
146+
});
147+
});
148+
149+
describe('withEventPostRetry', () => {
150+
beforeEach(() => {
151+
vi.useFakeTimers();
152+
});
153+
afterEach(() => {
154+
vi.useRealTimers();
155+
});
156+
157+
it('retries a retryable event past a transient blip and returns the result', async () => {
158+
let calls = 0;
159+
const fn = vi.fn(async () => {
160+
calls++;
161+
if (calls === 1) throw transportErr('ECONNRESET');
162+
return 'ok';
163+
});
164+
165+
const p = withEventPostRetry(fn, 'step_completed');
166+
await vi.runAllTimersAsync();
167+
168+
await expect(p).resolves.toBe('ok');
169+
expect(fn).toHaveBeenCalledTimes(2);
170+
});
171+
172+
it('surfaces a 409 that appears on a retry (original landed)', async () => {
173+
let calls = 0;
174+
const fn = vi.fn(async () => {
175+
calls++;
176+
if (calls === 1) throw transportErr('ECONNRESET');
177+
// The first attempt actually landed; the retry observes the conflict.
178+
throw new EntityConflictError('already completed');
179+
});
180+
181+
const p = withEventPostRetry(fn, 'step_completed').catch((e) => e);
182+
await vi.runAllTimersAsync();
183+
184+
const err = await p;
185+
expect(EntityConflictError.is(err)).toBe(true);
186+
expect(fn).toHaveBeenCalledTimes(2);
187+
});
188+
189+
it('gives up after MAX_EVENT_POST_RETRIES and throws the last error', async () => {
190+
const fn = vi.fn(async () => {
191+
throw transportErr('ECONNRESET');
192+
});
193+
194+
const p = withEventPostRetry(fn, 'step_completed').catch((e) => e);
195+
await vi.runAllTimersAsync();
196+
197+
const err = await p;
198+
expect((err as { code?: string }).code).toBe('ECONNRESET');
199+
expect(fn).toHaveBeenCalledTimes(MAX_EVENT_POST_RETRIES + 1);
200+
});
201+
202+
it('does not retry a definitive failure even for a retryable event', async () => {
203+
const fn = vi.fn(async () => {
204+
throw new WorkflowWorldError('bad request', { status: 400 });
205+
});
206+
207+
await expect(withEventPostRetry(fn, 'step_completed')).rejects.toThrow();
208+
expect(fn).toHaveBeenCalledTimes(1);
209+
});
210+
211+
it('does not retry an excluded event (step_started — protects the attempt counter)', async () => {
212+
const fn = vi.fn(async () => {
213+
throw transportErr('ECONNRESET');
214+
});
215+
216+
await expect(withEventPostRetry(fn, 'step_started')).rejects.toThrow();
217+
expect(fn).toHaveBeenCalledTimes(1);
218+
});
219+
220+
it('does not retry an excluded event (hook_received — avoids double-delivery)', async () => {
221+
const fn = vi.fn(async () => {
222+
throw transportErr('ECONNRESET');
223+
});
224+
225+
await expect(withEventPostRetry(fn, 'hook_received')).rejects.toThrow();
226+
expect(fn).toHaveBeenCalledTimes(1);
227+
});
228+
});

0 commit comments

Comments
 (0)