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..c40724cb04 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,25 @@ 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. + +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: diff --git a/packages/core/src/runtime/helpers.test.ts b/packages/core/src/runtime/helpers.test.ts index c0b8077fb6..0de1f58279 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,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) => + ({ + 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(); + }); +}); diff --git a/packages/core/src/runtime/helpers.ts b/packages/core/src/runtime/helpers.ts index 82d9e44016..a664032535 100644 --- a/packages/core/src/runtime/helpers.ts +++ b/packages/core/src/runtime/helpers.ts @@ -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[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); + // 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( + (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 }); + 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 559fdcbf0a..36da2650a0 100644 --- a/packages/core/src/runtime/suspension-handler.ts +++ b/packages/core/src/runtime/suspension-handler.ts @@ -60,6 +60,7 @@ import { type LoadedEventLog, maxEventSlot, queueMessage, + queueMessages, slotSnapshotParams, stepDispatchIdempotencyKey, } from './helpers.js'; @@ -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! + ), + }, + })) ); }; diff --git a/packages/core/src/telemetry/semantic-conventions.ts b/packages/core/src/telemetry/semantic-conventions.ts index dbd0f7b5ae..a4aaaaaa82 100644 --- a/packages/core/src/telemetry/semantic-conventions.ts +++ b/packages/core/src/telemetry/semantic-conventions.ts @@ -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( + '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..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, @@ -10,6 +20,7 @@ import { const { mockSend, + mockSendBatch, MockConsumerDiscoveryError, MockQueueClient, mockHandleCallback, @@ -22,6 +33,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 +41,14 @@ const { const MockQueueClient = vi.fn().mockImplementation(function () { return { send: mockSend, + experimental_sendBatch: mockSendBatch, handleCallback: mockHandleCallback, }; }); return { mockSend, + mockSendBatch, MockConsumerDiscoveryError, MockQueueClient, mockHandleCallback, @@ -1300,3 +1314,241 @@ 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(() => { + delete process.env.VERCEL_DEPLOYMENT_ID; + }); + + 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(); + }); +}); + +/** + * 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 0541c16de7..0156ab58a6 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, @@ -18,9 +19,60 @@ 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'; +/** + * 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; + +/** + * 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}. + * + * `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 @@ -431,9 +483,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 @@ -450,7 +509,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 @@ -459,7 +517,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 @@ -470,39 +564,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. @@ -510,6 +594,101 @@ export function createQueue(config?: APIConfig): Queue { }; }; + const queueBatch: NonNullable = async ( + queueName, + messages + ) => { + 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 + // 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, + { + 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, + // 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); + } + + 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 @@ -608,6 +787,7 @@ export function createQueue(config?: APIConfig): Queue { return { queue, + queueBatch, createQueueHandler, getDeploymentId, isDeploymentUnavailableError, 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 diff --git a/packages/world/src/queue.ts b/packages/world/src/queue.ts index af1261b2ec..01699b000b 100644 --- a/packages/world/src/queue.ts +++ b/packages/world/src/queue.ts @@ -455,6 +455,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; @@ -478,6 +495,41 @@ 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. 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 + * 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 6fe4e3b9ca..af1688daa2 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 @@ -1404,7 +1404,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 @@ -1459,7 +1459,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 @@ -1606,7 +1606,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 @@ -10274,8 +10274,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 @@ -27053,7 +27053,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 diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index 66e0613d2f..57b090d823 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