Skip to content
Closed
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
8 changes: 8 additions & 0 deletions .changeset/enforce-strict-concurrency.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
---
"@workflow/world-vercel": minor
"@workflow/builders": minor
"@workflow/next": minor
"@workflow/sveltekit": patch
---

Add opt-in `WORKFLOW_SEQUENTIAL_REPLAYS` env var (also enabled by the `WORKFLOW_SAFE_MODE=1` umbrella flag when not set explicitly). When set to `1`, flow (orchestrator) routes are limited to one invocation per run at a time via a per-run queue topic and `maxConcurrency: 1` on the flow trigger. Step routes are unaffected.
17 changes: 17 additions & 0 deletions docs/content/docs/deploying/world/vercel-world.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,23 @@ Vercel team ID for API requests. Automatically detected.

Custom base URL for the Vercel workflow API. Automatically detected.

### `WORKFLOW_SEQUENTIAL_REPLAYS`

Set `WORKFLOW_SEQUENTIAL_REPLAYS=1` to guarantee that **at most one orchestrator (flow) invocation runs at a time per workflow run**. This behavior is off by default; without it, the runtime relies on idempotency and the event log to tolerate concurrent flow invocations of the same run. It is also enabled by `WORKFLOW_SAFE_MODE=1` when `WORKFLOW_SEQUENTIAL_REPLAYS` is not set explicitly.

When enabled, each run is given its own queue topic and the flow route is configured with `maxConcurrency: 1`, so [Vercel Queues](https://vercel.com/docs/queues) processes flow messages for a given run strictly one at a time. Step routes are unaffected and continue to run with full concurrency.

<Callout type="warn">
This variable is read at **both build time and runtime**, so it must be set as a project-level environment variable that applies to your build and your deployed functions. Setting it for only one will produce an inconsistent configuration. The same applies to framework integrations that write their own queue trigger configuration instead of using `getWorkflowQueueTrigger()` from `@workflow/builders`: they only get the runtime half (per-run topics) unless they also emit `maxConcurrency: 1` on their flow trigger.

Enabling sequential replays has a cost. Per [Vercel Queues pricing](https://vercel.com/docs/queues/pricing), push deliveries under `maxConcurrency` are billed at **2x units** for that operation, so every flow-route delivery costs double while this is enabled. It also creates one queue topic per run, which increases the number of distinct queues surfaced in queue observability, and each flow invocation waits for a per-run concurrency slot before delivery, which can add queueing latency. Leave it off unless you specifically need the per-run serialization guarantee.


While a replay holds a run's slot — including time spent executing steps inline — other wake messages for that run (hook resumes, aborts and cancellations, and run-timeout enforcement) wait for the slot. Expect aborts and timeouts to be delayed by up to the duration of the longest single invocation.

The guarantee covers messages sent by the Workflow SDK itself. External producers that compute a flow topic name directly (rather than enqueueing through the SDK) still deliver, but bypass the per-run serialization slot.
</Callout>

### Programmatic configuration

{/*@skip-typecheck: incomplete code sample*/}
Expand Down
9 changes: 7 additions & 2 deletions docs/content/docs/how-it-works/framework-integrations.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -405,12 +405,17 @@ Two queue topics are created per deployment:
| `step.func` | `__wkf_step_*` | Step execution (long-running, `maxDuration: max`) |
| `flow.func` | `__wkf_workflow_*` | Workflow orchestration (`maxDuration: 60`) |

If you're building a framework integration that targets Vercel, you should write these triggers into the `.vc-config.json` for each generated function. The `STEP_QUEUE_TRIGGER` and `WORKFLOW_QUEUE_TRIGGER` constants are exported from `@workflow/builders` for this purpose:
If you're building a framework integration that targets Vercel, you should write these triggers into the `.vc-config.json` for each generated function. Use `getWorkflowQueueTrigger()` for flow functions so `WORKFLOW_SEQUENTIAL_REPLAYS=1` is reflected in the generated trigger configuration (it also accepts a `namespace` option, matching `createWorkflowQueueTrigger`); `STEP_QUEUE_TRIGGER` is exported for step functions:

```typescript
import { STEP_QUEUE_TRIGGER, WORKFLOW_QUEUE_TRIGGER } from "@workflow/builders";
import { getWorkflowQueueTrigger, STEP_QUEUE_TRIGGER } from "@workflow/builders";

const flowTriggers = [getWorkflowQueueTrigger()];
const stepTriggers = [STEP_QUEUE_TRIGGER];
```

If your integration constructs the flow trigger object itself instead of calling `getWorkflowQueueTrigger()`, it must add `maxConcurrency: 1` to that trigger when sequential replays are enabled at build time (`WORKFLOW_SEQUENTIAL_REPLAYS=1`, or `WORKFLOW_SAFE_MODE=1` when the specific variable is unset — the exported `isSequentialReplaysEnabled()` helper implements this check). The runtime half of the feature (per-run queue topics) activates from the environment variable alone — without the trigger half, those per-run topics are not serialized and the setting only adds queue-topic cardinality.


### Custom implementations

Expand Down
1 change: 1 addition & 0 deletions packages/builders/src/base-builder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1710,6 +1710,7 @@ export const OPTIONS = handler;`;
topic: string;
consumer: string;
maxDeliveries?: number;
maxConcurrency?: number;
retryAfterSeconds?: number;
initialDelaySeconds?: number;
}>;
Expand Down
80 changes: 79 additions & 1 deletion packages/builders/src/constants.test.ts
Original file line number Diff line number Diff line change
@@ -1,9 +1,87 @@
import { afterEach, describe, expect, it } from 'vitest';
import { afterEach, beforeEach, describe, expect, it } from 'vitest';

import {
createWorkflowEntrypointOptionsCode,
createWorkflowQueueTrigger,
getWorkflowQueueTrigger,
} from './constants.js';

describe('getWorkflowQueueTrigger', () => {
let originalStrict: string | undefined;
let originalSafeMode: string | undefined;

beforeEach(() => {
originalStrict = process.env.WORKFLOW_SEQUENTIAL_REPLAYS;
originalSafeMode = process.env.WORKFLOW_SAFE_MODE;
delete process.env.WORKFLOW_SAFE_MODE;
});

afterEach(() => {
if (originalStrict !== undefined) {
process.env.WORKFLOW_SEQUENTIAL_REPLAYS = originalStrict;
} else {
delete process.env.WORKFLOW_SEQUENTIAL_REPLAYS;
}
if (originalSafeMode !== undefined) {
process.env.WORKFLOW_SAFE_MODE = originalSafeMode;
} else {
delete process.env.WORKFLOW_SAFE_MODE;
}
});

it('omits maxConcurrency by default', () => {
delete process.env.WORKFLOW_SEQUENTIAL_REPLAYS;
const trigger = getWorkflowQueueTrigger();
expect(trigger.topic).toBe('__wkf_workflow_*');
expect('maxConcurrency' in trigger).toBe(false);
});

it('sets maxConcurrency: 1 when WORKFLOW_SEQUENTIAL_REPLAYS=1', () => {
process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '1';
const trigger = getWorkflowQueueTrigger();
expect(trigger).toMatchObject({
topic: '__wkf_workflow_*',
maxConcurrency: 1,
});
});

it('does not set maxConcurrency for non-"1" values', () => {
process.env.WORKFLOW_SEQUENTIAL_REPLAYS = 'true';
const trigger = getWorkflowQueueTrigger();
expect('maxConcurrency' in trigger).toBe(false);
});

it('WORKFLOW_SAFE_MODE=1 sets maxConcurrency when the specific variable is unset', () => {
delete process.env.WORKFLOW_SEQUENTIAL_REPLAYS;
process.env.WORKFLOW_SAFE_MODE = '1';
expect(getWorkflowQueueTrigger()).toMatchObject({ maxConcurrency: 1 });
});

it('an explicit WORKFLOW_SEQUENTIAL_REPLAYS=0 wins over WORKFLOW_SAFE_MODE', () => {
process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '0';
process.env.WORKFLOW_SAFE_MODE = '1';
expect('maxConcurrency' in getWorkflowQueueTrigger()).toBe(false);
});

it('composes with an explicit namespace option', () => {
process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '1';
expect(getWorkflowQueueTrigger({ namespace: 'custom' })).toMatchObject({
topic: '__custom_wkf_workflow_*',
maxConcurrency: 1,
});
});

it('resolves WORKFLOW_QUEUE_NAMESPACE at call time', () => {
delete process.env.WORKFLOW_SEQUENTIAL_REPLAYS;
process.env.WORKFLOW_QUEUE_NAMESPACE = 'callns';
try {
expect(getWorkflowQueueTrigger().topic).toBe('__callns_wkf_workflow_*');
} finally {
delete process.env.WORKFLOW_QUEUE_NAMESPACE;
}
});
});

describe('createWorkflowQueueTrigger', () => {
afterEach(() => {
delete process.env.WORKFLOW_QUEUE_NAMESPACE;
Expand Down
42 changes: 42 additions & 0 deletions packages/builders/src/constants.ts
Original file line number Diff line number Diff line change
Expand Up @@ -121,3 +121,45 @@ export const OPTIONS = POST;`;
* Default queue trigger (no namespace). Backward compatible.
*/
export const WORKFLOW_QUEUE_TRIGGER = createWorkflowQueueTrigger();

/**
* Returns the queue trigger configuration for workflow (flow) routes.
*
* Builds on `createWorkflowQueueTrigger()` — the namespace comes from
* `options` or `WORKFLOW_QUEUE_NAMESPACE`, resolved at call time. When
* `WORKFLOW_SEQUENTIAL_REPLAYS` is enabled, sets `maxConcurrency: 1` so the
* queue processes at most one flow invocation per concrete topic at a time.
* Paired with the per-run physical topic naming in `@workflow/world-vercel`
* (which appends the run id to the flow topic), this enforces at most one
* orchestrator invocation per run. Step routes are intentionally excluded.
*
* Integrations that write their own flow trigger config instead of calling
* this must mirror the conditional `maxConcurrency: 1` themselves — the
* runtime half (per-run topics) activates from the env var alone, and without
* the trigger half those topics are not serialized.
*
* Must be read at build time, where the env var gates what is written into
* the route's `experimentalTriggers` config.
*/
/**
* Whether sequential replays are enabled: `WORKFLOW_SEQUENTIAL_REPLAYS=1`,
* or `WORKFLOW_SAFE_MODE=1` when `WORKFLOW_SEQUENTIAL_REPLAYS` is not set
* explicitly (safe mode fills the default of every safety-over-performance
* flag; an explicit per-flag value always wins). Read at call time.
*/
export function isSequentialReplaysEnabled(): boolean {
const explicit = process.env.WORKFLOW_SEQUENTIAL_REPLAYS;
if (explicit !== undefined && explicit !== '') {
return explicit === '1';
}
return process.env.WORKFLOW_SAFE_MODE === '1';
}

export function getWorkflowQueueTrigger(options?: { namespace?: string }) {
return {
...createWorkflowQueueTrigger(options),
...(isSequentialReplaysEnabled() && {
maxConcurrency: 1,
}),
};
}
2 changes: 2 additions & 0 deletions packages/builders/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,8 @@ export {
createStepQueueTrigger,
createWorkflowEntrypointOptionsCode,
createWorkflowQueueTrigger,
getWorkflowQueueTrigger,
isSequentialReplaysEnabled,
STEP_QUEUE_TRIGGER,
WORKFLOW_QUEUE_TRIGGER,
} from './constants.js';
Expand Down
4 changes: 2 additions & 2 deletions packages/builders/src/vercel-build-output-api.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { copyFile, mkdir, writeFile } from 'node:fs/promises';
import { join, resolve } from 'node:path';
import { BaseBuilder } from './base-builder.js';
import { STEP_QUEUE_TRIGGER, WORKFLOW_QUEUE_TRIGGER } from './constants.js';
import { getWorkflowQueueTrigger, STEP_QUEUE_TRIGGER } from './constants.js';

export class VercelBuildOutputAPIBuilder extends BaseBuilder {
async build(): Promise<void> {
Expand Down Expand Up @@ -116,7 +116,7 @@ export class VercelBuildOutputAPIBuilder extends BaseBuilder {
await this.createPackageJson(workflowsFuncDir, 'commonjs');
await this.createVcConfig(workflowsFuncDir, {
maxDuration: 'max',
experimentalTriggers: [WORKFLOW_QUEUE_TRIGGER],
experimentalTriggers: [getWorkflowQueueTrigger()],
runtime: this.config.runtime,
});

Expand Down
19 changes: 15 additions & 4 deletions packages/core/e2e/event-log-race-repro.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -417,7 +417,10 @@ async function describeStuckRun(
// settled yet. `completed` past the poll budget is downgraded to a non-gating
// SLOW_COMPLETION (slow, not wedged); `failed`/`cancelled` keep their meaning.
function classifyTerminalRun(
runData: { status: string; errorCode?: string },
runData: {
status: string;
error?: { code?: string; message?: string };
},
context: {
runId: string;
scenario: Scenario;
Expand Down Expand Up @@ -448,10 +451,16 @@ function classifyTerminalRun(
}

if (runData.status === 'failed') {
// A failed WorkflowRun carries its reason in `error: { code, message }`
// — the run has no top-level `errorCode`. Reading the structured error
// is what lets us classify USER_ERROR/RUNTIME_ERROR/CORRUPTED_EVENT_LOG
// (vs. uncategorised `other`) and surface *why* it failed in the summary.
const structuredError = runData.error;
return {
...base,
outcome: classifyFailure(runData.errorCode),
errorCode: runData.errorCode,
outcome: classifyFailure(structuredError?.code),
errorCode: structuredError?.code,
errorMessage: structuredError?.message,
};
}

Expand Down Expand Up @@ -1019,7 +1028,9 @@ describe('event log race repro', () => {

test(
'event log races do not corrupt, stall, or take stale branches',
{ timeout: testTimeoutMs },
{
timeout: testTimeoutMs,
},
async () => {
const stepBiasedAttempts = Math.ceil(config.stepSleepRaceAttempts / 2);
const sleepBiasedAttempts = Math.floor(config.stepSleepRaceAttempts / 2);
Expand Down
4 changes: 2 additions & 2 deletions packages/next/src/builder-eager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ export async function getNextBuilderEager() {
const {
BaseBuilder: BaseBuilderClass,
STEP_QUEUE_TRIGGER,
WORKFLOW_QUEUE_TRIGGER,
getWorkflowQueueTrigger,
// biome-ignore lint/security/noGlobalEval: Need to use eval here to avoid TypeScript from transpiling the import statement into `require()`
} = (await eval(
'import("@workflow/builders")'
Expand Down Expand Up @@ -456,7 +456,7 @@ export async function getNextBuilderEager() {
},
workflows: {
maxDuration: 'max',
experimentalTriggers: [WORKFLOW_QUEUE_TRIGGER],
experimentalTriggers: [getWorkflowQueueTrigger()],
},
};

Expand Down
11 changes: 2 additions & 9 deletions packages/sveltekit/src/index.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import path from 'node:path';
import { getWorkflowQueueTrigger } from '@workflow/builders';
import fs from 'fs-extra';

import { SvelteKitBuilder } from './builder.js';
Expand All @@ -21,15 +22,7 @@ process.on('beforeExit', () => {
file: '.vercel/output/functions/.well-known/workflow/v1/flow.func/.vc-config.json',
config: {
maxDuration: 'max',
experimentalTriggers: [
{
type: 'queue/v2beta',
topic: '__wkf_workflow_*',
consumer: 'default',
retryAfterSeconds: 5,
initialDelaySeconds: 0,
},
],
experimentalTriggers: [getWorkflowQueueTrigger()],
},
},
{
Expand Down
Loading