From 304b0037ec949bf53038ef42261ac8cf467f8c90 Mon Sep 17 00:00:00 2001 From: devlikepro Date: Sun, 2 Nov 2025 14:42:25 +0700 Subject: [PATCH] [core] ChatWoot - queue management --- src/apps/chatwoot/chatwoot.module.ts | 4 + src/apps/chatwoot/cli/cmd.queue.ts | 52 ++++++++++ src/apps/chatwoot/cli/program.a.ts | 4 +- src/apps/chatwoot/cli/program.messages.ts | 15 ++- src/apps/chatwoot/cli/program.queue.ts | 66 +++++++++++++ src/apps/chatwoot/cli/types.ts | 2 + src/apps/chatwoot/consumers/inbox/commands.ts | 11 +-- .../chatwoot/consumers/task/messages.pull.ts | 67 ++++++++++--- src/apps/chatwoot/di/DIContainer.ts | 1 + src/apps/chatwoot/dto/config.dto.ts | 4 + src/apps/chatwoot/i18n/locales/en-US.yaml | 48 ++++++++- .../chatwoot/services/ChatWootQueueService.ts | 47 +++++---- .../services/ChatWootScheduleService.ts | 44 +++++---- .../services/ChatWootWAHAQueueService.ts | 29 ++---- src/apps/chatwoot/services/QueueManager.ts | 98 +++++++++++++++++++ src/apps/chatwoot/services/QueueRegistry.ts | 78 +++++++++++++++ 16 files changed, 483 insertions(+), 87 deletions(-) create mode 100644 src/apps/chatwoot/cli/cmd.queue.ts create mode 100644 src/apps/chatwoot/cli/program.queue.ts create mode 100644 src/apps/chatwoot/services/QueueManager.ts create mode 100644 src/apps/chatwoot/services/QueueRegistry.ts diff --git a/src/apps/chatwoot/chatwoot.module.ts b/src/apps/chatwoot/chatwoot.module.ts index 87a8d188..bac6db69 100644 --- a/src/apps/chatwoot/chatwoot.module.ts +++ b/src/apps/chatwoot/chatwoot.module.ts @@ -26,11 +26,13 @@ import { WAHASessionStatusConsumer } from './consumers/waha/session.status'; import { ChatWootQueueService } from './services/ChatWootQueueService'; import { ChatWootScheduleService } from './services/ChatWootScheduleService'; import { ChatWootWAHAQueueService } from './services/ChatWootWAHAQueueService'; +import { QueueRegistry } from './services/QueueRegistry'; import { ChatWootConversationCreatedConsumer } from './consumers/inbox/conversation_created'; import { ChatWootConversationStatusChangedConsumer } from '@waha/apps/chatwoot/consumers/inbox/conversation_status_changed'; import { TaskContactsPullConsumer } from '@waha/apps/chatwoot/consumers/task/contacts.pull'; import { TaskMessagesPullConsumer } from '@waha/apps/chatwoot/consumers/task/messages.pull'; import { BullModule } from '@nestjs/bullmq'; +import { QueueManager } from '@waha/apps/chatwoot/services/QueueManager'; const CONTROLLERS = [ChatwootWebhookController, ChatwootLocalesController]; @@ -130,6 +132,8 @@ const PROVIDERS = [ ChatWootQueueService, ChatWootScheduleService, ChatWootAppService, + QueueRegistry, + QueueManager, ]; export const ChatWootExports = { diff --git a/src/apps/chatwoot/cli/cmd.queue.ts b/src/apps/chatwoot/cli/cmd.queue.ts new file mode 100644 index 00000000..64b94c54 --- /dev/null +++ b/src/apps/chatwoot/cli/cmd.queue.ts @@ -0,0 +1,52 @@ +import * as lodash from 'lodash'; +import { QueueManager } from '@waha/apps/chatwoot/services/QueueManager'; +import { Locale } from '@waha/apps/chatwoot/i18n/locale'; +import { Conversation } from '@waha/apps/chatwoot/client/Conversation'; +import { QueueRegistry } from '@waha/apps/chatwoot/services/QueueRegistry'; + +function repr(name: string): string { + // No chatwoot in the name always (it start with chatwoot) + name = name.replace('chatwoot.', ''); + // No waha in the name + name = name.replace('waha |', 'whatsapp |'); + return name; +} + +export interface QueueCommandContext { + queues: { + registry: QueueRegistry; + }; + l: Locale; + conversation: Conversation; +} + +export async function QueueStatus(ctx: QueueCommandContext, name: string) { + const manager = new QueueManager(ctx.queues.registry); + const names = manager.resolve(name); + let result = await manager.status(names); + for (const status of result) { + status.name = repr(status.name); + } + // locked: true - last + result = lodash.sortBy(result, [(x) => !!x.locked, 'name']); + const msg = ctx.l.r('cli.cmd.queue.status.result', { + queues: result, + }); + await ctx.conversation.incoming(msg); +} + +export async function QueueStart(ctx: QueueCommandContext, name: string) { + const manager = new QueueManager(ctx.queues.registry); + const names = manager.resolve(name); + await manager.resume(names); + const msg = ctx.l.r('cli.cmd.queue.resumed'); + await ctx.conversation.activity(msg); +} + +export async function QueueStop(ctx: QueueCommandContext, name?: string) { + const manager = new QueueManager(ctx.queues.registry); + const names = manager.resolve(name); + await manager.pause(names); + const msg = ctx.l.r('cli.cmd.queue.paused'); + await ctx.conversation.activity(msg); +} diff --git a/src/apps/chatwoot/cli/program.a.ts b/src/apps/chatwoot/cli/program.a.ts index f8200751..79baf89f 100644 --- a/src/apps/chatwoot/cli/program.a.ts +++ b/src/apps/chatwoot/cli/program.a.ts @@ -9,6 +9,7 @@ import { AddSessionCommand } from '@waha/apps/chatwoot/cli/program.session'; import { AddContactsCommand } from '@waha/apps/chatwoot/cli/program.contacts'; import { AddMessagesCommand } from '@waha/apps/chatwoot/cli/program.messages'; import { AddServerCommand } from '@waha/apps/chatwoot/cli/program.server'; +import { AddQueueCommand } from '@waha/apps/chatwoot/cli/program.queue'; function Program(ctx: CommandContext, output: OutputConfiguration) { const l = ctx.l; @@ -62,7 +63,8 @@ export function BuildProgram( const program = Program(ctx, output); AddSessionCommand(program, ctx); AddContactsCommand(program, ctx); - AddMessagesCommand(program, ctx); + AddMessagesCommand(program, ctx, commands.queue); + AddQueueCommand(program, ctx, commands.queue); AddServerCommand(program, ctx, commands.server); return program; } diff --git a/src/apps/chatwoot/cli/program.messages.ts b/src/apps/chatwoot/cli/program.messages.ts index 198142ae..f02ebacc 100644 --- a/src/apps/chatwoot/cli/program.messages.ts +++ b/src/apps/chatwoot/cli/program.messages.ts @@ -15,7 +15,11 @@ import { MessagesPullOptions } from '@waha/apps/chatwoot/consumers/task/messages import { JobsOptions } from 'bullmq'; import { JobDataTimeout } from '@waha/apps/app_sdk/AppConsumer'; -export function AddMessagesCommand(program: Command, ctx: CommandContext) { +export function AddMessagesCommand( + program: Command, + ctx: CommandContext, + queue: boolean, +) { const l = ctx.l; const SyncGroup = l.r('cli.cmd.root.sub.sync'); program.commandsGroup(SyncGroup); @@ -53,6 +57,7 @@ export function AddMessagesCommand(program: Command, ctx: CommandContext) { .option('-s, --status', l.r('cli.cmd.messages.pull.option.status')) .option('--bc, --broadcast', l.r('cli.cmd.messages.pull.option.broadcast')) .option('-m, --media', l.r('cli.cmd.messages.pull.option.media')) + .option('--pause', l.r('cli.cmd.messages.pull.option.pause')) .option( '-b, --batch ', l.r('cli.cmd.messages.pull.option.batch'), @@ -92,6 +97,13 @@ export function AddMessagesCommand(program: Command, ctx: CommandContext) { start = tmp; } + if (opts.pause && !queue) { + await ctx.conversation.incoming( + l.r('cli.cmd.messages.pull.error.pause-no-queue'), + ); + return; + } + const options: MessagesPullOptions = { chat: opts.chat, progress: opts.progress, @@ -101,6 +113,7 @@ export function AddMessagesCommand(program: Command, ctx: CommandContext) { }, media: opts.media, force: opts.force, + pause: !!opts.pause, timeout: { media: opts.timeoutMedia, }, diff --git a/src/apps/chatwoot/cli/program.queue.ts b/src/apps/chatwoot/cli/program.queue.ts new file mode 100644 index 00000000..10db30de --- /dev/null +++ b/src/apps/chatwoot/cli/program.queue.ts @@ -0,0 +1,66 @@ +import { Argument, Command } from 'commander'; +import { CommandContext } from '@waha/apps/chatwoot/cli/types'; +import { CommandDisabled } from '@waha/apps/chatwoot/cli/cmd.disabled'; +import { + QueueStart, + QueueStatus, + QueueStop, +} from '@waha/apps/chatwoot/cli/cmd.queue'; + +export function AddQueueCommand( + program: Command, + ctx: CommandContext, + enabled: boolean, +) { + const l = ctx.l; + const QueueGroup = l.r('cli.cmd.root.sub.queue'); + program.commandsGroup(QueueGroup); + + program + .command('queue', { hidden: !enabled }) + .alias('q') + .summary(l.r('cli.cmd.queue.summary')) + .description(l.r('cli.cmd.queue.description')) + .helpGroup(QueueGroup) + .addArgument( + new Argument('[action]', l.r('cli.cmd.queue.action.description')) + .choices(['status', 'start', 'stop', 'help']) + .default('help'), + ) + .addArgument( + new Argument('[name]', l.r('cli.cmd.queue.argument.name')).default(''), + ) + .action(async function (this: Command, action: string, name: string) { + if (!action) { + this.outputHelp(); + return; + } + + if (action === 'help') { + this.outputHelp(); + return; + } + + if (!enabled) { + await CommandDisabled(ctx, 'queue'); + return; + } + + if (action === 'status') { + await QueueStatus(ctx, name); + return; + } + + if (action === 'start') { + await QueueStart(ctx, name); + return; + } + + if (action === 'stop') { + await QueueStop(ctx, name); + return; + } + + this.outputHelp(); + }); +} diff --git a/src/apps/chatwoot/cli/types.ts b/src/apps/chatwoot/cli/types.ts index e0c3d93c..5386ae0e 100644 --- a/src/apps/chatwoot/cli/types.ts +++ b/src/apps/chatwoot/cli/types.ts @@ -4,6 +4,7 @@ import { Conversation } from '../client/Conversation'; import { ILogger } from '@waha/apps/app_sdk/ILogger'; import { FlowProducer, Queue } from 'bullmq'; import { InboxData } from '@waha/apps/chatwoot/consumers/types'; +import { QueueRegistry } from '@waha/apps/chatwoot/services/QueueRegistry'; export interface CommandContext { data: InboxData; @@ -12,6 +13,7 @@ export interface CommandContext { waha: WAHASelf; conversation: Conversation; queues: { + registry: QueueRegistry; contactsPull: Queue; messagesPull: Queue; }; diff --git a/src/apps/chatwoot/consumers/inbox/commands.ts b/src/apps/chatwoot/consumers/inbox/commands.ts index 31c71c8d..82d72253 100644 --- a/src/apps/chatwoot/consumers/inbox/commands.ts +++ b/src/apps/chatwoot/consumers/inbox/commands.ts @@ -14,6 +14,7 @@ import { TKey } from '@waha/apps/chatwoot/i18n/templates'; import { CommandPrefix, runText } from '@waha/apps/chatwoot/cli'; import { CommandContext } from '@waha/apps/chatwoot/cli/types'; import { IsCommandsChat } from '@waha/apps/chatwoot/client/ids'; +import { QueueRegistry } from '@waha/apps/chatwoot/services/QueueRegistry'; @Processor(QueueName.INBOX_COMMANDS, { concurrency: JOB_CONCURRENCY }) export class ChatWootInboxCommandsConsumer extends ChatWootInboxMessageConsumer { @@ -21,10 +22,7 @@ export class ChatWootInboxCommandsConsumer extends ChatWootInboxMessageConsumer protected readonly manager: SessionManager, log: PinoLogger, rmutex: RMutexService, - @InjectQueue(QueueName.TASK_CONTACTS_PULL) - private readonly contactsPullQueue: Queue, - @InjectQueue(QueueName.TASK_MESSAGES_PULL) - private readonly messagesPullQueue: Queue, + private readonly queueRegistry: QueueRegistry, @InjectFlowProducer(FlowProducerName.MESSAGES_PULL_FLOW) private readonly messagesPullFlow: FlowProducer, ) { @@ -61,8 +59,9 @@ export class ChatWootInboxCommandsConsumer extends ChatWootInboxMessageConsumer waha: container.WAHASelf(), conversation: conversation, queues: { - contactsPull: this.contactsPullQueue, - messagesPull: this.messagesPullQueue, + registry: this.queueRegistry, + contactsPull: this.queueRegistry.queue(QueueName.TASK_CONTACTS_PULL), + messagesPull: this.queueRegistry.queue(QueueName.TASK_MESSAGES_PULL), }, flows: { messagesPull: this.messagesPullFlow, diff --git a/src/apps/chatwoot/consumers/task/messages.pull.ts b/src/apps/chatwoot/consumers/task/messages.pull.ts index 496e18b3..c80a879e 100644 --- a/src/apps/chatwoot/consumers/task/messages.pull.ts +++ b/src/apps/chatwoot/consumers/task/messages.pull.ts @@ -35,6 +35,7 @@ import { IgnoreJidConfig, isNullJid, JidFilter } from '@waha/core/utils/jids'; import { ErrorRenderer } from '@waha/apps/chatwoot/error/ErrorRenderer'; import { IsCommandsChat } from '@waha/apps/chatwoot/client/ids'; import { CommandPrefix } from '@waha/apps/chatwoot/cli'; +import { QueueManager } from '@waha/apps/chatwoot/services/QueueManager'; export enum ChatID { ALL = 'all', @@ -50,6 +51,7 @@ export type MessagesPullOptions = { end: number; }; force: boolean; + pause: boolean; timeout: { media: number; }; @@ -84,7 +86,12 @@ function total(progress: Progress) { */ @Processor(QueueName.TASK_MESSAGES_PULL, { concurrency: JOB_CONCURRENCY }) export class TaskMessagesPullConsumer extends ChatWootTaskConsumer { - constructor(manager: SessionManager, log: PinoLogger, rmutex: RMutexService) { + constructor( + manager: SessionManager, + log: PinoLogger, + rmutex: RMutexService, + protected queueManager: QueueManager, + ) { super(manager, log, rmutex, TaskMessagesPullConsumer.name); } @@ -103,6 +110,7 @@ export class TaskMessagesPullConsumer extends ChatWootTaskConsumer { const waha = container.WAHASelf(); const session = new WAHASessionAPI(job.data.session, waha); const handler = new MessagesPullHandler( + this.queueManager, container, signal, container.Logger(), @@ -110,11 +118,9 @@ export class TaskMessagesPullConsumer extends ChatWootTaskConsumer { session, locale, ); - if (options.chat === ChatID.SUMMARY) { - return await handler.summary(options, job); - } - - return await handler.handle(options, job); + await handler.start(options, job); + job = await handler.handle(options, job); + return await handler.end(options, job); } } @@ -122,6 +128,7 @@ class MessagesPullHandler { private activity: TaskActivity; constructor( + protected queueManager: QueueManager, protected container: DIContainer, protected signal: AbortSignal, protected logger: ILogger, @@ -158,14 +165,42 @@ class MessagesPullHandler { return progress; } - async summary(options: MessagesPullOptions, job: Job): Promise { + async start(options: MessagesPullOptions, job: Job) { + const { ignored, processed, unprocessed, failed } = + await job.getDependenciesCount(); + const total = ignored + processed + unprocessed + failed; + const hasChildren = total > 0; + if (hasChildren) { + return; + } + + // Perform some actions before any job started + if (options.pause) { + await this.queueManager.pause(); + await this.activity.queue(true); + } + } + + async end(options: MessagesPullOptions, job: Job) { + if (job.parentKey) { + return await this.summaryProgress(job); + } + // Perform some actions after all jobs are done const progress = await this.summaryProgress(job); await job.updateProgress(progress); await this.activity.completed(progress, options); + if (options.pause) { + await this.queueManager.resume(); + await this.activity.queue(false); + } return progress; } - async handle(options: MessagesPullOptions, job: Job): Promise { + async handle(options: MessagesPullOptions, job: Job): Promise { + if (options.chat === ChatID.SUMMARY) { + // Do nothing for "summary" chat + return job; + } const jids = new JidFilter(options.ignore); let progress = lodash.merge({}, NullProgress, job.progress); const batch = options.batch ?? 100; @@ -287,11 +322,8 @@ class MessagesPullHandler { // update progress before signal fails await job.updateProgress(progress); } - - if (job.parentKey) { - return await this.summaryProgress(job); - } - return await this.summary(options, job); + job.progress = progress; + return job; } } @@ -300,6 +332,15 @@ export class TaskActivity { private l: Locale, private conversation: Conversation, ) {} + public async queue(paused: boolean) { + let msg: string; + if (paused) { + msg = this.l.r('cli.cmd.queue.paused'); + } else { + msg = this.l.r('cli.cmd.queue.resumed'); + } + await this.conversation.activity(msg); + } public async details(data) { await this.conversation.activity( diff --git a/src/apps/chatwoot/di/DIContainer.ts b/src/apps/chatwoot/di/DIContainer.ts index 01b8b135..066cba13 100644 --- a/src/apps/chatwoot/di/DIContainer.ts +++ b/src/apps/chatwoot/di/DIContainer.ts @@ -212,6 +212,7 @@ export class DIContainer { linkPreview: LinkPreview.OFF, commands: { server: true, + queue: false, }, conversations: { sort: ConversationSort.created_newest, diff --git a/src/apps/chatwoot/dto/config.dto.ts b/src/apps/chatwoot/dto/config.dto.ts index 153e2c12..c576909f 100644 --- a/src/apps/chatwoot/dto/config.dto.ts +++ b/src/apps/chatwoot/dto/config.dto.ts @@ -20,6 +20,10 @@ export const DEFAULT_LOCALE = 'en-US'; export class ChatWootCommandsConfig { @IsBoolean() server: boolean = true; + + @IsOptional() + @IsBoolean() + queue?: boolean = false; } export enum LinkPreview { diff --git a/src/apps/chatwoot/i18n/locales/en-US.yaml b/src/apps/chatwoot/i18n/locales/en-US.yaml index cb619cc7..10ca37e2 100644 --- a/src/apps/chatwoot/i18n/locales/en-US.yaml +++ b/src/apps/chatwoot/i18n/locales/en-US.yaml @@ -353,8 +353,11 @@ cli.cmd.root.description: |- cli.cmd.root.sub.session: |- 🖥️ **Session**: +cli.cmd.root.sub.queue: |- + 🧩 **Queue**: + cli.cmd.root.sub.server: |- - 🖥️ **Session**: + 🔧 **Server**: cli.cmd.root.sub.sync: |- 🔄️ **Sync**: @@ -419,11 +422,50 @@ cli.cmd.messages.pull.option.broadcast: |- cli.cmd.messages.pull.option.media: |- Fetch media attachments (images, videos, files) when pulling history. +cli.cmd.messages.pull.option.pause: |- + Enqueue the messages pull job in a paused state, halting new WhatsApp message processing until the pull completes. + +cli.cmd.messages.pull.error.pause-no-queue: |- + ⛔ Cannot use `--pause` because the **queue** command is disabled in this app configuration. + +cli.cmd.queue.description: |- + Manage queues. + +cli.cmd.queue.summary: |- + Perform queue actions. + +cli.cmd.queue.action.description: |- + Choose an action. + +cli.cmd.queue.status.description: |- + Show the status of a queue. + +cli.cmd.queue.argument.name: |- + Queue name: `inbox` or `whatsapp` (optional). + +cli.cmd.queue.start.description: |- + Start a queue by name. + +cli.cmd.queue.stop.description: |- + Stop a queue by name. + +cli.cmd.queue.paused: |- + ⚪🧩 Queues has been paused. Send `queue status` to check the status. + +cli.cmd.queue.resumed: |- + 🟢🧩 Queues has been resumed. Send `queue status` to check the status. + +cli.cmd.queue.status.result: |- + 📊 Queue Status: + {{#queues}} + - **{{name}}**: {{#locked}}🔒{{/locked}}{{#paused}}⚪ Paused{{/paused}}{{^paused}}🟢 Running{{/paused}} + {{/queues}} + cli.cmd.messages.pull.option.timeout-media: |- - Set how long to try downloading media for each message (defaults to `30s`). + Set how long to try downloading media for each message. cli.cmd.messages.pull.argument.end: |- - Set how far back the pull window ends, e.g., `2d` stops at two days ago (defaults to `1d`). + Set how far back the pull window ends, e.g., `2d` stops at two days ago. cli.cmd.messages.pull.argument.start: |- Set where the pull window starts, e.g., `1d` begins one day ago (defaults to `0d`, meaning now). diff --git a/src/apps/chatwoot/services/ChatWootQueueService.ts b/src/apps/chatwoot/services/ChatWootQueueService.ts index 0175b6dc..7fa3d62e 100644 --- a/src/apps/chatwoot/services/ChatWootQueueService.ts +++ b/src/apps/chatwoot/services/ChatWootQueueService.ts @@ -1,10 +1,10 @@ -import { InjectQueue } from '@nestjs/bullmq'; import { Injectable } from '@nestjs/common'; import { EventName } from '@waha/apps/chatwoot/client/types'; import { InboxData } from '@waha/apps/chatwoot/consumers/types'; import { Queue } from 'bullmq'; import { QueueName } from '../consumers/QueueName'; +import { QueueRegistry } from './QueueRegistry'; /** * Service for managing ChatWoot queues for inbox events @@ -12,20 +12,7 @@ import { QueueName } from '../consumers/QueueName'; */ @Injectable() export class ChatWootQueueService { - constructor( - @InjectQueue(QueueName.INBOX_MESSAGE_CREATED) - private readonly messageCreatedQueue: Queue, - @InjectQueue(QueueName.INBOX_MESSAGE_UPDATED) - private readonly messageUpdatedQueue: Queue, - @InjectQueue(QueueName.INBOX_MESSAGE_DELETED) - private readonly messageDeletedQueue: Queue, - @InjectQueue(QueueName.INBOX_CONVERSATION_CREATED) - private readonly conversationCreatedQueue: Queue, - @InjectQueue(QueueName.INBOX_CONVERSATION_STATUS_CHANGED) - private readonly conversationStatusChanged: Queue, - @InjectQueue(QueueName.INBOX_COMMANDS) - private readonly commandsQueue: Queue, - ) {} + constructor(private readonly queueRegistry: QueueRegistry) {} /** * Generic method to add a job to a queue @@ -46,17 +33,19 @@ export class ChatWootQueueService { private getQueueForEvent(event: string): Queue | null { switch (event) { case EventName.CONVERSATION_CREATED: - return this.conversationCreatedQueue; + return this.queueRegistry.queue(QueueName.INBOX_CONVERSATION_CREATED); case EventName.MESSAGE_CREATED: - return this.messageCreatedQueue; + return this.queueRegistry.queue(QueueName.INBOX_MESSAGE_CREATED); case EventName.MESSAGE_UPDATED: - return this.messageUpdatedQueue; + return this.queueRegistry.queue(QueueName.INBOX_MESSAGE_UPDATED); case EventName.CONVERSATION_STATUS_CHANGED: - return this.conversationStatusChanged; + return this.queueRegistry.queue( + QueueName.INBOX_CONVERSATION_STATUS_CHANGED, + ); case 'message_deleted': - return this.messageDeletedQueue; + return this.queueRegistry.queue(QueueName.INBOX_MESSAGE_DELETED); case 'commands': - return this.commandsQueue; + return this.queueRegistry.queue(QueueName.INBOX_COMMANDS); default: return null; } @@ -67,7 +56,7 @@ export class ChatWootQueueService { */ async addMessageCreatedJob(data: InboxData): Promise { return await this.add( - this.messageCreatedQueue, + this.queueRegistry.queue(QueueName.INBOX_MESSAGE_CREATED), EventName.MESSAGE_CREATED, data, ); @@ -78,7 +67,7 @@ export class ChatWootQueueService { */ async addMessageUpdatedJob(data: InboxData): Promise { return await this.add( - this.messageUpdatedQueue, + this.queueRegistry.queue(QueueName.INBOX_MESSAGE_UPDATED), EventName.MESSAGE_UPDATED, data, ); @@ -88,14 +77,22 @@ export class ChatWootQueueService { * Add a job to the message deleted queue */ async addMessageDeletedJob(data: InboxData): Promise { - return await this.add(this.messageDeletedQueue, 'message_deleted', data); + return await this.add( + this.queueRegistry.queue(QueueName.INBOX_MESSAGE_DELETED), + 'message_deleted', + data, + ); } /** * Add a job to the commands queue */ async addCommandsJob(event: string, data: InboxData): Promise { - return await this.add(this.commandsQueue, event, data); + return await this.add( + this.queueRegistry.queue(QueueName.INBOX_COMMANDS), + event, + data, + ); } /** diff --git a/src/apps/chatwoot/services/ChatWootScheduleService.ts b/src/apps/chatwoot/services/ChatWootScheduleService.ts index 138dbe91..ecfe24d4 100644 --- a/src/apps/chatwoot/services/ChatWootScheduleService.ts +++ b/src/apps/chatwoot/services/ChatWootScheduleService.ts @@ -1,11 +1,10 @@ -import { InjectQueue } from '@nestjs/bullmq'; -import { Injectable, Logger } from '@nestjs/common'; -import { Queue } from 'bullmq'; +import { Injectable } from '@nestjs/common'; import { QueueName } from '../consumers/QueueName'; import { InjectPinoLogger, PinoLogger } from 'nestjs-pino'; import { ContactsPullRemove } from '@waha/apps/chatwoot/cli/cmd.contacts'; import { MessagesPullRemove } from '@waha/apps/chatwoot/cli/cmd.messages'; +import { QueueRegistry } from './QueueRegistry'; /** * Service for scheduling ChatWoot tasks @@ -16,14 +15,7 @@ export class ChatWootScheduleService { constructor( @InjectPinoLogger('DashboardConfigService') protected logger: PinoLogger, - @InjectQueue(QueueName.SCHEDULED_MESSAGE_CLEANUP) - private readonly messageCleanupQueue: Queue, - @InjectQueue(QueueName.SCHEDULED_CHECK_VERSION) - private readonly checkVersionQueue: Queue, - @InjectQueue(QueueName.TASK_CONTACTS_PULL) - private readonly contactsPullQueue: Queue, - @InjectQueue(QueueName.TASK_MESSAGES_PULL) - private readonly messagesPullQueue: Queue, + private readonly queueRegistry: QueueRegistry, ) {} /** @@ -33,7 +25,10 @@ export class ChatWootScheduleService { */ async schedule(appId: string, sessionName: string): Promise { // Message Cleanup - await this.messageCleanupQueue.upsertJobScheduler( + const messageCleanupQueue = this.queueRegistry.queue( + QueueName.SCHEDULED_MESSAGE_CLEANUP, + ); + await messageCleanupQueue.upsertJobScheduler( this.JobId(QueueName.SCHEDULED_MESSAGE_CLEANUP, appId), // Every day at 17:00 { pattern: '0 0 17 * * *' }, @@ -45,7 +40,10 @@ export class ChatWootScheduleService { }, ); // Check the version - await this.checkVersionQueue.upsertJobScheduler( + const checkVersionQueue = this.queueRegistry.queue( + QueueName.SCHEDULED_CHECK_VERSION, + ); + await checkVersionQueue.upsertJobScheduler( this.JobId(QueueName.SCHEDULED_CHECK_VERSION, appId), // Every Wednesday (3) at 18:00 { pattern: '0 0 18 * * 3' }, @@ -60,16 +58,25 @@ export class ChatWootScheduleService { async unschedule(appId: string, sessionName: string): Promise { // Message Cleanup - await this.messageCleanupQueue.removeJobScheduler( + const messageCleanupQueue = this.queueRegistry.queue( + QueueName.SCHEDULED_MESSAGE_CLEANUP, + ); + await messageCleanupQueue.removeJobScheduler( this.JobId(QueueName.SCHEDULED_MESSAGE_CLEANUP, appId), ); // Check the version - await this.checkVersionQueue.removeJobScheduler( + const checkVersionQueue = this.queueRegistry.queue( + QueueName.SCHEDULED_CHECK_VERSION, + ); + await checkVersionQueue.removeJobScheduler( this.JobId(QueueName.SCHEDULED_CHECK_VERSION, appId), ); // contacts - ContactsPullRemove(this.contactsPullQueue, appId, this.logger).catch( + const contactsPullQueue = this.queueRegistry.queue( + QueueName.TASK_CONTACTS_PULL, + ); + ContactsPullRemove(contactsPullQueue, appId, this.logger).catch( (reason) => { // Ignore errors this.logger.warn( @@ -78,7 +85,10 @@ export class ChatWootScheduleService { }, ); - MessagesPullRemove(this.messagesPullQueue, appId, this.logger).catch( + const messagesPullQueue = this.queueRegistry.queue( + QueueName.TASK_MESSAGES_PULL, + ); + MessagesPullRemove(messagesPullQueue, appId, this.logger).catch( (reason) => { this.logger.warn( `Failed to remove "messages" job for app ${appId}, session ${sessionName}: ${reason}`, diff --git a/src/apps/chatwoot/services/ChatWootWAHAQueueService.ts b/src/apps/chatwoot/services/ChatWootWAHAQueueService.ts index a1a8be9f..5790371d 100644 --- a/src/apps/chatwoot/services/ChatWootWAHAQueueService.ts +++ b/src/apps/chatwoot/services/ChatWootWAHAQueueService.ts @@ -1,4 +1,3 @@ -import { InjectQueue } from '@nestjs/bullmq'; import { Injectable } from '@nestjs/common'; import { ListenEventsForChatWoot } from '@waha/apps/chatwoot/consumers/waha/base'; import { populateSessionInfo } from '@waha/core/abc/manager.abc'; @@ -7,6 +6,7 @@ import { WAHAEvents } from '@waha/structures/enums.dto'; import { Queue } from 'bullmq'; import { QueueName } from '../consumers/QueueName'; +import { QueueRegistry } from './QueueRegistry'; /** * Service for managing ChatWoot queues for WAHA events @@ -14,20 +14,7 @@ import { QueueName } from '../consumers/QueueName'; */ @Injectable() export class ChatWootWAHAQueueService { - constructor( - @InjectQueue(QueueName.WAHA_MESSAGE_ANY) - private readonly queueMessageAny: Queue, - @InjectQueue(QueueName.WAHA_MESSAGE_REACTION) - private readonly queueMessageReaction: Queue, - @InjectQueue(QueueName.WAHA_MESSAGE_EDITED) - private readonly queueMessageEdited: Queue, - @InjectQueue(QueueName.WAHA_MESSAGE_REVOKED) - private readonly queueMessageRevoked: Queue, - @InjectQueue(QueueName.WAHA_MESSAGE_ACK) - private readonly queueMessageAck: Queue, - @InjectQueue(QueueName.WAHA_SESSION_STATUS) - private readonly queueSessionStatus: Queue, - ) {} + constructor(private readonly queueRegistry: QueueRegistry) {} /** * Get the specific queue for an event @@ -37,17 +24,17 @@ export class ChatWootWAHAQueueService { private getQueueForEvent(event: WAHAEvents): Queue | null { switch (event) { case WAHAEvents.MESSAGE_ANY: - return this.queueMessageAny; + return this.queueRegistry.queue(QueueName.WAHA_MESSAGE_ANY); case WAHAEvents.MESSAGE_REACTION: - return this.queueMessageReaction; + return this.queueRegistry.queue(QueueName.WAHA_MESSAGE_REACTION); case WAHAEvents.MESSAGE_EDITED: - return this.queueMessageEdited; + return this.queueRegistry.queue(QueueName.WAHA_MESSAGE_EDITED); case WAHAEvents.MESSAGE_REVOKED: - return this.queueMessageRevoked; + return this.queueRegistry.queue(QueueName.WAHA_MESSAGE_REVOKED); case WAHAEvents.MESSAGE_ACK: - return this.queueMessageAck; + return this.queueRegistry.queue(QueueName.WAHA_MESSAGE_ACK); case WAHAEvents.SESSION_STATUS: - return this.queueSessionStatus; + return this.queueRegistry.queue(QueueName.WAHA_SESSION_STATUS); default: return null; } diff --git a/src/apps/chatwoot/services/QueueManager.ts b/src/apps/chatwoot/services/QueueManager.ts new file mode 100644 index 00000000..ab6b86cc --- /dev/null +++ b/src/apps/chatwoot/services/QueueManager.ts @@ -0,0 +1,98 @@ +import { QueueName } from '../consumers/QueueName'; +import { QueueRegistry } from '@waha/apps/chatwoot/services/QueueRegistry'; +import { Injectable } from '@nestjs/common'; + +const Managable = true; +const Locked = false; + +export interface QueueStatus { + name: string; + paused: boolean; + locked: boolean; +} + +@Injectable() +export class QueueManager { + private readonly queues: Record; + + constructor(private readonly registry: QueueRegistry) { + this.queues = { + [QueueName.SCHEDULED_MESSAGE_CLEANUP]: Locked, + [QueueName.SCHEDULED_CHECK_VERSION]: Locked, + [QueueName.TASK_CONTACTS_PULL]: Locked, + [QueueName.TASK_MESSAGES_PULL]: Locked, + [QueueName.WAHA_SESSION_STATUS]: Locked, + [QueueName.WAHA_MESSAGE_ANY]: Managable, + [QueueName.WAHA_MESSAGE_REACTION]: Managable, + [QueueName.WAHA_MESSAGE_EDITED]: Managable, + [QueueName.WAHA_MESSAGE_REVOKED]: Managable, + [QueueName.WAHA_MESSAGE_ACK]: Managable, + [QueueName.INBOX_MESSAGE_CREATED]: Managable, + [QueueName.INBOX_MESSAGE_UPDATED]: Managable, + [QueueName.INBOX_CONVERSATION_CREATED]: Managable, + [QueueName.INBOX_CONVERSATION_STATUS_CHANGED]: Managable, + [QueueName.INBOX_MESSAGE_DELETED]: Managable, + [QueueName.INBOX_COMMANDS]: Locked, + } satisfies Record; + } + + async pause(queues: QueueName[] = null) { + queues = queues || Object.values(QueueName); + queues = this.managable(queues); + for (const name of queues) { + const queue = this.registry.queue(name); + await queue.pause(); + } + } + + async resume(queues: QueueName[] = null) { + queues = queues || Object.values(QueueName); + queues = this.managable(queues); + for (const name of queues) { + const queue = this.registry.queue(name); + await queue.resume(); + } + } + + resolve(shortcut: string | null): QueueName[] { + const queues = Object.values(QueueName); + switch (shortcut) { + case 'inbox': + return queues.filter((q) => q.startsWith('chatwoot.inbox')); + case 'whatsapp': + case 'waha': + return queues.filter((q) => q.startsWith('chatwoot.waha')); + case 'scheduled': + return queues.filter((q) => q.startsWith('chatwoot.scheduled')); + case 'task': + return queues.filter((q) => q.startsWith('chatwoot.task')); + case 'all': + case '': + case null: + case undefined: + return queues; + + default: + return [shortcut as QueueName]; + } + } + + protected managable(queues) { + return queues.filter((q) => this.queues[q] === Managable); + } + + async status(queues: QueueName[] = null): Promise { + queues = queues || Object.values(QueueName); + const result: QueueStatus[] = []; + for (const name of queues) { + const queue = this.registry.queue(name); + const paused = await queue.isPaused(); + result.push({ + name: name, + paused: paused, + locked: this.queues[name] === Locked, + }); + } + return result; + } +} diff --git a/src/apps/chatwoot/services/QueueRegistry.ts b/src/apps/chatwoot/services/QueueRegistry.ts new file mode 100644 index 00000000..ba1bd72e --- /dev/null +++ b/src/apps/chatwoot/services/QueueRegistry.ts @@ -0,0 +1,78 @@ +import { InjectQueue } from '@nestjs/bullmq'; +import { Injectable } from '@nestjs/common'; +import { Queue } from 'bullmq'; + +import { QueueName } from '../consumers/QueueName'; + +/** + * Central registry for ChatWoot queues backed by direct @InjectQueue bindings. + */ +@Injectable() +export class QueueRegistry { + private readonly queues: Record; + + constructor( + @InjectQueue(QueueName.SCHEDULED_MESSAGE_CLEANUP) + private readonly scheduledMessageCleanupQueue: Queue, + @InjectQueue(QueueName.SCHEDULED_CHECK_VERSION) + private readonly scheduledCheckVersionQueue: Queue, + @InjectQueue(QueueName.TASK_CONTACTS_PULL) + private readonly taskContactsPullQueue: Queue, + @InjectQueue(QueueName.TASK_MESSAGES_PULL) + private readonly taskMessagesPullQueue: Queue, + @InjectQueue(QueueName.WAHA_SESSION_STATUS) + private readonly wahaSessionStatusQueue: Queue, + @InjectQueue(QueueName.WAHA_MESSAGE_ANY) + private readonly wahaMessageAnyQueue: Queue, + @InjectQueue(QueueName.WAHA_MESSAGE_REACTION) + private readonly wahaMessageReactionQueue: Queue, + @InjectQueue(QueueName.WAHA_MESSAGE_EDITED) + private readonly wahaMessageEditedQueue: Queue, + @InjectQueue(QueueName.WAHA_MESSAGE_REVOKED) + private readonly wahaMessageRevokedQueue: Queue, + @InjectQueue(QueueName.WAHA_MESSAGE_ACK) + private readonly wahaMessageAckQueue: Queue, + @InjectQueue(QueueName.INBOX_MESSAGE_CREATED) + private readonly inboxMessageCreatedQueue: Queue, + @InjectQueue(QueueName.INBOX_MESSAGE_UPDATED) + private readonly inboxMessageUpdatedQueue: Queue, + @InjectQueue(QueueName.INBOX_CONVERSATION_CREATED) + private readonly inboxConversationCreatedQueue: Queue, + @InjectQueue(QueueName.INBOX_CONVERSATION_STATUS_CHANGED) + private readonly inboxConversationStatusChangedQueue: Queue, + @InjectQueue(QueueName.INBOX_MESSAGE_DELETED) + private readonly inboxMessageDeletedQueue: Queue, + @InjectQueue(QueueName.INBOX_COMMANDS) + private readonly inboxCommandsQueue: Queue, + ) { + // Strictly typed object literal + this.queues = { + [QueueName.SCHEDULED_MESSAGE_CLEANUP]: this.scheduledMessageCleanupQueue, + [QueueName.SCHEDULED_CHECK_VERSION]: this.scheduledCheckVersionQueue, + [QueueName.TASK_CONTACTS_PULL]: this.taskContactsPullQueue, + [QueueName.TASK_MESSAGES_PULL]: this.taskMessagesPullQueue, + [QueueName.WAHA_SESSION_STATUS]: this.wahaSessionStatusQueue, + [QueueName.WAHA_MESSAGE_ANY]: this.wahaMessageAnyQueue, + [QueueName.WAHA_MESSAGE_REACTION]: this.wahaMessageReactionQueue, + [QueueName.WAHA_MESSAGE_EDITED]: this.wahaMessageEditedQueue, + [QueueName.WAHA_MESSAGE_REVOKED]: this.wahaMessageRevokedQueue, + [QueueName.WAHA_MESSAGE_ACK]: this.wahaMessageAckQueue, + [QueueName.INBOX_MESSAGE_CREATED]: this.inboxMessageCreatedQueue, + [QueueName.INBOX_MESSAGE_UPDATED]: this.inboxMessageUpdatedQueue, + [QueueName.INBOX_CONVERSATION_CREATED]: + this.inboxConversationCreatedQueue, + [QueueName.INBOX_CONVERSATION_STATUS_CHANGED]: + this.inboxConversationStatusChangedQueue, + [QueueName.INBOX_MESSAGE_DELETED]: this.inboxMessageDeletedQueue, + [QueueName.INBOX_COMMANDS]: this.inboxCommandsQueue, + } satisfies Record; + } + + queue(name: QueueName): Queue { + const queue = this.queues[name]; + if (!queue) { + throw new Error(`Queue ${name} is not registered`); + } + return queue; + } +}