From 68b61ac38108ac34d612a6f1ce490b5529fd9a69 Mon Sep 17 00:00:00 2001 From: Pranay Prakash Date: Wed, 26 Aug 2026 17:00:11 -0700 Subject: [PATCH 1/3] perf(world-vercel): batch a fan-out's step-execution queue publishes MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A `Promise.all` fan-out dispatched one queue message per branch. Those publishes ride the shared default undici agent (8 connections, HTTP/1.1, `pipelining: 1` — see `getQueueDispatcher`), and `handleSuspension` is awaited in full before the first inline step body runs, so an N-branch fan-out paid ~N/8 serialized round trips straight onto time-to-first-step. The `step_created` writes were already batched and HTTP/2-multiplexed; the publishes were the remaining per-branch round trip. Adds an optional `Queue.queueBatch`, implemented on `@vercel/queue`'s `experimental_sendBatch` (0.5.1), and uses it for the batched fan-out fold's publishes. Each commit chunk now publishes in one request instead of up to 32. `queueBatch` reports per-entry outcomes rather than throwing, because a batch can partially fail. `queueMessages` in core keeps the previous all-or-nothing behavior for this call site: it rejects if any entry failed, so the delivery is redelivered and republishes the set, deduped by the per-step `idempotencyKey` the caller already passed. Worlds without `queueBatch` fall back to concurrent single sends. Co-Authored-By: Claude Opus 5 (1M context) --- .changeset/queue-batch-fanout-dispatch.md | 7 + docs/content/worlds/v5/building-a-world.mdx | 23 ++ packages/core/src/runtime/helpers.test.ts | 90 ++++++++ packages/core/src/runtime/helpers.ts | 63 ++++++ .../core/src/runtime/suspension-handler.ts | 50 +++-- .../src/telemetry/semantic-conventions.ts | 9 + packages/world-vercel/src/queue.test.ts | 146 +++++++++++++ packages/world-vercel/src/queue.ts | 197 +++++++++++++++--- packages/world/src/queue.ts | 48 +++++ pnpm-lock.yaml | 23 +- pnpm-workspace.yaml | 2 +- 11 files changed, 592 insertions(+), 66 deletions(-) create mode 100644 .changeset/queue-batch-fanout-dispatch.md diff --git a/.changeset/queue-batch-fanout-dispatch.md b/.changeset/queue-batch-fanout-dispatch.md new file mode 100644 index 0000000000..b7866b42ea --- /dev/null +++ b/.changeset/queue-batch-fanout-dispatch.md @@ -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`. diff --git a/docs/content/worlds/v5/building-a-world.mdx b/docs/content/worlds/v5/building-a-world.mdx index e71cdb66ee..6e784f2fb7 100644 --- a/docs/content/worlds/v5/building-a-world.mdx +++ b/docs/content/worlds/v5/building-a-world.mdx @@ -229,6 +229,12 @@ 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; + createQueueHandler( queueNamePrefix: QueuePrefix, handler: (message: unknown, meta: { attempt: number; queueName: ValidQueueName; messageId: MessageId }) => Promise @@ -236,6 +242,23 @@ interface Queue { } ``` +### 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. + +{/* @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: diff --git a/packages/core/src/runtime/helpers.test.ts b/packages/core/src/runtime/helpers.test.ts index 21852fac4b..0577880726 100644 --- a/packages/core/src/runtime/helpers.test.ts +++ b/packages/core/src/runtime/helpers.test.ts @@ -22,6 +22,7 @@ import { memoizeEncryptionKey, mergeReportedEvents, preconditionEventDelta, + queueMessages, SLOT_GAP_RECHECK_ATTEMPTS, settleEventSlotGap, slotSnapshotParams, @@ -1068,3 +1069,92 @@ 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) => + ({ + 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('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(); + }); +}); diff --git a/packages/core/src/runtime/helpers.ts b/packages/core/src/runtime/helpers.ts index 8c7a4ce530..2117d61148 100644 --- a/packages/core/src/runtime/helpers.ts +++ b/packages/core/src/runtime/helpers.ts @@ -1200,6 +1200,69 @@ 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[0], + messages: readonly { + message: Parameters[1]; + opts?: Parameters[2]; + }[] +): Promise { + 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); + const failures = results.filter((result) => result.error !== undefined); + 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}` + ); + // Surfaced so a caller (and the delivery-level retry above it) can tell + // a transient partial batch from a permanent rejection. + Object.assign(error, { retryable }); + throw error; + } + ); +} + /** * Calculates the queue overhead time in milliseconds for a given message. */ diff --git a/packages/core/src/runtime/suspension-handler.ts b/packages/core/src/runtime/suspension-handler.ts index f397315902..a0371a707c 100644 --- a/packages/core/src/runtime/suspension-handler.ts +++ b/packages/core/src/runtime/suspension-handler.ts @@ -59,6 +59,7 @@ import { type LoadedEventLog, maxEventSlot, queueMessage, + queueMessages, slotSnapshotParams, stepDispatchIdempotencyKey, } from './helpers.js'; @@ -1602,29 +1603,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! + ), + }, + })) ); }; diff --git a/packages/core/src/telemetry/semantic-conventions.ts b/packages/core/src/telemetry/semantic-conventions.ts index fb6754017f..d418f0fa35 100644 --- a/packages/core/src/telemetry/semantic-conventions.ts +++ b/packages/core/src/telemetry/semantic-conventions.ts @@ -389,6 +389,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( + 'messaging.batch.message_count' +); + /** Time taken to enqueue the message in milliseconds (workflow-specific) */ export const QueueOverheadMs = SemanticConvention( 'workflow.queue.overhead_ms' diff --git a/packages/world-vercel/src/queue.test.ts b/packages/world-vercel/src/queue.test.ts index efa74129c5..9c85a0bf36 100644 --- a/packages/world-vercel/src/queue.test.ts +++ b/packages/world-vercel/src/queue.test.ts @@ -10,6 +10,7 @@ import { const { mockSend, + mockSendBatch, MockConsumerDiscoveryError, MockQueueClient, mockHandleCallback, @@ -22,6 +23,7 @@ const { } const mockSend = vi.fn(); + const mockSendBatch = vi.fn(); const mockHandleCallback = vi.fn(); // Must be a `function` (not an arrow): queue.ts calls `new QueueClient(...)`, // and an arrow function cannot be used as a constructor. @@ -29,12 +31,14 @@ const { const MockQueueClient = vi.fn().mockImplementation(function () { return { send: mockSend, + experimental_sendBatch: mockSendBatch, handleCallback: mockHandleCallback, }; }); return { mockSend, + mockSendBatch, MockConsumerDiscoveryError, MockQueueClient, mockHandleCallback, @@ -1300,3 +1304,145 @@ describe('createQueue', () => { }); }); }); + +describe('queueBatch', () => { + const RUN = 'wrun_01ARZ3NDEKTSV4RRFFQ69G5FAV'; + const sent = (id: string) => ({ status: 'sent' as const, messageId: id }); + + beforeEach(() => { + vi.clearAllMocks(); + process.env.VERCEL_DEPLOYMENT_ID = 'dpl_batch'; + }); + afterEach(() => { + process.env.VERCEL_DEPLOYMENT_ID = undefined; + }); + + const entries = (n: number, runId = RUN) => + Array.from({ length: n }, (_, i) => ({ + message: { runId, stepId: `step-${i}`, stepName: 'myStep' }, + opts: { idempotencyKey: `key-${i}` }, + })); + + it('publishes a whole fan-out in one request and preserves input order', async () => { + mockSendBatch.mockResolvedValueOnce( + Array.from({ length: 5 }, (_, i) => sent(`m${i}`)) + ); + const queue = createQueue(); + assert(queue.queueBatch); + + const results = await queue.queueBatch('__wkf_workflow_test', entries(5)); + + expect(mockSendBatch).toHaveBeenCalledTimes(1); + expect(mockSend).not.toHaveBeenCalled(); + const [topic, messages] = mockSendBatch.mock.calls[0]; + expect(topic).toBe('__wkf_workflow_test'); + expect(messages).toHaveLength(5); + // Each message keeps its own idempotency key: the recovery for a failed + // batch is to republish it, which must not redeliver what already landed. + expect( + messages.map((m: { idempotencyKey?: string }) => m.idempotencyKey) + ).toEqual(['key-0', 'key-1', 'key-2', 'key-3', 'key-4']); + expect(results.map((r) => r.messageId)).toEqual([ + 'm0', + 'm1', + 'm2', + 'm3', + 'm4', + ]); + }); + + it('splits at the 100-message VQS cap', async () => { + mockSendBatch + .mockResolvedValueOnce( + Array.from({ length: 100 }, (_, i) => sent(`a${i}`)) + ) + .mockResolvedValueOnce( + Array.from({ length: 40 }, (_, i) => sent(`b${i}`)) + ); + const queue = createQueue(); + assert(queue.queueBatch); + + const results = await queue.queueBatch('__wkf_workflow_test', entries(140)); + + expect(mockSendBatch).toHaveBeenCalledTimes(2); + expect(mockSendBatch.mock.calls[0][1]).toHaveLength(100); + expect(mockSendBatch.mock.calls[1][1]).toHaveLength(40); + // The split must not be observable in the returned order. + expect(results).toHaveLength(140); + expect(results[0].messageId).toBe('a0'); + expect(results[99].messageId).toBe('a99'); + expect(results[100].messageId).toBe('b0'); + expect(results[139].messageId).toBe('b39'); + }); + + it('reports per-entry failures without rejecting', async () => { + mockSendBatch.mockResolvedValueOnce([ + sent('m0'), + { + status: 'failed', + statusCode: 429, + error: 'rate limited', + retryable: true, + }, + { status: 'deferred', messageId: null }, + ]); + const queue = createQueue(); + assert(queue.queueBatch); + + const results = await queue.queueBatch('__wkf_workflow_test', entries(3)); + + expect(results[0]).toEqual({ messageId: 'm0' }); + expect(results[1]).toEqual({ + messageId: null, + error: 'rate limited', + retryable: true, + }); + // Deferred is an acceptance, not a failure: no `error`, so callers that + // test `error === undefined` treat it as sent. + expect(results[2]).toEqual({ messageId: null }); + }); + + it('flags a short result array as a retryable per-entry failure', async () => { + mockSendBatch.mockResolvedValueOnce([sent('m0')]); + const queue = createQueue(); + assert(queue.queueBatch); + + const results = await queue.queueBatch('__wkf_workflow_test', entries(2)); + + expect(results[0]).toEqual({ messageId: 'm0' }); + expect(results[1]?.error).toMatch(/no result/i); + assert(results[1]?.error !== undefined); + expect(results[1].retryable).toBe(true); + }); + + it('routes messages for different regions through separate requests', async () => { + const { encode } = await import('./run-id/index.js'); + const sfo = `wrun_${encode('01ARZ3NDEKTSV4RRFFQ69G5FAV', 'sfo1')}`; + const fra = `wrun_${encode('01ARZ3NDEKTSV4RRFFQ69G5FAV', 'fra1')}`; + mockSendBatch.mockResolvedValue([sent('x'), sent('y')]); + const queue = createQueue(); + assert(queue.queueBatch); + + const results = await queue.queueBatch('__wkf_workflow_test', [ + ...entries(2, sfo), + ...entries(2, fra), + ]); + + expect(mockSendBatch).toHaveBeenCalledTimes(2); + const regions = ( + MockQueueClient as unknown as { mock: { calls: [{ region?: string }][] } } + ).mock.calls.map((call) => call[0].region); + expect(new Set(regions)).toEqual(new Set(['sfo1', 'fra1'])); + expect(results).toHaveLength(4); + expect(results.every((r) => r.error === undefined)).toBe(true); + }); + + it('returns an empty result set without touching the transport', async () => { + const queue = createQueue(); + assert(queue.queueBatch); + await expect(queue.queueBatch('__wkf_workflow_test', [])).resolves.toEqual( + [] + ); + expect(mockSendBatch).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/world-vercel/src/queue.ts b/packages/world-vercel/src/queue.ts index 0f1eec291e..b700761e7b 100644 --- a/packages/world-vercel/src/queue.ts +++ b/packages/world-vercel/src/queue.ts @@ -5,6 +5,7 @@ import { globalSingleton } from '@workflow/utils'; import { MessageId, type Queue, + type QueueBatchResult, type QueueOptions, type QueuePayload, QueuePayloadSchema, @@ -21,6 +22,46 @@ import { isKnownRegionCode, REGION_IDS } from './run-id/regions.js'; import { type APIConfig, getHeaders, getHttpUrl } from './utils.js'; import { isWsEventsTransportEnabled } from './ws-transport-enabled.js'; +/** + * Messages per `experimental_sendBatch` request. VQS caps a batch at 100 and + * rejects the whole request above it, so this is the API's ceiling rather + * than a tuning knob; `queueBatch` splits anything larger. + */ +const MAX_QUEUE_SEND_BATCH = 100; + +/** + * Maps one `experimental_sendBatch` outcome onto the World's + * {@link QueueBatchResult}. + * + * `undefined` means the server returned fewer results than the batch carried. + * That is reported as a retryable failure rather than left as a hole the + * caller would read as success: republishing under the same idempotency keys + * is safe, silently dropping a step's message is not. + */ +function toBatchResult( + outcome: + | Awaited>[number] + | undefined +): QueueBatchResult { + if (outcome === undefined) { + return { + messageId: null, + error: 'Queue batch returned no result for this message', + retryable: true, + }; + } + if (outcome.status === 'failed') { + return { + messageId: null, + error: outcome.error, + retryable: outcome.retryable, + }; + } + return { + messageId: outcome.messageId ? MessageId.parse(outcome.messageId) : null, + }; +} + /** * CBOR-based queue transport. Encodes values with cbor-x on send and * decodes on receive, preserving Uint8Array values natively (workflow @@ -429,9 +470,16 @@ export function createQueue(config?: APIConfig): Queue { headers: Object.fromEntries(headers.entries()), }; - const queue: QueueFunction = async ( - queueName, - payload, + /** + * Resolves everything a send needs from one (payload, opts) pair: the + * routing dimensions that decide WHICH client the message goes through + * (region / deploymentId / transport / physical topic) and the per-message + * arguments. Shared by `queue` and `queueBatch` so a batched send routes + * byte-for-byte the same way the single send would have. + */ + const prepareSend = ( + queueName: ValidQueueName, + payload: QueuePayload, opts?: QueueOptions ) => { // Check if we have a deployment ID either from options or environment @@ -448,7 +496,6 @@ export function createQueue(config?: APIConfig): Queue { const useCbor = (opts?.specVersion ?? SPEC_VERSION_CURRENT) >= SPEC_VERSION_SUPPORTS_CBOR_QUEUE_TRANSPORT; - const transport = useCbor ? cborTransport : jsonTransport; // Resolve the destination region. Explicit `opts.region` wins, otherwise // we decode it from the payload's tagged run ID so messages produced by @@ -457,7 +504,43 @@ export function createQueue(config?: APIConfig): Queue { // behavior for legacy / untagged run IDs. const region = resolveTargetRegion(payload, opts); - const client = new QueueClient({ + const topic = getPhysicalQueueName(queueName, payload).replace( + /[^A-Za-z0-9-_]/g, + '-' + ); + + return { + deploymentId, + useCbor, + region, + topic, + // The CborTransport handles CBOR encoding inside serialize(), + // preserving Uint8Array values (workflow input in specVersion >= 2). + 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, + }, + sendOptions: { + idempotencyKey: opts?.idempotencyKey, + delaySeconds: opts?.delaySeconds, + headers: { + ...getHeadersFromPayload(payload), + ...opts?.headers, + }, + }, + }; + }; + + const clientFor = (route: { + region: string; + deploymentId: string; + useCbor: boolean; + }) => + new QueueClient({ ...clientOptions, // When sending through the api.vercel.com proxy, the fixed // `resolveBaseUrl` above replaces the queue SDK's own @@ -468,39 +551,29 @@ export function createQueue(config?: APIConfig): Queue { ...(usingProxy && { headers: { ...clientOptions.headers, - 'x-vercel-queue-region': region, + 'x-vercel-queue-region': route.region, }, }), - region, - deploymentId, - transport, + region: route.region, + deploymentId: route.deploymentId, + transport: route.useCbor ? cborTransport : jsonTransport, }); - // The CborTransport handles CBOR encoding inside serialize(), - // 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 = getPhysicalQueueName(queueName, payload).replace( - /[^A-Za-z0-9-_]/g, - '-' - ); + const queue: QueueFunction = async ( + queueName, + payload, + opts?: QueueOptions + ) => { + const prepared = prepareSend(queueName, payload, opts); + const client = clientFor(prepared); // A repeated `idempotencyKey` is accepted rather than rejected: the send // returns a fresh message ID and only one of the messages is delivered, so // there is no conflict for the caller to handle here. - const { messageId } = await client.send(sanitizedQueueName, wrapper, { - idempotencyKey: opts?.idempotencyKey, - delaySeconds: opts?.delaySeconds, - headers: { - ...getHeadersFromPayload(payload), - ...opts?.headers, - }, - }); + const { messageId } = await client.send( + prepared.topic, + prepared.wrapper, + prepared.sendOptions + ); return { // messageId may be null when the queue fails over to a different region: // the event is ingested but the responding region cannot return an ID. @@ -508,6 +581,67 @@ export function createQueue(config?: APIConfig): Queue { }; }; + const queueBatch: NonNullable = async ( + queueName, + messages + ) => { + const results = new Array(messages.length); + if (messages.length === 0) return results; + + // Group by the routing dimensions a single VQS request cannot span. In + // the case this exists for — one run's fan-out to one logical queue — + // every message lands in one group, so this is one request per + // MAX_QUEUE_SEND_BATCH messages. Mixed input still works, it just costs + // one request per distinct route. + const groups = new Map< + string, + { + route: { region: string; deploymentId: string; useCbor: boolean }; + entries: { + index: number; + topic: string; + message: Parameters[1][number]; + }[]; + } + >(); + for (const [index, entry] of messages.entries()) { + const prepared = prepareSend(queueName, entry.message, entry.opts); + const key = `${prepared.region}${prepared.deploymentId}${prepared.useCbor}${prepared.topic}`; + const group = groups.get(key) ?? { + route: prepared, + entries: [], + }; + group.entries.push({ + index, + topic: prepared.topic, + message: { payload: prepared.wrapper, ...prepared.sendOptions }, + }); + groups.set(key, group); + } + + await Promise.all( + [...groups.values()].map(async ({ route, entries }) => { + const client = clientFor(route); + for ( + let offset = 0; + offset < entries.length; + offset += MAX_QUEUE_SEND_BATCH + ) { + const chunk = entries.slice(offset, offset + MAX_QUEUE_SEND_BATCH); + const sent = await client.experimental_sendBatch( + // biome-ignore lint/style/noNonNullAssertion: chunks are non-empty + chunk[0]!.topic, + chunk.map((entry) => entry.message) + ); + for (const [position, entry] of chunk.entries()) { + results[entry.index] = toBatchResult(sent[position]); + } + } + }) + ); + return results; + }; + const createQueueHandler: Queue['createQueueHandler'] = ( _prefix, handler @@ -606,6 +740,7 @@ export function createQueue(config?: APIConfig): Queue { return { queue, + queueBatch, createQueueHandler, getDeploymentId, isDeploymentUnavailableError, diff --git a/packages/world/src/queue.ts b/packages/world/src/queue.ts index 912c19f065..f5f54b5390 100644 --- a/packages/world/src/queue.ts +++ b/packages/world/src/queue.ts @@ -427,6 +427,23 @@ export interface QueueOptions { region?: string; } +/** + * Outcome of one message in a {@link Queue.queueBatch} call, in input order. + * + * `messageId: null` with no `error` means accepted for deferred processing + * (the same "accepted, no ID yet" case {@link Queue.queue} reports), so + * `error === undefined` is the success test, not a non-null `messageId`. + */ +export type QueueBatchResult = + | { messageId: MessageId | null; error?: undefined } + | { + messageId: null; + /** Human-readable failure description for this message. */ + error: string; + /** Republishing the batch may succeed. */ + retryable: boolean; + }; + export interface Queue { getDeploymentId(): Promise; @@ -450,6 +467,37 @@ export interface Queue { opts?: QueueOptions ): Promise<{ messageId: MessageId | null }>; + /** + * Enqueues several messages to the SAME logical queue in as few round trips + * as the backing transport allows. Optional: callers MUST fall back to + * per-message {@link Queue.queue} when a World does not implement it. + * + * Exists for wide fan-outs. Publishing an N-branch `Promise.all` one message + * at a time costs N round trips through a bounded connection pool, and that + * cost is paid before the fan-out's first step body runs, so it lands + * directly on time-to-first-step. + * + * Contract: + * + * - Results are returned in input order, one per input message. + * - Partial failure is normal. A rejected entry reports `error`; a + * `retryable` entry may succeed if the whole batch is published again. + * Implementations MUST NOT throw for a per-entry failure — reserve + * rejection for request-level failures where no entry outcome is known. + * - Callers are expected to pass `opts.idempotencyKey` per message, because + * the recovery for both a request-level failure and a retryable entry is + * to republish the batch: without keys that redelivers the entries that + * already succeeded. + * - Every message must target one logical `queueName`. Implementations are + * free to split the batch (by transport cap, or by any per-message routing + * dimension they derive from the payload, such as region or physical + * topic); the split must not be observable in the returned order. + */ + queueBatch?( + queueName: ValidQueueName, + messages: readonly { message: QueuePayload; opts?: QueueOptions }[] + ): Promise; + /** * Creates an HTTP queue handler for processing messages from a specific queue. * A rejected handler must retry the same message with an incremented attempt. diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index b98f4765a2..f7629b19f0 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -22,8 +22,8 @@ catalogs: specifier: 3.2.0 version: 3.2.0 '@vercel/queue': - specifier: 0.5.0 - version: 0.5.0 + specifier: 0.5.1 + version: 0.5.1 '@vitest/coverage-v8': specifier: ^4.1.10 version: 4.1.10 @@ -1389,7 +1389,7 @@ importers: dependencies: '@vercel/queue': specifier: 'catalog:' - version: 0.5.0(@opentelemetry/api@1.9.1) + version: 0.5.1(@opentelemetry/api@1.9.1) '@workflow/errors': specifier: workspace:* version: link:../errors @@ -1444,7 +1444,7 @@ importers: dependencies: '@vercel/queue': specifier: 'catalog:' - version: 0.5.0(@opentelemetry/api@1.9.1) + version: 0.5.1(@opentelemetry/api@1.9.1) '@workflow/errors': specifier: workspace:* version: link:../errors @@ -1591,7 +1591,7 @@ importers: version: 3.2.0 '@vercel/queue': specifier: 'catalog:' - version: 0.5.0(@opentelemetry/api@1.9.1) + version: 0.5.1(@opentelemetry/api@1.9.1) '@workflow/errors': specifier: workspace:* version: link:../errors @@ -10078,8 +10078,8 @@ packages: resolution: {integrity: sha512-6pjdXyNfdCQnj1nyeB1rPtll/XUhmciyeZJD0rpIUUwmcIfR+utDl6+iFvlHvsBdqAVT8UC6ydObT0v+xldfUQ==} engines: {node: '>=20.0.0'} - '@vercel/queue@0.5.0': - resolution: {integrity: sha512-TqvoMuhZK9Hw2yal6ee6EZ+wQw55+qPrpev6O6jJrpJMC1myrCrt7bLZvb9GQDIR35h40EiETVaMZT4eFWBVpA==} + '@vercel/queue@0.5.1': + resolution: {integrity: sha512-7gO8hOtIh15UDzbjCg7oJDOwY1Hs3XC4xSK6HjMVo93Ge3oEhPZJ18ogamwZTXwcFRYwvgN1vp2w3+zgqfIpoA==} engines: {node: '>=20.0.0'} peerDependencies: '@opentelemetry/api': 1.9.1 @@ -20839,7 +20839,7 @@ snapshots: '@types/shimmer': 1.2.0 import-in-the-middle: 1.15.0 require-in-the-middle: 7.5.2 - semver: 7.8.2 + semver: 7.8.5 shimmer: 1.2.1 transitivePeerDependencies: - supports-color @@ -26707,7 +26707,7 @@ snapshots: picocolors: 1.1.1 optional: true - '@vercel/queue@0.5.0(@opentelemetry/api@1.9.1)': + '@vercel/queue@0.5.1(@opentelemetry/api@1.9.1)': dependencies: '@vercel/oidc': 3.2.0 minimatch: 10.2.5 @@ -32505,7 +32505,7 @@ snapshots: node-abi@3.89.0: dependencies: - semver: 7.8.2 + semver: 7.8.5 optional: true node-abort-controller@3.1.1: {} @@ -34778,8 +34778,7 @@ snapshots: semver@7.8.2: {} - semver@7.8.5: - optional: true + semver@7.8.5: {} send@1.2.0: dependencies: diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index c708418589..507f87d47b 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -12,7 +12,7 @@ catalog: "@types/node": 22.19.0 "@vercel/functions": ^3.8.0 "@vercel/oidc": 3.2.0 - "@vercel/queue": 0.5.0 + "@vercel/queue": 0.5.1 "@vitest/coverage-v8": ^4.1.10 "@vitest/runner": ^4.1.10 ai: 6.0.116 From 480b34f5652e76b2896700c2e97d13c98c177191 Mon Sep 17 00:00:00 2001 From: Karthik Kalyanaraman Date: Fri, 11 Sep 2026 09:50:14 -0700 Subject: [PATCH 2/3] fix(core): reject a short queueBatch result set instead of reading it as success `queueMessages` only inspected `error`, so a World whose `queueBatch` returned fewer results than it was given messages reported success for the whole batch. The omitted entries were never published and nothing raised: `handleSuspension` resolved, the delivery was acked, and those steps were never dispatched, so the run stalls with no error recorded anywhere. Reproduced at 64 branches against a World returning half its results: 32 of 63 steps silently lost. world-vercel guards this internally and `@vercel/queue` length-checks its own response, so it was not reachable through the world added here. It is reachable through the interface `building-a-world` opens to third-party worlds, which is where the check belongs. Documented on the interface and in the guide alongside it. Also notes that the batch grouping degenerates to one request per message under WORKFLOW_SEQUENTIAL_REPLAYS=1 (per-step physical topics are one of the routing dimensions groups split on), and corrects the comment claiming the error's `retryable` flag is consumed downstream: nothing reads it yet. Co-Authored-By: Claude Opus 5 (1M context) --- docs/content/worlds/v5/building-a-world.mdx | 2 ++ packages/core/src/runtime/helpers.test.ts | 23 +++++++++++++++++++++ packages/core/src/runtime/helpers.ts | 23 +++++++++++++++++++-- packages/world-vercel/src/queue.test.ts | 2 +- packages/world-vercel/src/queue.ts | 10 +++++++++ packages/world/src/queue.ts | 6 +++++- 6 files changed, 62 insertions(+), 4 deletions(-) diff --git a/docs/content/worlds/v5/building-a-world.mdx b/docs/content/worlds/v5/building-a-world.mdx index 6e784f2fb7..c40724cb04 100644 --- a/docs/content/worlds/v5/building-a-world.mdx +++ b/docs/content/worlds/v5/building-a-world.mdx @@ -248,6 +248,8 @@ interface Queue { 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 = diff --git a/packages/core/src/runtime/helpers.test.ts b/packages/core/src/runtime/helpers.test.ts index 8a0595496a..0de1f58279 100644 --- a/packages/core/src/runtime/helpers.test.ts +++ b/packages/core/src/runtime/helpers.test.ts @@ -1148,6 +1148,29 @@ describe('queueMessages', () => { ).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 }); diff --git a/packages/core/src/runtime/helpers.ts b/packages/core/src/runtime/helpers.ts index 785f1dcc75..a664032535 100644 --- a/packages/core/src/runtime/helpers.ts +++ b/packages/core/src/runtime/helpers.ts @@ -1262,6 +1262,21 @@ export async function queueMessages( }, 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); if (failures.length === 0) return; const retryable = failures.some( @@ -1271,8 +1286,12 @@ export async function queueMessages( `Failed to publish ${failures.length} of ${messages.length} queue ` + `message(s) to ${queueName}: ${failures[0]?.error}` ); - // Surfaced so a caller (and the delivery-level retry above it) can tell - // a transient partial batch from a permanent rejection. + // 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 }); throw error; } diff --git a/packages/world-vercel/src/queue.test.ts b/packages/world-vercel/src/queue.test.ts index 9c85a0bf36..dfb1e403f0 100644 --- a/packages/world-vercel/src/queue.test.ts +++ b/packages/world-vercel/src/queue.test.ts @@ -1314,7 +1314,7 @@ describe('queueBatch', () => { process.env.VERCEL_DEPLOYMENT_ID = 'dpl_batch'; }); afterEach(() => { - process.env.VERCEL_DEPLOYMENT_ID = undefined; + delete process.env.VERCEL_DEPLOYMENT_ID; }); const entries = (n: number, runId = RUN) => diff --git a/packages/world-vercel/src/queue.ts b/packages/world-vercel/src/queue.ts index d96caf8980..d8c999c8e7 100644 --- a/packages/world-vercel/src/queue.ts +++ b/packages/world-vercel/src/queue.ts @@ -595,6 +595,16 @@ export function createQueue(config?: APIConfig): Queue { // every message lands in one group, so this is one request per // MAX_QUEUE_SEND_BATCH messages. Mixed input still works, it just costs // one request per distinct route. + // + // `topic` is one of those dimensions, which makes this a no-op under + // WORKFLOW_SEQUENTIAL_REPLAYS=1: step dispatches ride the flow topic, + // `getPhysicalQueueName` gives each one a per-step physical topic, and + // every message therefore lands in a group of its own. The fan-out still + // publishes correctly, just at one request per step as before, and each + // goes through the batch endpoint rather than `send()` — which does not + // map 502 `consumer_discovery_failed` to ConsumerDiscoveryError (only + // 503). No caller on this path classifies that error today, so nothing + // changes behaviorally; worth knowing before one starts. const groups = new Map< string, { diff --git a/packages/world/src/queue.ts b/packages/world/src/queue.ts index 0424968ce0..01699b000b 100644 --- a/packages/world/src/queue.ts +++ b/packages/world/src/queue.ts @@ -507,7 +507,11 @@ export interface Queue { * * Contract: * - * - Results are returned in input order, one per input message. + * - Results are returned in input order, one per input message. Returning + * a different number of results than there were messages is a contract + * violation the runtime rejects the whole batch on: an omitted result is + * indistinguishable from a message that was never published, and reading + * it as success would strand that step with no error anywhere. * - Partial failure is normal. A rejected entry reports `error`; a * `retryable` entry may succeed if the whole batch is published again. * Implementations MUST NOT throw for a per-entry failure — reserve From 367e5e096d243a041928e89650804757ec61313f Mon Sep 17 00:00:00 2001 From: Karthik Kalyanaraman Date: Fri, 11 Sep 2026 10:44:12 -0700 Subject: [PATCH 3/3] fix(world-vercel): carry trace context on each batched queue message MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `experimental_sendBatch` injects the active trace context into the multipart REQUEST headers, and the per-part headers it builds never see it. VQS stores headers per message and re-emits a stored `traceparent` at delivery as `x-vercel-queue-traceparent`, which is what lets a consumer attach a span link back to its producer, so a batched message arrived with no producer context and its `vqs.process` span got no link. `send()` is unaffected: for a single message the request headers ARE that message's headers. At 64 branches that was 63 of 64 step dispatches losing the transport-level producer link. The run's own step tracing was never affected: that carrier travels in the message payload (`WorkflowInvokePayload.traceCarrier`), which is what the consumer builds its trace context from, not a header. Injects the active context into each entry's headers in `queueBatch` — last, so it wins over caller-supplied `opts.headers` exactly as the SDK's own injection does — and honors VERCEL_QUEUE_TRACE_PROPAGATION so that kill switch still covers both paths. `getTraceContextHeaders()` is factored out of `injectTraceContextIntoHeaders` so the two share one source. Verified on the wire against a stub VQS speaking the real batch endpoint: `traceparent` carrying the producer's traceId/spanId lands on all 64 multipart parts through the real SDK, with the per-message idempotency keys still alongside it. Co-Authored-By: Claude Opus 5 (1M context) --- packages/world-vercel/src/queue.test.ts | 106 ++++++++++++++++++++++++ packages/world-vercel/src/queue.ts | 37 ++++++++- packages/world-vercel/src/telemetry.ts | 24 +++++- 3 files changed, 162 insertions(+), 5 deletions(-) diff --git a/packages/world-vercel/src/queue.test.ts b/packages/world-vercel/src/queue.test.ts index dfb1e403f0..d9ded8226c 100644 --- a/packages/world-vercel/src/queue.test.ts +++ b/packages/world-vercel/src/queue.test.ts @@ -1,6 +1,16 @@ +import { context, trace as otelTrace, propagation } from '@opentelemetry/api'; +import { AsyncLocalStorageContextManager } from '@opentelemetry/context-async-hooks'; +import { W3CTraceContextPropagator } from '@opentelemetry/core'; import { + BasicTracerProvider, + InMemorySpanExporter, + SimpleSpanProcessor, +} from '@opentelemetry/sdk-trace-base'; +import { + afterAll, afterEach, assert, + beforeAll, beforeEach, describe, expect, @@ -1446,3 +1456,99 @@ describe('queueBatch', () => { expect(mockSendBatch).not.toHaveBeenCalled(); }); }); + +/** + * A batched message carries its producer context on its OWN headers. The SDK + * injects into the multipart request's headers, which VQS does not store per + * message, so world-vercel injects per entry; without it a consumer's + * `vqs.process` span has no link back to the producer (vqs-server re-emits a + * stored `traceparent` as `x-vercel-queue-traceparent` at delivery). + */ +describe('queueBatch trace propagation', () => { + const exporter = new InMemorySpanExporter(); + const provider = new BasicTracerProvider(); + const contextManager = new AsyncLocalStorageContextManager(); + + beforeAll(() => { + provider.addSpanProcessor(new SimpleSpanProcessor(exporter)); + contextManager.enable(); + context.setGlobalContextManager(contextManager); + propagation.setGlobalPropagator(new W3CTraceContextPropagator()); + otelTrace.setGlobalTracerProvider(provider); + }); + + afterAll(async () => { + await provider.shutdown(); + context.disable(); + propagation.disable(); + otelTrace.disable(); + }); + + beforeEach(() => { + vi.clearAllMocks(); + process.env.VERCEL_DEPLOYMENT_ID = 'dpl_trace'; + }); + + afterEach(() => { + delete process.env.VERCEL_DEPLOYMENT_ID; + delete process.env.VERCEL_QUEUE_TRACE_PROPAGATION; + }); + + const entries = (n: number) => + Array.from({ length: n }, (_, i) => ({ + message: { runId: 'wrun_trace', stepId: `step-${i}` }, + opts: { idempotencyKey: `key-${i}` }, + })); + + /** Publishes inside an active span and returns the sent message headers. */ + async function publishInSpan( + count: number + ): Promise< + { headers: Record | undefined; spanId: string }[] + > { + mockSendBatch.mockResolvedValueOnce( + Array.from({ length: count }, (_, i) => ({ + status: 'sent' as const, + messageId: `m${i}`, + })) + ); + const queue = createQueue(); + assert(queue.queueBatch); + const span = provider.getTracer('test').startSpan('publish'); + const spanId = span.spanContext().spanId; + await context.with(otelTrace.setSpan(context.active(), span), async () => { + await queue.queueBatch?.('__wkf_workflow_test', entries(count)); + }); + span.end(); + const sent = mockSendBatch.mock.calls[0]?.[1] as + | { headers?: Record }[] + | undefined; + return (sent ?? []).map((m) => ({ headers: m.headers, spanId })); + } + + it('puts the producer traceparent on EVERY message in the batch', async () => { + const sent = await publishInSpan(64); + + expect(sent).toHaveLength(64); + for (const { headers, spanId } of sent) { + // Same span on every entry: one publish, one producer context. + expect(headers?.traceparent).toMatch(/^00-[0-9a-f]{32}-[0-9a-f]{16}-/); + expect(headers?.traceparent).toContain(spanId); + } + // The payload-derived headers the single send also carries survive it. + expect(sent[0].headers?.['x-vercel-workflow-run-id']).toBe('wrun_trace'); + expect(sent[0].headers?.['x-vercel-workflow-step-id']).toBe('step-0'); + }); + + it('honors VERCEL_QUEUE_TRACE_PROPAGATION=off, like the SDK does', async () => { + process.env.VERCEL_QUEUE_TRACE_PROPAGATION = 'off'; + const sent = await publishInSpan(2); + + expect(sent).toHaveLength(2); + for (const { headers } of sent) { + expect(headers?.traceparent).toBeUndefined(); + // The kill switch is trace-only; message routing headers stay. + expect(headers?.['x-vercel-workflow-run-id']).toBe('wrun_trace'); + } + }); +}); diff --git a/packages/world-vercel/src/queue.ts b/packages/world-vercel/src/queue.ts index d8c999c8e7..0156ab58a6 100644 --- a/packages/world-vercel/src/queue.ts +++ b/packages/world-vercel/src/queue.ts @@ -19,6 +19,7 @@ import { missingDeploymentIdMessage } from './deployment-id.js'; import { getQueueDispatcher } from './http-client.js'; import { decode as decodeTaggedRunId } from './run-id/index.js'; import { isKnownRegionCode, REGION_IDS } from './run-id/regions.js'; +import { getTraceContextHeaders } from './telemetry.js'; import { type APIConfig, getHeaders, getHttpUrl } from './utils.js'; import { isWsEventsTransportEnabled } from './ws-transport-enabled.js'; @@ -29,6 +30,16 @@ import { isWsEventsTransportEnabled } from './ws-transport-enabled.js'; */ const MAX_QUEUE_SEND_BATCH = 100; +/** + * Mirrors `@vercel/queue`'s own kill switch. `queueBatch` injects trace + * context itself (see below), so without this check `off` would still + * disable it on the single send and not on the batched one. + */ +function isQueueTracePropagationDisabled(): boolean { + const value = process.env.VERCEL_QUEUE_TRACE_PROPAGATION?.toLowerCase(); + return value === 'off' || value === '0' || value === 'false'; +} + /** * Maps one `experimental_sendBatch` outcome onto the World's * {@link QueueBatchResult}. @@ -590,6 +601,24 @@ export function createQueue(config?: APIConfig): Queue { const results = new Array(messages.length); if (messages.length === 0) return results; + // Trace context has to be attached per MESSAGE here, not per request. + // + // `send()` gets this for free: the SDK injects into the headers it is + // about to send, and for a single message those headers ARE the message's + // headers, so VQS stores the `traceparent` and re-emits it at delivery as + // `x-vercel-queue-traceparent` — which is what lets a consumer attach a + // span link back to this producer. `experimental_sendBatch` injects into + // the multipart REQUEST headers instead, and the per-part headers it + // builds never see it, so a batched message would arrive with no producer + // context and the consumer's `vqs.process` span would have no link. + // + // One injection for the whole call: every message in a batch is published + // under the same active span, which is exactly what the single send would + // have recorded on each of them. + const traceHeaders = isQueueTracePropagationDisabled() + ? {} + : await getTraceContextHeaders(); + // Group by the routing dimensions a single VQS request cannot span. In // the case this exists for — one run's fan-out to one logical queue — // every message lands in one group, so this is one request per @@ -626,7 +655,13 @@ export function createQueue(config?: APIConfig): Queue { group.entries.push({ index, topic: prepared.topic, - message: { payload: prepared.wrapper, ...prepared.sendOptions }, + message: { + payload: prepared.wrapper, + ...prepared.sendOptions, + // Trace headers last, matching the single send: the SDK injects + // after it has applied the caller's `opts.headers`. + headers: { ...prepared.sendOptions.headers, ...traceHeaders }, + }, }); groups.set(key, group); } diff --git a/packages/world-vercel/src/telemetry.ts b/packages/world-vercel/src/telemetry.ts index 815d95a5a9..68ae75beba 100644 --- a/packages/world-vercel/src/telemetry.ts +++ b/packages/world-vercel/src/telemetry.ts @@ -186,13 +186,29 @@ export async function getSpanKind( export async function injectTraceContextIntoHeaders( headers: Headers ): Promise { + for (const [key, value] of Object.entries(await getTraceContextHeaders())) { + headers.set(key, value); + } +} + +/** + * The active W3C trace context as a plain header record. + * + * The same source as {@link injectTraceContextIntoHeaders}, shaped for APIs + * that take a header map rather than a `Headers` — a batched queue send + * carries its context on each message's own headers, not on the request's. + * + * Empty when `@opentelemetry/api` is unavailable or no propagator is + * registered, so callers can spread it unconditionally. + */ +export async function getTraceContextHeaders(): Promise< + Record +> { const otel = await getOtelApi(); - if (!otel) return; + if (!otel) return {}; const carrier: Record = {}; otel.propagation.inject(otel.context.active(), carrier); - for (const [key, value] of Object.entries(carrier)) { - headers.set(key, value); - } + return carrier; } // Semantic conventions for World/Storage tracing