From 693b8c0b58db5b572ec3e2b351bd3c94828b355a Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Mon, 13 Jul 2026 21:07:59 +0000 Subject: [PATCH] [world-vercel] [builders] Add WORKFLOW_SEQUENTIAL_REPLAYS option to limit flow route concurrency to one (#2193) Signed-off-by: Peter Wielander --- .changeset/enforce-strict-concurrency.md | 8 + .../docs/deploying/world/vercel-world.mdx | 17 ++ .../how-it-works/framework-integrations.mdx | 9 +- packages/builders/src/base-builder.ts | 1 + packages/builders/src/constants.test.ts | 80 ++++++- packages/builders/src/constants.ts | 42 ++++ packages/builders/src/index.ts | 2 + .../builders/src/vercel-build-output-api.ts | 4 +- .../core/e2e/event-log-race-repro.test.ts | 19 +- packages/next/src/builder-eager.ts | 4 +- packages/sveltekit/src/index.ts | 11 +- packages/world-vercel/src/queue.test.ts | 225 ++++++++++++++++++ packages/world-vercel/src/queue.ts | 89 ++++++- 13 files changed, 490 insertions(+), 21 deletions(-) create mode 100644 .changeset/enforce-strict-concurrency.md diff --git a/.changeset/enforce-strict-concurrency.md b/.changeset/enforce-strict-concurrency.md new file mode 100644 index 0000000000..dea61c3a7c --- /dev/null +++ b/.changeset/enforce-strict-concurrency.md @@ -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. diff --git a/docs/content/docs/deploying/world/vercel-world.mdx b/docs/content/docs/deploying/world/vercel-world.mdx index 8746c76f19..eaf69f3e40 100644 --- a/docs/content/docs/deploying/world/vercel-world.mdx +++ b/docs/content/docs/deploying/world/vercel-world.mdx @@ -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. + + + 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. + + ### Programmatic configuration {/*@skip-typecheck: incomplete code sample*/} diff --git a/docs/content/docs/how-it-works/framework-integrations.mdx b/docs/content/docs/how-it-works/framework-integrations.mdx index 487538670c..cc64459a73 100644 --- a/docs/content/docs/how-it-works/framework-integrations.mdx +++ b/docs/content/docs/how-it-works/framework-integrations.mdx @@ -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 diff --git a/packages/builders/src/base-builder.ts b/packages/builders/src/base-builder.ts index 31f50ce81a..8f037f28c5 100644 --- a/packages/builders/src/base-builder.ts +++ b/packages/builders/src/base-builder.ts @@ -1710,6 +1710,7 @@ export const OPTIONS = handler;`; topic: string; consumer: string; maxDeliveries?: number; + maxConcurrency?: number; retryAfterSeconds?: number; initialDelaySeconds?: number; }>; diff --git a/packages/builders/src/constants.test.ts b/packages/builders/src/constants.test.ts index 3363ec9b14..9c88c7837d 100644 --- a/packages/builders/src/constants.test.ts +++ b/packages/builders/src/constants.test.ts @@ -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; diff --git a/packages/builders/src/constants.ts b/packages/builders/src/constants.ts index 6ca43d847d..fa7534f685 100644 --- a/packages/builders/src/constants.ts +++ b/packages/builders/src/constants.ts @@ -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, + }), + }; +} diff --git a/packages/builders/src/index.ts b/packages/builders/src/index.ts index 8b30ecbd2f..af53230cb2 100644 --- a/packages/builders/src/index.ts +++ b/packages/builders/src/index.ts @@ -15,6 +15,8 @@ export { createStepQueueTrigger, createWorkflowEntrypointOptionsCode, createWorkflowQueueTrigger, + getWorkflowQueueTrigger, + isSequentialReplaysEnabled, STEP_QUEUE_TRIGGER, WORKFLOW_QUEUE_TRIGGER, } from './constants.js'; diff --git a/packages/builders/src/vercel-build-output-api.ts b/packages/builders/src/vercel-build-output-api.ts index d2c6a7838a..9525e119ea 100644 --- a/packages/builders/src/vercel-build-output-api.ts +++ b/packages/builders/src/vercel-build-output-api.ts @@ -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 { @@ -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, }); diff --git a/packages/core/e2e/event-log-race-repro.test.ts b/packages/core/e2e/event-log-race-repro.test.ts index 515d08e292..3be145ae48 100644 --- a/packages/core/e2e/event-log-race-repro.test.ts +++ b/packages/core/e2e/event-log-race-repro.test.ts @@ -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; @@ -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, }; } @@ -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); diff --git a/packages/next/src/builder-eager.ts b/packages/next/src/builder-eager.ts index 432c458347..4b2b845d63 100644 --- a/packages/next/src/builder-eager.ts +++ b/packages/next/src/builder-eager.ts @@ -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")' @@ -456,7 +456,7 @@ export async function getNextBuilderEager() { }, workflows: { maxDuration: 'max', - experimentalTriggers: [WORKFLOW_QUEUE_TRIGGER], + experimentalTriggers: [getWorkflowQueueTrigger()], }, }; diff --git a/packages/sveltekit/src/index.ts b/packages/sveltekit/src/index.ts index 55ba9b6f70..12eab34b6d 100644 --- a/packages/sveltekit/src/index.ts +++ b/packages/sveltekit/src/index.ts @@ -1,4 +1,5 @@ import path from 'node:path'; +import { getWorkflowQueueTrigger } from '@workflow/builders'; import fs from 'fs-extra'; import { SvelteKitBuilder } from './builder.js'; @@ -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()], }, }, { diff --git a/packages/world-vercel/src/queue.test.ts b/packages/world-vercel/src/queue.test.ts index 3f6eb80377..cb5b31df4b 100644 --- a/packages/world-vercel/src/queue.test.ts +++ b/packages/world-vercel/src/queue.test.ts @@ -330,6 +330,231 @@ describe('createQueue', () => { }); }); + describe('strict concurrency (WORKFLOW_SEQUENTIAL_REPLAYS)', () => { + let originalDeploymentId: string | undefined; + let originalStrict: string | undefined; + let originalSafeMode: string | undefined; + + beforeEach(() => { + originalDeploymentId = process.env.VERCEL_DEPLOYMENT_ID; + originalStrict = process.env.WORKFLOW_SEQUENTIAL_REPLAYS; + originalSafeMode = process.env.WORKFLOW_SAFE_MODE; + delete process.env.WORKFLOW_SAFE_MODE; + process.env.VERCEL_DEPLOYMENT_ID = 'dpl_test'; + mockSend.mockResolvedValue({ messageId: 'msg-123' }); + }); + + afterEach(() => { + if (originalDeploymentId !== undefined) { + process.env.VERCEL_DEPLOYMENT_ID = originalDeploymentId; + } else { + delete process.env.VERCEL_DEPLOYMENT_ID; + } + 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('appends runId to the physical flow topic while keeping the logical queueName', async () => { + process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '1'; + + const queue = createQueue(); + await queue.queue('__wkf_workflow_test', { runId: 'wrun_abc' }); + + // send(physicalTopic, wrapper, options) + expect(mockSend.mock.calls[0][0]).toBe('__wkf_workflow_test_wrun_abc'); + // The logical queue name is preserved so the handler + re-enqueue path + // resolves the same per-run physical topic on the next invocation. + expect(mockSend.mock.calls[0][1].queueName).toBe('__wkf_workflow_test'); + }); + + it('re-enqueues delayed flow messages to the same per-run physical topic', async () => { + process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '1'; + + let capturedHandler: ( + message: unknown, + metadata: unknown + ) => Promise; + mockHandleCallback.mockImplementation((handler) => { + capturedHandler = handler; + return async () => new Response('ok'); + }); + + const queue = createQueue(); + queue.createQueueHandler('__wkf_workflow_', async () => ({ + timeoutSeconds: 300, + })); + + await capturedHandler!( + { + payload: { runId: 'wrun_abc' }, + queueName: '__wkf_workflow_test', + deploymentId: 'dpl_original', + }, + { messageId: 'msg-123', deliveryCount: 1, createdAt: new Date() } + ); + + expect(mockSend.mock.calls[0][0]).toBe('__wkf_workflow_test_wrun_abc'); + }); + + it('does not rewrite the topic when the flag is unset', async () => { + delete process.env.WORKFLOW_SEQUENTIAL_REPLAYS; + + const queue = createQueue(); + await queue.queue('__wkf_workflow_test', { runId: 'wrun_abc' }); + + expect(mockSend.mock.calls[0][0]).toBe('__wkf_workflow_test'); + }); + + it('does not rewrite step topics even when the flag is set', async () => { + process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '1'; + + const queue = createQueue(); + await queue.queue('__wkf_step_myStep', { + workflowName: 'test-workflow', + workflowRunId: 'wrun_abc', + workflowStartedAt: Date.now(), + stepId: 'step_xyz', + }); + + expect(mockSend.mock.calls[0][0]).toBe('__wkf_step_myStep'); + }); + + it('WORKFLOW_SAFE_MODE=1 routes to per-run topics when the specific variable is unset', async () => { + delete process.env.WORKFLOW_SEQUENTIAL_REPLAYS; + process.env.WORKFLOW_SAFE_MODE = '1'; + + const queue = createQueue(); + await queue.queue('__wkf_workflow_test', { runId: 'wrun_abc' }); + + expect(mockSend.mock.calls[0][0]).toBe('__wkf_workflow_test_wrun_abc'); + }); + + it('an explicit WORKFLOW_SEQUENTIAL_REPLAYS=0 wins over WORKFLOW_SAFE_MODE', async () => { + process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '0'; + process.env.WORKFLOW_SAFE_MODE = '1'; + + const queue = createQueue(); + await queue.queue('__wkf_workflow_test', { runId: 'wrun_abc' }); + + expect(mockSend.mock.calls[0][0]).toBe('__wkf_workflow_test'); + }); + + it('gives inline step executions (flow topic + stepId) a per-step topic for full parallelism', async () => { + process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '1'; + + const queue = createQueue(); + // Inline step executions ride the flow topic as WorkflowInvokePayload + // with a stepId. They must NOT share the per-run serialized topic, or + // a run's parallel steps would execute one at a time. + await queue.queue('__wkf_workflow_test', { + runId: 'wrun_abc', + stepId: 'step_one', + }); + await queue.queue('__wkf_workflow_test', { + runId: 'wrun_abc', + stepId: 'step_two', + }); + + expect(mockSend.mock.calls[0][0]).toBe( + '__wkf_workflow_test_wrun_abc_step_one' + ); + expect(mockSend.mock.calls[1][0]).toBe( + '__wkf_workflow_test_wrun_abc_step_two' + ); + // The wrapper keeps the logical queue name for handler dispatch. + expect(mockSend.mock.calls[0][1].queueName).toBe('__wkf_workflow_test'); + }); + + it('gives each health check its own physical topic so concurrent probes never serialize', async () => { + process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '1'; + + const queue = createQueue(); + // Concurrent probes: with maxConcurrency: 1 applied per concrete topic, + // a single shared `…_health_check` topic would process probes one at a + // time and let a slow probe time out its successors. Distinct + // per-correlation topics keep them independent. + await queue.queue('__wkf_workflow_health_check', { + __healthCheck: true as const, + correlationId: 'corr_123', + }); + await queue.queue('__wkf_workflow_health_check', { + __healthCheck: true as const, + correlationId: 'corr_456', + }); + + expect(mockSend.mock.calls[0][0]).toBe( + '__wkf_workflow_health_check_corr_123' + ); + expect(mockSend.mock.calls[1][0]).toBe( + '__wkf_workflow_health_check_corr_456' + ); + // The wrapper keeps the logical queue name for handler dispatch. + expect(mockSend.mock.calls[0][1].queueName).toBe( + '__wkf_workflow_health_check' + ); + }); + + it('does not rewrite health check topics when the flag is unset', async () => { + delete process.env.WORKFLOW_SEQUENTIAL_REPLAYS; + + const queue = createQueue(); + await queue.queue('__wkf_workflow_health_check', { + __healthCheck: true as const, + correlationId: 'corr_123', + }); + + expect(mockSend.mock.calls[0][0]).toBe('__wkf_workflow_health_check'); + }); + + it('does not rewrite step health check topics even when the flag is set', async () => { + process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '1'; + + const queue = createQueue(); + await queue.queue('__wkf_step_health_check', { + __healthCheck: true as const, + correlationId: 'corr_123', + }); + + expect(mockSend.mock.calls[0][0]).toBe('__wkf_step_health_check'); + }); + + it('appends runId to namespaced flow topics so it composes with WORKFLOW_QUEUE_NAMESPACE', async () => { + process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '1'; + + const queue = createQueue(); + await queue.queue('__custom_wkf_workflow_test', { runId: 'wrun_abc' }); + + expect(mockSend.mock.calls[0][0]).toBe( + '__custom_wkf_workflow_test_wrun_abc' + ); + expect(mockSend.mock.calls[0][1].queueName).toBe( + '__custom_wkf_workflow_test' + ); + }); + + it('does not rewrite namespaced step topics even when the flag is set', async () => { + process.env.WORKFLOW_SEQUENTIAL_REPLAYS = '1'; + + const queue = createQueue(); + await queue.queue('__custom_wkf_step_myStep', { + workflowName: 'test-workflow', + workflowRunId: 'wrun_abc', + workflowStartedAt: Date.now(), + stepId: 'step_xyz', + }); + + expect(mockSend.mock.calls[0][0]).toBe('__custom_wkf_step_myStep'); + }); + }); + describe('createQueueHandler()', () => { const setupHandler = ({ timeoutSeconds }: { timeoutSeconds: number }) => { let capturedHandler: ( diff --git a/packages/world-vercel/src/queue.ts b/packages/world-vercel/src/queue.ts index 8938865f7b..3f024921aa 100644 --- a/packages/world-vercel/src/queue.ts +++ b/packages/world-vercel/src/queue.ts @@ -186,6 +186,88 @@ function getHeadersFromPayload( return Object.keys(headers).length > 0 ? headers : undefined; } +/** + * Resolves the physical VQS topic for a message. + * + * Normally this is just the logical queue name. When + * `WORKFLOW_SEQUENTIAL_REPLAYS` is enabled, messages on flow (workflow) + * topics get a payload-dependent physical topic. VQS scopes `maxConcurrency` + * per concrete topic, so combined with `maxConcurrency: 1` on the flow + * trigger: + * + * - Orchestrator replays (`WorkflowInvokePayload` without a `stepId`) get a + * per-run topic — at most one replay per run at a time. + * - Inline step executions (`WorkflowInvokePayload` WITH a `stepId` — they + * ride the flow topic in the combined handler model) get a per-step topic + * so steps keep full parallelism across a run; only redeliveries of the + * same step serialize. + * - Health checks get a per-probe topic (their correlation id) so concurrent + * probes never queue behind one shared `…_health_check` slot. + * + * Legacy `*_wkf_step_*` topics are intentionally excluded. + * + * The flow-topic match allows an optional queue namespace prefix + * (`___wkf_workflow_`, see `@workflow/builders` constants) so the + * behavior composes with `WORKFLOW_QUEUE_NAMESPACE`. + * + * This rewrite only serializes messages sent through this adapter (the + * wrapper's logical `queueName` keeps handler dispatch and re-enqueues on the + * same physical topic). A producer that computes the shared topic name itself + * still delivers (the trigger subscribes with a wildcard), but bypasses the + * per-run concurrency slot. + */ +const FLOW_TOPIC_PATTERN = /^__([a-z][a-z0-9]*_)?wkf_workflow_/; + +let loggedSequentialReplays = false; + +/** + * 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). Mirrors + * `isSequentialReplaysEnabled` in `@workflow/builders` — world-vercel must + * not depend on the build-time package, so the check is duplicated. + */ +function isSequentialReplaysEnabled(): boolean { + const explicit = process.env.WORKFLOW_SEQUENTIAL_REPLAYS; + if (explicit !== undefined && explicit !== '') { + return explicit === '1'; + } + return process.env.WORKFLOW_SAFE_MODE === '1'; +} + +function getPhysicalQueueName( + queueName: ValidQueueName, + payload: QueuePayload +): string { + if (!isSequentialReplaysEnabled() || !FLOW_TOPIC_PATTERN.test(queueName)) { + return queueName; + } + if (!loggedSequentialReplays) { + loggedSequentialReplays = true; + // One-time breadcrumb so a half-applied configuration (env var set without + // a maxConcurrency-bearing flow trigger, or vice versa) is diagnosable + // from function logs. Must go to stderr: this code also runs inside CLI + // commands whose stdout is a machine-parsed JSON contract (e.g. + // `workflow health --json`). + console.warn( + '[workflow] WORKFLOW_SEQUENTIAL_REPLAYS=1: routing flow messages to per-run queue topics' + ); + } + if ('runId' in payload && typeof payload.runId === 'string') { + // Inline step execution: full parallelism via a per-step topic. + if ('stepId' in payload && typeof payload.stepId === 'string') { + return `${queueName}_${payload.runId}_${payload.stepId}`; + } + // Orchestrator replay: serialize per run. + return `${queueName}_${payload.runId}`; + } + if ('__healthCheck' in payload && typeof payload.correlationId === 'string') { + return `${queueName}_${payload.correlationId}`; + } + return queueName; +} + type QueueFunction = ( queueName: ValidQueueName, payload: QueuePayload, @@ -247,11 +329,16 @@ export function createQueue(config?: APIConfig): Queue { // preserving Uint8Array values (workflow input in specVersion >= 2). const wrapper = { payload, + // Keep the logical queue name so the handler and re-enqueue path + // resolve the same per-run physical topic on the next invocation. queueName, // Store deploymentId in the message so it can be preserved when re-enqueueing deploymentId: opts?.deploymentId, }; - const sanitizedQueueName = queueName.replace(/[^A-Za-z0-9-_]/g, '-'); + const sanitizedQueueName = getPhysicalQueueName(queueName, payload).replace( + /[^A-Za-z0-9-_]/g, + '-' + ); try { const { messageId } = await client.send(sanitizedQueueName, wrapper, { idempotencyKey: opts?.idempotencyKey,