Skip to content
Merged
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
7 changes: 7 additions & 0 deletions .changeset/queue-batch-fanout-dispatch.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
---
'@workflow/world-vercel': patch
'@workflow/world': patch
'@workflow/core': patch
---

Publish a fan-out's step-execution messages in one batched queue request instead of one per step, via a new optional `Queue.queueBatch` implemented on `@vercel/queue`'s `experimental_sendBatch`.
25 changes: 25 additions & 0 deletions docs/content/worlds/v5/building-a-world.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -229,13 +229,38 @@ interface Queue {
opts?: QueueOptions
): Promise<{ messageId: MessageId | null }>;

// Optional. Omit it and the runtime publishes one message at a time.
queueBatch?(
queueName: ValidQueueName,
messages: readonly { message: QueuePayload; opts?: QueueOptions }[]
): Promise<QueueBatchResult[]>;

createQueueHandler(
queueNamePrefix: QueuePrefix,
handler: (message: unknown, meta: { attempt: number; queueName: ValidQueueName; messageId: MessageId }) => Promise<void | { timeoutSeconds: number }>
): (req: Request) => Promise<Response>;
}
```

### Batched publishing

`queueBatch` is optional. Implement it when your transport can accept several messages in one round trip, and the runtime will use it to dispatch a wide `Promise.all` fan-out. A fan-out otherwise costs one round trip per branch, paid before the first branch's step body runs, so it lands directly on time-to-first-step.

Its contract differs from `queue` in one important way: **a message that fails is reported, not thrown.** Return one result per input message, in input order, where `error` is set on the entries that were rejected and `retryable` says whether republishing might succeed. Reject the promise only for a request-level failure, where you cannot say what happened to any individual message.

The count must match: return exactly as many results as you were given messages. The runtime rejects the whole batch if it does not, because an omitted result is indistinguishable from a message that was never published, and treating it as success would strand that step with nothing reporting an error.

{/* @skip-typecheck - interface definition, not runnable code */}
```typescript
type QueueBatchResult =
| { messageId: MessageId | null; error?: undefined }
| { messageId: null; error: string; retryable: boolean };
```

`messageId: null` with no `error` means accepted without an ID yet, exactly as for `queue`, so callers test `error === undefined` for success rather than a non-null `messageId`.

Every message targets one logical `queueName`, but you may split the batch internally — by a transport cap, or by any per-message routing dimension you derive from the payload — as long as the returned order still matches the input. The runtime passes a distinct `idempotencyKey` per message, because its recovery for both a request-level failure and a retryable entry is to republish the whole batch.

### Queue names

Queue names follow a specific pattern:
Expand Down
113 changes: 113 additions & 0 deletions packages/core/src/runtime/helpers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import {
memoizeEncryptionKey,
mergeReportedEvents,
preconditionEventDelta,
queueMessages,
SLOT_GAP_RECHECK_ATTEMPTS,
settleEventSlotGap,
slotSnapshotParams,
Expand Down Expand Up @@ -1068,3 +1069,115 @@ describe('health check run public key', () => {
expect(response.workflowCoreVersion).toBeDefined();
});
});

describe('queueMessages', () => {
const entries = (n: number) =>
Array.from({ length: n }, (_, i) => ({
message: { runId: 'wrun_1', stepId: `step-${i}` },
opts: { idempotencyKey: `key-${i}` },
}));

const makeWorld = (over: Partial<World>) =>
({
queue: vi.fn().mockResolvedValue({ messageId: null }),
...over,
}) as unknown as World;

it('uses the World batch send when available', async () => {
const queueBatch = vi
.fn()
.mockResolvedValue([{ messageId: 'a' }, { messageId: 'b' }]);
const world = makeWorld({ queueBatch });

await queueMessages(world, '__wkf_workflow_t', entries(2));

expect(queueBatch).toHaveBeenCalledTimes(1);
expect(queueBatch.mock.calls[0][1]).toHaveLength(2);
expect(world.queue).not.toHaveBeenCalled();
});

it('falls back to single sends on a World with no batch support', async () => {
const world = makeWorld({});

await queueMessages(world, '__wkf_workflow_t', entries(3));

expect(world.queue).toHaveBeenCalledTimes(3);
// Each fallback send keeps its own key and payload.
expect(vi.mocked(world.queue).mock.calls.map((c) => c[2])).toEqual([
{ idempotencyKey: 'key-0' },
{ idempotencyKey: 'key-1' },
{ idempotencyKey: 'key-2' },
]);
});

it('rejects when any entry failed, naming the shortfall', async () => {
const queueBatch = vi
.fn()
.mockResolvedValue([
{ messageId: 'a' },
{ messageId: null, error: 'rate limited', retryable: true },
]);
const world = makeWorld({ queueBatch });

await expect(
queueMessages(world, '__wkf_workflow_t', entries(2))
).rejects.toThrow(/Failed to publish 1 of 2/);
});

it('marks the rejection retryable only when a failed entry is', async () => {
const world = makeWorld({
queueBatch: vi
.fn()
.mockResolvedValue([
{ messageId: null, error: 'bad request', retryable: false },
]),
});

await expect(
queueMessages(world, '__wkf_workflow_t', entries(1))
).rejects.toMatchObject({ retryable: false });
});

it('treats a deferred acceptance (null id, no error) as success', async () => {
const world = makeWorld({
queueBatch: vi.fn().mockResolvedValue([{ messageId: null }]),
});

await expect(
queueMessages(world, '__wkf_workflow_t', entries(1))
).resolves.toBeUndefined();
});

it('rejects when the World returns fewer results than messages', async () => {
// A short array says nothing about the messages it omits. Reading it as
// success would ack the delivery with those steps never dispatched, and
// the run would stall with no error recorded anywhere.
const world = makeWorld({
queueBatch: vi.fn().mockResolvedValue([{ messageId: 'a' }]),
});

await expect(
queueMessages(world, '__wkf_workflow_t', entries(3))
).rejects.toThrow(/returned 1 result\(s\) for 3 message\(s\)/);
});

it('marks a short-result rejection retryable', async () => {
const world = makeWorld({
queueBatch: vi.fn().mockResolvedValue([]),
});

await expect(
queueMessages(world, '__wkf_workflow_t', entries(2))
).rejects.toMatchObject({ retryable: true });
});

it('does not touch the World for an empty message set', async () => {
const queueBatch = vi.fn();
const world = makeWorld({ queueBatch });

await queueMessages(world, '__wkf_workflow_t', []);

expect(queueBatch).not.toHaveBeenCalled();
expect(world.queue).not.toHaveBeenCalled();
});
});
82 changes: 82 additions & 0 deletions packages/core/src/runtime/helpers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1216,6 +1216,88 @@ export async function queueMessage(
);
}

/**
* Publishes several messages to one logical queue, using the World's batch
* send when it has one and falling back to concurrent single sends when it
* does not.
*
* Rejects if ANY message failed to publish, because every caller so far wants
* all-or-nothing: the recovery is to fail the delivery and let redelivery
* republish the whole set, deduped by the per-message `idempotencyKey`. That
* means a partial batch can leave some messages already out — which is
* exactly why the keys are required rather than advisory.
*/
export async function queueMessages(
world: World,
queueName: Parameters<typeof world.queue>[0],
messages: readonly {
message: Parameters<typeof world.queue>[1];
opts?: Parameters<typeof world.queue>[2];
}[]
): Promise<void> {
if (messages.length === 0) return;
const batch = world.queueBatch?.bind(world);
if (!batch) {
await Promise.all(
messages.map((entry) =>
queueMessage(world, queueName, entry.message, entry.opts)
)
);
return;
}
await trace(
'queue.publish',
{
attributes: {
...Attribute.MessagingSystem('vercel-queue'),
...Attribute.MessagingDestinationName(queueName),
...Attribute.MessagingOperationType('publish'),
...Attribute.MessagingBatchMessageCount(messages.length),
...Attribute.PeerService('vercel-queue'),
...Attribute.RpcSystem('vercel-queue'),
...Attribute.RpcService('vqs'),
...Attribute.RpcMethod('publishBatch'),
},
kind: await getSpanKind('PRODUCER'),
},
async () => {
const results = await batch(queueName, messages);
// A World that answers with the wrong number of results has told us
// nothing about the messages it left out. Treated as a failure of the
// whole batch rather than read as success for the entries that ARE
// present: republishing under the same idempotency keys is safe,
// silently never dispatching a step is not (the run makes no progress
// and nothing surfaces an error).
if (results.length !== messages.length) {
throw Object.assign(
new Error(
`Queue batch for ${queueName} returned ${results.length} ` +
`result(s) for ${messages.length} message(s)`
),
{ retryable: true }
);
}
const failures = results.filter((result) => result.error !== undefined);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

AI Review: queueMessages only inspects error, never results.length. A World whose queueBatch returns a short array reports success here: handleSuspension resolves, the delivery is acked, and those steps are never dispatched — the run wedges with no error anywhere in the system.

world-vercel guards this internally (toBatchResult(undefined)) and the SDK length-checks too, so it's unreachable through the world in this PR. But queueBatch is documented in building-a-world.mdx for third-party worlds, which makes it an unenforced contract at exactly the boundary that publishes it — and unlike a rejection, this failure is silent.

I reproduced it with a 64-branch fan-out against a world returning half the results: 32 of 63 steps silently lost, handleSuspension resolved without error. Adding the length check flips it to a clean rejection that converges on redelivery, and all 143 tests in helpers.test.ts + suspension-handler.test.ts still pass:

const results = await batch(queueName, messages);
if (results.length !== messages.length) {
  throw Object.assign(
    new Error(
      `Queue batch for ${queueName} returned ${results.length} result(s) ` +
        `for ${messages.length} message(s)`
    ),
    { retryable: true }
  );
}

if (failures.length === 0) return;
const retryable = failures.some(
(failure) => failure.error !== undefined && failure.retryable
);
const error = new Error(
`Failed to publish ${failures.length} of ${messages.length} queue ` +
`message(s) to ${queueName}: ${failures[0]?.error}`
);
// Carried on the error so a caller CAN tell a transient partial batch
// from a permanent rejection. Nothing reads it yet: today every caller
// rejects the delivery either way, so a permanently rejected entry
// still costs the full redelivery budget. Left in place because the
// information is only available here, and a fast-fail on
// `retryable: false` needs it.
Object.assign(error, { retryable });

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

AI Review: retryable is never read anywhere in packages/core/src/runtime, so the comment above it overstates what happens: queueMessages throws identically for retryable and non-retryable, and the delivery retries either way.

That has a measurable cost. With one permanently-rejected entry in a 64-branch fan-out, the delivery burns all 48 redeliveries and ~2,900 redundant republishes before giving up. The unbatched path does the same, so this isn't a regression — but the batched path is the one that now has the information and discards it, which makes a fast-fail on retryable: false cheap to add later.

throw error;
}
);
}

/**
* Calculates the queue overhead time in milliseconds for a given message.
*/
Expand Down
50 changes: 28 additions & 22 deletions packages/core/src/runtime/suspension-handler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ import {
type LoadedEventLog,
maxEventSlot,
queueMessage,
queueMessages,
slotSnapshotParams,
stepDispatchIdempotencyKey,
} from './helpers.js';
Expand Down Expand Up @@ -1723,29 +1724,34 @@ export async function handleSuspension({
const stepEntries = chunk.filter((entry) => entry.kind === 'step');
if (stepEntries.length === 0) return;
const traceCarrier = await getStepDispatchTraceCarrier();
await Promise.all(
stepEntries.map((entry) =>
queueMessage(
world,
// biome-ignore lint/style/noNonNullAssertion: publishEagerSteps implies presence
stepDispatch!.queueName,
{
runId,
stepId: entry.correlationId,
// One batched publish per chunk instead of one round trip per step.
// The publishes are the fan-out's serialization point: they ride the
// shared 8-connection HTTP/1.1 agent (see `getQueueDispatcher` in
// world-vercel), and the caller awaits all of them before running
// the first inline step body, so N round trips land directly on
// time-to-first-step. `queueMessages` falls back to concurrent
// single sends on a World with no batch support.
await queueMessages(
world,
// biome-ignore lint/style/noNonNullAssertion: publishEagerSteps implies presence
stepDispatch!.queueName,
stepEntries.map((entry) => ({
message: {
runId,
stepId: entry.correlationId,
// biome-ignore lint/style/noNonNullAssertion: set on every 'step' entry at enqueue
stepName: entry.stepName!,
traceCarrier,
requestedAt: new Date(),
},
opts: {
idempotencyKey: stepDispatchIdempotencyKey(
entry.correlationId,
// biome-ignore lint/style/noNonNullAssertion: set on every 'step' entry at enqueue
stepName: entry.stepName!,
traceCarrier,
requestedAt: new Date(),
},
{
idempotencyKey: stepDispatchIdempotencyKey(
entry.correlationId,
// biome-ignore lint/style/noNonNullAssertion: set on every 'step' entry at enqueue
entry.stepName!
),
}
)
)
entry.stepName!
),
},
}))
);
};

Expand Down
9 changes: 9 additions & 0 deletions packages/core/src/telemetry/semantic-conventions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -413,6 +413,15 @@ export const MessagingOperationType = SemanticConvention<
'publish' | 'receive' | 'process'
>('messaging.operation.type');

/**
* Messages carried by one batched publish (standard OTEL:
* messaging.batch.message_count). Set only on the batch send, so a span
* without it is a single-message publish.
*/
export const MessagingBatchMessageCount = SemanticConvention<number>(
'messaging.batch.message_count'
);

/** Time taken to enqueue the message in milliseconds (workflow-specific) */
export const QueueOverheadMs = SemanticConvention<number>(
'workflow.queue.overhead_ms'
Expand Down
Loading
Loading