diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/constants/ai-stream-lock-duration.constant.ts b/packages/twenty-server/src/engine/core-modules/message-queue/constants/ai-stream-lock-duration.constant.ts deleted file mode 100644 index c739aa8943..0000000000 --- a/packages/twenty-server/src/engine/core-modules/message-queue/constants/ai-stream-lock-duration.constant.ts +++ /dev/null @@ -1 +0,0 @@ -export const AI_STREAM_LOCK_DURATION_MS = 10 * 60 * 1000; 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 0aefd7278d..8fb5080f4d 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 @@ -30,7 +30,7 @@ import { import { type MessageQueueWorkerOptions } from 'src/engine/core-modules/message-queue/interfaces/message-queue-worker-options.interface'; import { QUEUE_RETENTION } from 'src/engine/core-modules/message-queue/constants/queue-retention.constants'; -import { MESSAGE_QUEUE_PRIORITY } from 'src/engine/core-modules/message-queue/message-queue-priority.constant'; +import { MESSAGE_QUEUE_WORKER_CONFIG } from 'src/engine/core-modules/message-queue/message-queue-worker-config.constant'; import { type MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; import { getJobKey } from 'src/engine/core-modules/message-queue/utils/get-job-key.util'; import { type MetricsService } from 'src/engine/core-modules/metrics/metrics.service'; @@ -337,7 +337,8 @@ export class BullMQDriver 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], + priority: + options?.priority ?? MESSAGE_QUEUE_WORKER_CONFIG[queueName].priority, attempts: 1 + (options?.retryLimit || 0), removeOnComplete: { age: QUEUE_RETENTION.completedMaxAge, diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-priority.constant.ts b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-priority.constant.ts deleted file mode 100644 index 5a8a158153..0000000000 --- a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-priority.constant.ts +++ /dev/null @@ -1,21 +0,0 @@ -import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; - -export const MESSAGE_QUEUE_PRIORITY = { - [MessageQueue.billingQueue]: 1, - [MessageQueue.entityEventsToDbQueue]: 1, - [MessageQueue.emailQueue]: 1, - [MessageQueue.workflowQueue]: 2, - [MessageQueue.webhookQueue]: 2, - [MessageQueue.messagingQueue]: 2, - [MessageQueue.delayedJobsQueue]: 3, - [MessageQueue.calendarQueue]: 4, - [MessageQueue.contactCreationQueue]: 4, - [MessageQueue.taskAssignedQueue]: 4, - [MessageQueue.logicFunctionQueue]: 4, - [MessageQueue.workspaceQueue]: 5, - [MessageQueue.triggerQueue]: 5, - [MessageQueue.deleteCascadeQueue]: 6, - [MessageQueue.cronQueue]: 7, - [MessageQueue.aiQueue]: 5, - [MessageQueue.aiStreamQueue]: 2, -}; diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-worker-config.constant.ts b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-worker-config.constant.ts new file mode 100644 index 0000000000..624b2e96d0 --- /dev/null +++ b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-worker-config.constant.ts @@ -0,0 +1,180 @@ +import { type MessageQueueWorkerOptions } from 'src/engine/core-modules/message-queue/interfaces/message-queue-worker-options.interface'; +import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; + +// Single source of truth to pilot worker behavior per queue. Every value is +// explicit on purpose (Record + Required): adding a queue without declaring +// its full configuration is a compile error, and no queue silently relies on +// implicit BullMQ defaults. +// +// priority: applied when enqueuing, lower value is processed first +// concurrency: max jobs processed in parallel per worker process +// lockDuration: ms a job may run before BullMQ considers it stalled +// maxStalledCount: times a stalled job is re-queued before failing permanently +// boundedShutdownDrain: on shutdown, abort still-active jobs after +// AI_STREAM_SHUTDOWN_DRAIN_MS instead of waiting for them to finish + +export type MessageQueueWorkerConfig = { + priority: number; + workerOptions: Required; +}; + +export const MESSAGE_QUEUE_WORKER_CONFIG: Record< + MessageQueue, + MessageQueueWorkerConfig +> = { + [MessageQueue.taskAssignedQueue]: { + priority: 4, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.messagingQueue]: { + priority: 2, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.webhookQueue]: { + priority: 2, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.cronQueue]: { + priority: 7, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.emailQueue]: { + priority: 1, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.calendarQueue]: { + priority: 4, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.contactCreationQueue]: { + priority: 4, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.billingQueue]: { + priority: 1, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.workspaceQueue]: { + priority: 5, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.entityEventsToDbQueue]: { + priority: 1, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.workflowQueue]: { + priority: 2, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.delayedJobsQueue]: { + priority: 3, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.deleteCascadeQueue]: { + priority: 6, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.logicFunctionQueue]: { + priority: 4, + workerOptions: { + concurrency: 10, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.triggerQueue]: { + priority: 5, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.aiQueue]: { + priority: 5, + workerOptions: { + concurrency: 1, + lockDuration: 30_000, + maxStalledCount: 1, + boundedShutdownDrain: false, + }, + }, + [MessageQueue.aiStreamQueue]: { + priority: 2, + workerOptions: { + concurrency: 20, + // 10 minutes: a stream job holds its lock for the whole stream duration + lockDuration: 600_000, + // A stalled stream cannot be resumed client-side, never re-queue it + maxStalledCount: 0, + boundedShutdownDrain: true, + }, + }, +}; diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-worker-options.constant.ts b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-worker-options.constant.ts deleted file mode 100644 index 8d5ee0defa..0000000000 --- a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-worker-options.constant.ts +++ /dev/null @@ -1,15 +0,0 @@ -import { AI_STREAM_LOCK_DURATION_MS } from 'src/engine/core-modules/message-queue/constants/ai-stream-lock-duration.constant'; -import { type MessageQueueWorkerOptions } from 'src/engine/core-modules/message-queue/interfaces/message-queue-worker-options.interface'; -import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; - -export const QUEUE_WORKER_OPTIONS: Partial< - Record -> = { - [MessageQueue.aiStreamQueue]: { - concurrency: 20, - lockDuration: AI_STREAM_LOCK_DURATION_MS, - maxStalledCount: 0, - boundedShutdownDrain: true, - }, - [MessageQueue.logicFunctionQueue]: { concurrency: 10 }, -}; diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.explorer.ts b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.explorer.ts index 6030ccf8e6..333e4a2f79 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.explorer.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.explorer.ts @@ -18,8 +18,8 @@ import { type MessageQueueWorkerOptions } from 'src/engine/core-modules/message- import { ExceptionHandlerService } from 'src/engine/core-modules/exception-handler/exception-handler.service'; import { MessageQueueMetadataAccessor } from 'src/engine/core-modules/message-queue/message-queue-metadata.accessor'; import { type MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; +import { MESSAGE_QUEUE_WORKER_CONFIG } from 'src/engine/core-modules/message-queue/message-queue-worker-config.constant'; import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; -import { QUEUE_WORKER_OPTIONS } from 'src/engine/core-modules/message-queue/message-queue-worker-options.constant'; import { getQueueToken } from 'src/engine/core-modules/message-queue/utils/get-queue-token.util'; import { shouldCreateWorkerForQueue } from 'src/engine/core-modules/message-queue/utils/should-create-worker-for-queue.util'; import { TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; @@ -103,7 +103,7 @@ export class MessageQueueExplorer implements OnModuleInit { this.handleProcessorGroupCollection( processorGroupCollection, messageQueueService, - QUEUE_WORKER_OPTIONS[queueName], + MESSAGE_QUEUE_WORKER_CONFIG[queueName].workerOptions, ); } } @@ -181,7 +181,7 @@ export class MessageQueueExplorer implements OnModuleInit { private handleProcessorGroupCollection( processorGroupCollection: ProcessorGroup[], queue: MessageQueueService, - options?: MessageQueueWorkerOptions, + options: MessageQueueWorkerOptions, ) { queue.work(async (job) => { for (const processorGroup of processorGroupCollection) {