Skip to content

Commit 2a446af

Browse files
[core] Exclude inline step execution from replay timeout (#2013)
* [core] Exclude inline step execution from replay timeout The v5 combined workflow+step handler wraps inline step bodies in the same setTimeout(..., REPLAY_TIMEOUT_MS) guard that previously only bounded the v4 'workflows' function's fast deterministic replay. As a result, any workflow with a single step exceeding 240s hard-fails with FatalError: Workflow replay exceeded maximum duration (240s) after 4 attempts — even though the step could legitimately run for the full function maxDuration (up to 800s on Pro Fluid). Replace the setTimeout guard with a per-invocation budget that only accumulates non-step time. pauseReplayBudget() / resumeReplayBudget() bracket each executeStep() call, and the loop checks the budget at iteration boundaries. The retry-then-fail semantics from #1567 are preserved verbatim for the pure-replay case. Also adds a WORKFLOW_REPLAY_TIMEOUT_MS env var override (clamped to 30s..780s) so operators can adjust the bare-replay ceiling without patching @workflow/core. Fixes #2009. * Address PR review feedback - Extract budget bookkeeping into ReplayBudget class (replay-budget.ts) with sentinel-protected idempotent pause()/resume() to avoid double-counting in future refactors that nest step execution - Restore VERCEL_URL gate around process.exit(1) so a long pure-replay in local dev/non-Vercel runtimes can't hard-kill the host process - Warn (once per distinct raw value) when WORKFLOW_REPLAY_TIMEOUT_MS is clamped or rejected, so misconfiguration is observable - Correct Hobby maxDuration comment (60s standard / 300s Fluid) - Document budget-check responsiveness trade-off vs. old setTimeout - Tighten describe-error test assertions to match the full new hint - Shorten changeset description - Add ReplayBudget unit tests (9) including 8-minute step regression - Add warn-once tests for getReplayTimeoutMs (extended) * Replace VERCEL_URL gate with World capability Per review feedback, gating runtime behavior on process.env.VERCEL_URL leaks deployment-environment concerns into @workflow/core. Replace the check with a new optional capability on the World interface: processExitTriggersQueueRedelivery (default false). - @workflow/world: declare the new optional capability on World - @workflow/world-vercel: set it to true (Vercel fails the function invocation on non-zero exit and VQS redelivers via fresh invocation) - @workflow/core: handleReplayBudgetExhausted reads world.processExitTriggersQueueRedelivery instead of process.env.VERCEL_URL; behavior is otherwise unchanged - Add 4 unit tests for handleReplayBudgetExhausted covering both branches (exit-for-redelivery and best-effort run_failed) --------- Co-authored-by: Peter Wielander <mittgfu@gmail.com>
1 parent 3c50f8c commit 2a446af

10 files changed

Lines changed: 840 additions & 115 deletions

File tree

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
---
2+
"@workflow/core": patch
3+
"@workflow/world": patch
4+
"@workflow/world-vercel": patch
5+
---
6+
7+
Exclude inline step execution from the workflow replay timeout. Long-running steps no longer hit `REPLAY_TIMEOUT` (fixes #2009). Adds a `WORKFLOW_REPLAY_TIMEOUT_MS` env var override and a new optional `World.processExitTriggersQueueRedelivery` capability used to gate the runtime's `process.exit(1)` failure path.

‎packages/core/src/describe-error.test.ts‎

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -83,7 +83,9 @@ describe('describeError', () => {
8383
const result = describeError(undefined, RUN_ERROR_CODES.REPLAY_TIMEOUT);
8484
expect(result.attribution).toBe('sdk');
8585
expect(result.errorCode).toBe(RUN_ERROR_CODES.REPLAY_TIMEOUT);
86-
expect(result.hint).toContain('replay took too long');
86+
expect(result.hint).toContain(
87+
'replay between step boundaries took too long'
88+
);
8789
});
8890

8991
test('MAX_DELIVERIES_EXCEEDED via precomputed errorCode is attributed to the SDK', () => {
@@ -161,7 +163,9 @@ describe('describeRunError', () => {
161163
errorCode: RUN_ERROR_CODES.REPLAY_TIMEOUT,
162164
});
163165
expect(result.attribution).toBe('sdk');
164-
expect(result.hint).toContain('replay took too long');
166+
expect(result.hint).toContain(
167+
'replay between step boundaries took too long'
168+
);
165169
});
166170

167171
test('CORRUPTED_EVENT_LOG errorCode is attributed to the SDK', () => {
@@ -290,7 +294,7 @@ describe('describeError — payload shape snapshots', () => {
290294
{
291295
"attribution": "sdk",
292296
"errorCode": "REPLAY_TIMEOUT",
293-
"hint": "The workflow replay took too long. This usually means the event log is unusually large or the workflow function is doing heavy synchronous work between step boundaries.",
297+
"hint": "The workflow replay between step boundaries took too long. This bounds workflow-VM and event-log replay time only — step bodies (\`"use step"\` functions) are excluded. This usually means the event log is unusually large or the workflow function is doing heavy synchronous work in workflow code outside of step bodies. Override the default budget via the WORKFLOW_REPLAY_TIMEOUT_MS env var if needed.",
294298
}
295299
`);
296300
});

‎packages/core/src/describe-error.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -80,7 +80,7 @@ const RUNTIME_ERROR_HINT =
8080
const CORRUPTED_EVENT_LOG_HINT =
8181
'The workflow event log contains orphaned or mismatched events and cannot be replayed. This is an internal workflow SDK error; please report it with the runId.';
8282
const REPLAY_TIMEOUT_HINT =
83-
'The workflow replay took too long. This usually means the event log is unusually large or the workflow function is doing heavy synchronous work between step boundaries.';
83+
'The workflow replay between step boundaries took too long. This bounds workflow-VM and event-log replay time only — step bodies (`"use step"` functions) are excluded. This usually means the event log is unusually large or the workflow function is doing heavy synchronous work in workflow code outside of step bodies. Override the default budget via the WORKFLOW_REPLAY_TIMEOUT_MS env var if needed.';
8484
const MAX_DELIVERIES_HINT =
8585
'The workflow queue exceeded its max-delivery budget. This usually indicates a persistent runtime failure — check the most recent stack traces for the underlying cause.';
8686
const WORLD_CONTRACT_HINT =

‎packages/core/src/runtime.ts‎

Lines changed: 89 additions & 104 deletions
Original file line numberDiff line numberDiff line change
@@ -19,11 +19,7 @@ import { classifyRunError, isWorldContractError } from './classify-error.js';
1919
import { describeError } from './describe-error.js';
2020
import { WorkflowSuspension } from './global.js';
2121
import { runtimeLogger } from './logger.js';
22-
import {
23-
MAX_QUEUE_DELIVERIES,
24-
REPLAY_TIMEOUT_MAX_RETRIES,
25-
REPLAY_TIMEOUT_MS,
26-
} from './runtime/constants.js';
22+
import { MAX_QUEUE_DELIVERIES } from './runtime/constants.js';
2723
import {
2824
getQueueOverhead,
2925
getWorkflowQueueName,
@@ -34,6 +30,10 @@ import {
3430
queueMessage,
3531
withHealthCheck,
3632
} from './runtime/helpers.js';
33+
import {
34+
handleReplayBudgetExhausted,
35+
ReplayBudget,
36+
} from './runtime/replay-budget.js';
3737
import { executeStep } from './runtime/step-executor.js';
3838
import { handleSuspension } from './runtime/suspension-handler.js';
3939
import {
@@ -284,84 +284,31 @@ export function workflowEntrypoint(
284284

285285
const spanLinks = await linkToCurrentContext();
286286

287-
// --- Replay timeout guard ---
288-
// If the replay takes longer than the timeout, fail the run and exit.
289-
// This must be lower than the function's maxDuration to ensure
290-
// the failure is recorded before the platform kills the function.
291-
let replayTimeout: NodeJS.Timeout | undefined;
292-
if (process.env.VERCEL_URL !== undefined) {
293-
replayTimeout = setTimeout(async () => {
294-
// Allow a few retries before permanently failing the run.
295-
// On early attempts, just exit so the queue retries the message.
296-
if (metadata.attempt <= REPLAY_TIMEOUT_MAX_RETRIES) {
297-
runLogger.warn(
298-
'Workflow replay exceeded timeout but will be re-attempted (attempt < maxRetries)',
299-
{
300-
timeoutMs: REPLAY_TIMEOUT_MS,
301-
attempt: metadata.attempt,
302-
maxRetries: REPLAY_TIMEOUT_MAX_RETRIES,
303-
}
304-
);
305-
process.exit(1);
306-
}
307-
308-
const replayTimeoutDescription = describeError(
309-
undefined,
310-
RUN_ERROR_CODES.REPLAY_TIMEOUT
311-
);
312-
runLogger.error(
313-
'Workflow replay exceeded timeout and max retries exceeded. Failing the run',
314-
{
315-
timeoutMs: REPLAY_TIMEOUT_MS,
316-
attempt: metadata.attempt,
317-
maxRetries: REPLAY_TIMEOUT_MAX_RETRIES,
318-
errorCode: replayTimeoutDescription.errorCode,
319-
errorAttribution: replayTimeoutDescription.attribution,
320-
}
321-
);
322-
323-
try {
324-
const world = await getWorld();
325-
const getEncryptionKey = memoizeEncryptionKey(world, runId);
326-
const timeoutErr = new FatalError(
327-
`Workflow replay exceeded maximum duration (${REPLAY_TIMEOUT_MS / 1000}s) after ${metadata.attempt} attempts`
328-
);
329-
await world.events.create(
330-
runId,
331-
{
332-
eventType: 'run_failed',
333-
specVersion: SPEC_VERSION_CURRENT,
334-
eventData: {
335-
error: await dehydrateRunError(
336-
timeoutErr,
337-
runId,
338-
await getEncryptionKey()
339-
),
340-
errorCode: RUN_ERROR_CODES.REPLAY_TIMEOUT,
341-
},
342-
},
343-
{ requestId }
344-
);
345-
} catch (err) {
346-
// Best effort — process exits regardless. Surface why so
347-
// operators can diagnose repeat timeouts against the backend.
348-
runLogger.warn(
349-
'Unable to mark run as failed. The queue will continue to retry',
350-
{
351-
attempt: metadata.attempt,
352-
errorName: err instanceof Error ? err.name : 'UnknownError',
353-
errorMessage:
354-
err instanceof Error ? err.message : String(err),
355-
errorStack: err instanceof Error ? err.stack : undefined,
356-
}
357-
);
358-
}
359-
// Note that this also prevents the runtime from acking the queue message,
360-
// so the queue will call back once, after which a 410 will get it to exit early.
361-
process.exit(1);
362-
}, REPLAY_TIMEOUT_MS);
363-
replayTimeout.unref();
364-
}
287+
// --- Replay budget bookkeeping ---
288+
// The replay budget bounds the *non-step* portion of a single
289+
// handler invocation: deterministic event-log replay, workflow-VM
290+
// execution between step boundaries, suspension handling, queue
291+
// round-trips, etc. Inline step bodies (`"use step"` functions
292+
// invoked via `executeStep`) are intentionally excluded — they are
293+
// bounded by the platform's function `maxDuration` and the
294+
// `NO_INLINE_REPLAY_AFTER_MS` early-return guard below.
295+
//
296+
// The budget is checked at loop boundaries (top of each `while`
297+
// iteration). Note this is *less responsive* than the old
298+
// `setTimeout`-based approach: a single pathological `runWorkflow`
299+
// call processing a huge event log can overshoot the budget by up
300+
// to one iteration before bailing. In practice the headroom built
301+
// into `MAX_REPLAY_TIMEOUT_MS` (and the platform `maxDuration`
302+
// SIGTERM as ultimate backstop) gives us slack — the previous
303+
// `setTimeout` approach also relied on the platform kill as the
304+
// hard backstop. Do *not* "fix" this by adding a `setInterval`;
305+
// it would risk the same bug we just removed (bounding step
306+
// bodies).
307+
//
308+
// Earlier versions (pre-#2009 fix) used a single `setTimeout`
309+
// that also bounded step bodies, which broke any workflow with a
310+
// single step longer than the budget.
311+
const replayBudget = new ReplayBudget();
365312

366313
return await withTraceContext(traceContext, async () => {
367314
return await withWorkflowBaggage(
@@ -418,14 +365,24 @@ export function workflowEntrypoint(
418365
const bgStartedAt = bgRun.startedAt
419366
? +bgRun.startedAt
420367
: Date.now();
421-
const stepResult = await executeStep({
422-
world,
423-
workflowRunId: runId,
424-
workflowName,
425-
workflowStartedAt: bgStartedAt,
426-
stepId: incomingStepId,
427-
stepName: incomingStepName,
428-
});
368+
// Pause the replay budget while the step body runs —
369+
// step duration is bounded by the platform's function
370+
// maxDuration, not by the replay timeout. See the
371+
// ReplayBudget docs for the contract.
372+
replayBudget.pause();
373+
let stepResult: Awaited<ReturnType<typeof executeStep>>;
374+
try {
375+
stepResult = await executeStep({
376+
world,
377+
workflowRunId: runId,
378+
workflowName,
379+
workflowStartedAt: bgStartedAt,
380+
stepId: incomingStepId,
381+
stepName: incomingStepName,
382+
});
383+
} finally {
384+
replayBudget.resume();
385+
}
429386
if (stepResult.type === 'retry') {
430387
return { timeoutSeconds: stepResult.timeoutSeconds };
431388
}
@@ -665,6 +622,26 @@ export function workflowEntrypoint(
665622
while (true) {
666623
loopIteration++;
667624

625+
// Replay-budget check: bail out (retry or fail) if
626+
// non-step time within this invocation has exceeded
627+
// the configured budget. Step bodies are excluded
628+
// because replayBudget.pause()/resume() bracket every
629+
// `executeStep` call.
630+
if (replayBudget.isExhausted()) {
631+
await handleReplayBudgetExhausted({
632+
runId,
633+
workflowName,
634+
requestId,
635+
attempt: metadata.attempt,
636+
limitMs: replayBudget.configuredLimitMs,
637+
});
638+
// On Vercel, handleReplayBudgetExhausted always
639+
// exits the process. On local dev it returns; we
640+
// fall through and the request ends normally
641+
// (run_failed has been written best-effort).
642+
return;
643+
}
644+
668645
// Check timeout before replay
669646
if (
670647
Date.now() - invocationStartTime >=
@@ -1074,15 +1051,27 @@ export function workflowEntrypoint(
10741051
return;
10751052
}
10761053

1077-
// Execute inline step
1078-
const stepResult = await executeStep({
1079-
world,
1080-
workflowRunId: runId,
1081-
workflowName,
1082-
workflowStartedAt,
1083-
stepId: inlineStep.correlationId,
1084-
stepName: inlineStep.stepName,
1085-
});
1054+
// Execute inline step. Pause the replay budget
1055+
// for the duration of the step body — step
1056+
// duration is bounded by the platform's function
1057+
// maxDuration, not by the replay timeout. Without
1058+
// this the replay-budget check at the top of the
1059+
// next loop iteration would (incorrectly) charge
1060+
// the step body against the budget.
1061+
replayBudget.pause();
1062+
let stepResult: Awaited<ReturnType<typeof executeStep>>;
1063+
try {
1064+
stepResult = await executeStep({
1065+
world,
1066+
workflowRunId: runId,
1067+
workflowName,
1068+
workflowStartedAt,
1069+
stepId: inlineStep.correlationId,
1070+
stepName: inlineStep.stepName,
1071+
});
1072+
} finally {
1073+
replayBudget.resume();
1074+
}
10861075

10871076
if (stepResult.type === 'retry') {
10881077
// Step needs retry — queue self with stepId for retry
@@ -1274,10 +1263,6 @@ export function workflowEntrypoint(
12741263
); // End trace
12751264
}
12761265
); // End withWorkflowBaggage
1277-
}).finally(() => {
1278-
if (replayTimeout) {
1279-
clearTimeout(replayTimeout);
1280-
}
12811266
}); // End withTraceContext
12821267
}
12831268
);

0 commit comments

Comments
 (0)