MCPcopy Create free account
hub / github.com/vercel/vercel / dispatchToConsumer

Method dispatchToConsumer

packages/cli/src/util/dev/queue-broker.ts:400–475  ·  view source on GitHub ↗
(
    message: StoredMessage,
    group: ConsumerGroup
  )

Source from the content-addressed store, hash-verified

398 }
399
400 private async dispatchToConsumer(
401 message: StoredMessage,
402 group: ConsumerGroup
403 ): Promise<void> {
404 const groupDeliveries = this.deliveryState.get(group.id);
405 if (!groupDeliveries) return;
406
407 const state = groupDeliveries.get(message.messageId);
408 if (!state || state.status === 'acked') return;
409
410 if (state.deliveryCount >= group.maxDeliveries) {
411 output.debug(
412 `queues: message ${message.messageId} exceeded maxDeliveries (${group.maxDeliveries}) for group "${group.name}", dropping`
413 );
414 groupDeliveries.delete(message.messageId);
415 this.maybeCleanupMessage(message.messageId);
416 return;
417 }
418
419 const upstream = group.serviceOriginFn();
420 if (!upstream) {
421 // Service not ready yet, retry later
422 state.visibleAt = Date.now() + group.retryAfterMs;
423 return;
424 }
425
426 const receiptHandle = randomBytes(16).toString('hex');
427 state.status = 'in-flight';
428 state.receiptHandle = receiptHandle;
429 state.deliveryCount++;
430 state.leaseExpiresAt = Date.now() + DEFAULT_VISIBILITY_TIMEOUT;
431
432 const now = new Date().toISOString();
433 const expiresAt = new Date(
434 new Date(message.createdAt).getTime() + message.retentionMs
435 ).toISOString();
436
437 output.debug(
438 `queues: dispatching v2beta callback to worker "${group.name}" at ${upstream}`
439 );
440
441 try {
442 const response = await directFetch(`${upstream}/`, {
443 method: 'POST',
444 headers: {
445 'content-type': message.contentType,
446 'ce-type': 'com.vercel.queue.v2beta',
447 'ce-specversion': '1.0',
448 'ce-source': `/topic/${message.queueName}/consumer/${group.name}`,
449 'ce-id': message.messageId,
450 'ce-time': now,
451 'ce-vqsmessageid': message.messageId,
452 'ce-vqsqueuename': message.queueName,
453 'ce-vqsconsumergroup': group.name,
454 'ce-vqsreceipthandle': receiptHandle,
455 'ce-vqsdeliverycount': String(state.deliveryCount),
456 'ce-vqscreatedat': message.createdAt,
457 'ce-vqsexpiresat': expiresAt,

Callers 2

enqueueMethod · 0.95
tickMethod · 0.95

Calls 7

maybeCleanupMessageMethod · 0.95
handleDeliveryFailureMethod · 0.95
directFetchFunction · 0.90
debugMethod · 0.80
getMethod · 0.65
deleteMethod · 0.45
toStringMethod · 0.45

Tested by

no test coverage detected