From 4d09c400a45241b58d9349d8cf7db7d80b739bae Mon Sep 17 00:00:00 2001 From: martmull Date: Thu, 23 Jul 2026 17:50:43 +0200 Subject: [PATCH] Remove DATABASE_EVENT_JOBS_CHUNK_SIZE and Promise.all from logic function trigger jobs (#23205) ## What - `LogicFunctionTriggerJob` now processes a single `LogicFunctionTriggerJobData` payload instead of an array processed with `Promise.all`. - Removed `DATABASE_EVENT_JOBS_CHUNK_SIZE` and the `lodash.chunk` usage in `CallDatabaseEventTriggerJobsJob`. - Added `bulkAdd` to `MessageQueueService` and both drivers (BullMQ driver uses native `queue.addBulk`, sync driver processes payloads sequentially). `CallDatabaseEventTriggerJobsJob` uses it to enqueue all payloads in one call. - Updated the other producers (`ServerRouteTriggerService`, `ApplicationInstallService`, `ConnectionProviderOauthFlowService`, `CronTriggerCronJob`) to enqueue a single payload instead of a one-element array, and updated the corresponding specs. ## Why Each logic function execution now gets its own queue job, so a failing execution only retries itself instead of re-running the whole chunk, and job-level retry/metrics apply per execution. --- _Generated by [Claude Code](https://claude.ai/code/session_018VCs2kopnDiZCL41eToxQF)_ Review in cubic --- .../application-install.service.ts | 14 ++-- ...ection-provider-oauth-flow.service.spec.ts | 18 +++-- .../connection-provider-oauth-flow.service.ts | 20 +++--- .../jobs/logic-function-trigger.job.ts | 49 +++++++------ .../triggers/cron/cron-trigger.cron.job.ts | 14 ++-- .../call-database-event-trigger-jobs.job.ts | 20 +----- .../message-queue/drivers/bullmq.driver.ts | 71 +++++++++++++++---- .../message-queue-driver.interface.ts | 6 ++ .../message-queue/drivers/sync.driver.ts | 24 +++++++ .../services/message-queue.service.ts | 8 +++ .../server-route-trigger.service.spec.ts | 12 ++-- .../server-route-trigger.service.ts | 14 ++-- 12 files changed, 164 insertions(+), 106 deletions(-) diff --git a/packages/twenty-server/src/engine/core-modules/application/application-install/application-install.service.ts b/packages/twenty-server/src/engine/core-modules/application/application-install/application-install.service.ts index 51e3a442f9..76a40d354a 100644 --- a/packages/twenty-server/src/engine/core-modules/application/application-install/application-install.service.ts +++ b/packages/twenty-server/src/engine/core-modules/application/application-install/application-install.service.ts @@ -557,15 +557,13 @@ export class ApplicationInstallService { ); if (!shouldRunSynchronously) { - await this.messageQueueService.add( + await this.messageQueueService.add( LogicFunctionTriggerJob.name, - [ - { - logicFunctionId: flatLogicFunction.id, - workspaceId, - payload, - }, - ], + { + logicFunctionId: flatLogicFunction.id, + workspaceId, + payload, + }, { retryLimit: 3 }, ); return; diff --git a/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-oauth-flow.service.spec.ts b/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-oauth-flow.service.spec.ts index 05c4d3ea48..2d1956493f 100644 --- a/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-oauth-flow.service.spec.ts +++ b/packages/twenty-server/src/engine/core-modules/application/connection-provider/__tests__/connection-provider-oauth-flow.service.spec.ts @@ -501,17 +501,15 @@ describe('ConnectionProviderOAuthFlowService', () => { ); expect(messageQueueService.add).toHaveBeenCalledWith( 'LogicFunctionTriggerJob', - [ - { - logicFunctionId: 'logic-function-1', - workspaceId: 'workspace-1', - payload: { - connectionProviderId: 'provider-1', - connectionProviderName: 'linear', - connectedAccountId: result.connectedAccountId, - }, + { + logicFunctionId: 'logic-function-1', + workspaceId: 'workspace-1', + payload: { + connectionProviderId: 'provider-1', + connectionProviderName: 'linear', + connectedAccountId: result.connectedAccountId, }, - ], + }, { retryLimit: 3 }, ); }); diff --git a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-oauth-flow.service.ts b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-oauth-flow.service.ts index 48a548a952..efce22153e 100644 --- a/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-oauth-flow.service.ts +++ b/packages/twenty-server/src/engine/core-modules/application/connection-provider/connection-provider-oauth-flow.service.ts @@ -252,19 +252,17 @@ export class ConnectionProviderOAuthFlowService { ); } - await this.messageQueueService.add( + await this.messageQueueService.add( LogicFunctionTriggerJob.name, - [ - { - logicFunctionId: flatLogicFunction.id, - workspaceId, - payload: { - connectionProviderId: provider.id, - connectionProviderName: provider.name, - connectedAccountId, - }, + { + logicFunctionId: flatLogicFunction.id, + workspaceId, + payload: { + connectionProviderId: provider.id, + connectionProviderName: provider.name, + connectedAccountId, }, - ], + }, { retryLimit: 3 }, ); } catch (error) { diff --git a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/jobs/logic-function-trigger.job.ts b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/jobs/logic-function-trigger.job.ts index 61f5d855a3..f3890cb2a9 100644 --- a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/jobs/logic-function-trigger.job.ts +++ b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/jobs/logic-function-trigger.job.ts @@ -27,30 +27,33 @@ export class LogicFunctionTriggerJob { ) {} @Process(LogicFunctionTriggerJob.name) - async handle(logicFunctionPayloads: LogicFunctionTriggerJobData[]) { - await Promise.all( - logicFunctionPayloads.map(async (logicFunctionPayload) => { - try { - await this.logicFunctionExecutorService.execute({ - logicFunctionId: logicFunctionPayload.logicFunctionId, - workspaceId: logicFunctionPayload.workspaceId, - payload: logicFunctionPayload.payload ?? {}, - userId: logicFunctionPayload.userId, - userWorkspaceId: logicFunctionPayload.userWorkspaceId, - }); - } catch (error) { - // A stopped application must not fail the job: failing would make - // the queue retry an execution that is intentionally blocked. - if ( - error instanceof LogicFunctionException && - error.code === LogicFunctionExceptionCode.LOGIC_FUNCTION_DISABLED - ) { - return; - } + async handle( + jobData: LogicFunctionTriggerJobData | LogicFunctionTriggerJobData[], + ) { + // Jobs enqueued in version <=2.24.x carry arrays, remove this case once those jobs are drained + const logicFunctionPayloads = Array.isArray(jobData) ? jobData : [jobData]; - throw error; + for (const logicFunctionPayload of logicFunctionPayloads) { + try { + await this.logicFunctionExecutorService.execute({ + logicFunctionId: logicFunctionPayload.logicFunctionId, + workspaceId: logicFunctionPayload.workspaceId, + payload: logicFunctionPayload.payload ?? {}, + userId: logicFunctionPayload.userId, + userWorkspaceId: logicFunctionPayload.userWorkspaceId, + }); + } catch (error) { + // A stopped application must not fail the job: failing would make + // the queue retry an execution that is intentionally blocked. + if ( + error instanceof LogicFunctionException && + error.code === LogicFunctionExceptionCode.LOGIC_FUNCTION_DISABLED + ) { + continue; } - }), - ); + + throw error; + } + } } } diff --git a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/cron/cron-trigger.cron.job.ts b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/cron/cron-trigger.cron.job.ts index 37adbad8d6..f9a891e06b 100644 --- a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/cron/cron-trigger.cron.job.ts +++ b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/cron/cron-trigger.cron.job.ts @@ -85,15 +85,13 @@ export class CronTriggerCronJob { continue; } - await this.messageQueueService.add( + await this.messageQueueService.add( LogicFunctionTriggerJob.name, - [ - { - logicFunctionId: logicFunction.id, - workspaceId: activeWorkspace.id, - payload: {}, - }, - ], + { + logicFunctionId: logicFunction.id, + workspaceId: activeWorkspace.id, + payload: {}, + }, { retryLimit: 10 }, ); } diff --git a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/database-event/call-database-event-trigger-jobs.job.ts b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/database-event/call-database-event-trigger-jobs.job.ts index 7c340de4f2..dbdd5d4439 100644 --- a/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/database-event/call-database-event-trigger-jobs.job.ts +++ b/packages/twenty-server/src/engine/core-modules/logic-function/logic-function-trigger/triggers/database-event/call-database-event-trigger-jobs.job.ts @@ -1,4 +1,3 @@ -import chunk from 'lodash.chunk'; import { isDefined } from 'twenty-shared/utils'; import type { ObjectRecordEvent } from 'twenty-shared/database-events'; @@ -16,8 +15,6 @@ import { import { WorkspaceCacheService } from 'src/engine/workspace-cache/services/workspace-cache.service'; import { WorkspaceEventBatch } from 'src/engine/workspace-event-emitter/types/workspace-event-batch.type'; -const DATABASE_EVENT_JOBS_CHUNK_SIZE = 20; - @Processor(MessageQueue.triggerQueue) export class CallDatabaseEventTriggerJobsJob { constructor( @@ -59,22 +56,11 @@ export class CallDatabaseEventTriggerJobsJob { workspaceEventBatch, }); - if (logicFunctionPayloads.length === 0) { - return; - } - - const logicFunctionPayloadsChunks = chunk( + await this.messageQueueService.bulkAdd( + LogicFunctionTriggerJob.name, logicFunctionPayloads, - DATABASE_EVENT_JOBS_CHUNK_SIZE, + { retryLimit: 3 }, ); - - for (const logicFunctionPayloadsChunk of logicFunctionPayloadsChunks) { - await this.messageQueueService.add( - LogicFunctionTriggerJob.name, - logicFunctionPayloadsChunk, - { retryLimit: 3 }, - ); - } } private shouldTriggerJob({ diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts b/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts index 177e33e042..0aefd7278d 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts @@ -327,6 +327,30 @@ export class BullMQDriver ); } + private buildJobsOptions({ + queueName, + options, + }: { + queueName: MessageQueue; + options?: QueueJobOptions; + }): JobsOptions { + return { + // We suffix the id with V4() to make sure ids are unique so we can add a waiting job when a job related with the same option.id is running + jobId: options?.id ? `${options.id}-${v4()}` : undefined, + priority: options?.priority ?? MESSAGE_QUEUE_PRIORITY[queueName], + attempts: 1 + (options?.retryLimit || 0), + removeOnComplete: { + age: QUEUE_RETENTION.completedMaxAge, + count: QUEUE_RETENTION.completedMaxCount, + }, + removeOnFail: { + age: QUEUE_RETENTION.failedMaxAge, + count: QUEUE_RETENTION.failedMaxCount, + }, + delay: options?.delay, + }; + } + async add( queueName: MessageQueue, jobName: string, @@ -352,24 +376,43 @@ export class BullMQDriver } } - const queueOptions: JobsOptions = { - jobId: options?.id ? `${options.id}-${v4()}` : undefined, // We add V4() to id to make sure ids are uniques so we can add a waiting job when a job related with the same option.id is running - priority: options?.priority ?? MESSAGE_QUEUE_PRIORITY[queueName], - attempts: 1 + (options?.retryLimit || 0), - removeOnComplete: { - age: QUEUE_RETENTION.completedMaxAge, - count: QUEUE_RETENTION.completedMaxCount, - }, - removeOnFail: { - age: QUEUE_RETENTION.failedMaxAge, - count: QUEUE_RETENTION.failedMaxCount, - }, - delay: options?.delay, - }; + const queueOptions = this.buildJobsOptions({ queueName, options }); await this.queueMap[queueName].add(jobName, data, queueOptions); } + async bulkAdd( + queueName: MessageQueue, + jobName: string, + dataItems: T[], + options?: QueueJobOptions, + ): Promise { + if (!this.queueMap[queueName]) { + throw new Error( + `Queue ${queueName} is not registered, make sure you have added it as a queue provider`, + ); + } + + if (dataItems.length === 0) { + return; + } + + const queueOptions = this.buildJobsOptions({ queueName, options }); + + await this.queueMap[queueName].addBulk( + dataItems.map((data, index) => ({ + name: jobName, + data, + opts: { + ...queueOptions, + jobId: queueOptions.jobId + ? `${queueOptions.jobId}-${index}` + : undefined, + }, + })), + ); + } + async getInFlightJobs( queueName: MessageQueue, ): Promise[]> { diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/drivers/interfaces/message-queue-driver.interface.ts b/packages/twenty-server/src/engine/core-modules/message-queue/drivers/interfaces/message-queue-driver.interface.ts index 8107f31c68..1ee089e480 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/drivers/interfaces/message-queue-driver.interface.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/drivers/interfaces/message-queue-driver.interface.ts @@ -14,6 +14,12 @@ export interface MessageQueueDriver { data: T, options?: QueueJobOptions, ): Promise; + bulkAdd( + queueName: MessageQueue, + jobName: string, + dataItems: T[], + options?: QueueJobOptions, + ): Promise; work( queueName: MessageQueue, handler: ({ data, id }: { data: T; id: string }) => Promise | void, diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/drivers/sync.driver.ts b/packages/twenty-server/src/engine/core-modules/message-queue/drivers/sync.driver.ts index 7e2447f621..47537815d6 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/drivers/sync.driver.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/drivers/sync.driver.ts @@ -1,5 +1,7 @@ import { Logger } from '@nestjs/common'; +import { isDefined } from 'twenty-shared/utils'; + import { type MessageQueueDriver } from 'src/engine/core-modules/message-queue/drivers/interfaces/message-queue-driver.interface'; import { type MessageQueueJob, @@ -25,6 +27,28 @@ export class SyncDriver implements MessageQueueDriver { await this.processJob(queueName, { id: '', name: jobName, data }); } + async bulkAdd( + queueName: MessageQueue, + jobName: string, + dataItems: T[], + ): Promise { + let firstError: unknown = undefined; + + // Each payload is an independent job in BullMQ, so a failing one must not + // prevent the others from being processed + for (const data of dataItems) { + try { + await this.processJob(queueName, { id: '', name: jobName, data }); + } catch (error) { + firstError = firstError ?? error; + } + } + + if (isDefined(firstError)) { + throw firstError; + } + } + async addCron({ queueName, jobName, diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/services/message-queue.service.ts b/packages/twenty-server/src/engine/core-modules/message-queue/services/message-queue.service.ts index 71bc29801d..1f36a63c49 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/services/message-queue.service.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/services/message-queue.service.ts @@ -38,6 +38,14 @@ export class MessageQueueService { return this.driver.add(this.queueName, jobName, data, options); } + bulkAdd( + jobName: string, + dataItems: T[], + options?: QueueJobOptions, + ): Promise { + return this.driver.bulkAdd(this.queueName, jobName, dataItems, options); + } + getInFlightJobs(): Promise< InFlightQueueJob[] > { diff --git a/packages/twenty-server/src/engine/core-modules/server-route-trigger/__tests__/server-route-trigger.service.spec.ts b/packages/twenty-server/src/engine/core-modules/server-route-trigger/__tests__/server-route-trigger.service.spec.ts index 529ba9f28f..75593f1e57 100644 --- a/packages/twenty-server/src/engine/core-modules/server-route-trigger/__tests__/server-route-trigger.service.spec.ts +++ b/packages/twenty-server/src/engine/core-modules/server-route-trigger/__tests__/server-route-trigger.service.spec.ts @@ -146,13 +146,11 @@ describe('ServerRouteTriggerService', () => { ); expect(messageQueueService.add).toHaveBeenCalledWith( LogicFunctionTriggerJob.name, - [ - { - logicFunctionId: 'target-id', - workspaceId: 'target-ws', - payload: { from: 'resolver' }, - }, - ], + { + logicFunctionId: 'target-id', + workspaceId: 'target-ws', + payload: { from: 'resolver' }, + }, { retryLimit: 3 }, ); expect(result).toEqual( diff --git a/packages/twenty-server/src/engine/core-modules/server-route-trigger/server-route-trigger.service.ts b/packages/twenty-server/src/engine/core-modules/server-route-trigger/server-route-trigger.service.ts index 33ecaafd20..d65ab3c5de 100644 --- a/packages/twenty-server/src/engine/core-modules/server-route-trigger/server-route-trigger.service.ts +++ b/packages/twenty-server/src/engine/core-modules/server-route-trigger/server-route-trigger.service.ts @@ -189,15 +189,13 @@ export class ServerRouteTriggerService { applicationRegistrationId, }); - await this.messageQueueService.add( + await this.messageQueueService.add( LogicFunctionTriggerJob.name, - [ - { - logicFunctionId: logicFunction.id, - workspaceId, - payload, - }, - ], + { + logicFunctionId: logicFunction.id, + workspaceId, + payload, + }, { retryLimit: QUEUED_TARGET_RETRY_LIMIT }, );