diff --git a/apps/webapp/app/services/events/deliverEvent.server.ts b/apps/webapp/app/services/events/deliverEvent.server.ts index e6cbf525c..12b20f80e 100644 --- a/apps/webapp/app/services/events/deliverEvent.server.ts +++ b/apps/webapp/app/services/events/deliverEvent.server.ts @@ -74,9 +74,19 @@ export class DeliverEventService { if (eventDispatcher.batcher) { const { maxPayloads, runAt } = this.#getBatchEnqueueOptions(eventDispatcher.batcher); + const jobKeyParts = [eventDispatcher.id]; + + if (eventRecord.isTest) { + jobKeyParts.push(String(eventRecord.isTest)); + } + + if (eventRecord.externalAccountId) { + jobKeyParts.push(eventRecord.externalAccountId); + } + return workerQueue.batchEnqueue("events.invokeDispatchBatcher", [eventRecord.id], { tx, - jobKey: eventDispatcher.id, + jobKey: jobKeyParts.join(":"), maxPayloads, runAt, }); diff --git a/apps/webapp/app/services/worker.server.ts b/apps/webapp/app/services/worker.server.ts index fd031f810..1a1b7f455 100644 --- a/apps/webapp/app/services/worker.server.ts +++ b/apps/webapp/app/services/worker.server.ts @@ -255,9 +255,11 @@ function getWorkerQueue() { throw new Error("Job key is required for batch jobs."); } + const batcherId = job.key.split(":")[0] + const service = new DispatchBatcherService(); - await service.call(job.key, payload); + await service.call(batcherId, payload); }, }, "events.invokeBatchDispatcher": { @@ -351,9 +353,11 @@ function getWorkerQueue() { throw new Error("Job key is required for batch jobs."); } + const batcherId = job.key.split(":")[0] + const service = new WebhookDeliveryBatcherService(); - await service.call(job.key, payload); + await service.call(batcherId, payload); }, }, deliverMultipleWebhookRequests: {