[core] ChatWoot - queue management

This commit is contained in:
devlikepro committed 2025-11-04 14:00:28 +07:00
1 parent 482dbbeb0e
commit 304b0037ec
16 files changed
+483 -87

No files matched your search

+4
View File
@@ -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 = {
+52
View File
@@ -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);
}
+3 -1
View File
@@ -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;
}
+14 -1
View File
@@ -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 <number>',
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,
},
+66
View File
@@ -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();
});
}
+2
View File
@@ -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;
};
@@ -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,
@@ -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<Progress> {
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<Progress> {
async handle(options: MessagesPullOptions, job: Job): Promise<Job> {
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(
+1
View File
@@ -212,6 +212,7 @@ export class DIContainer {
linkPreview: LinkPreview.OFF,
commands: {
server: true,
queue: false,
},
conversations: {
sort: ConversationSort.created_newest,
+4
View File
@@ -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 {
+45 -3
View File
@@ -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).
@@ -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<any> {
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<any> {
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<any> {
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<any> {
return await this.add(this.commandsQueue, event, data);
return await this.add(
this.queueRegistry.queue(QueueName.INBOX_COMMANDS),
event,
data,
);
}
/**
@@ -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<void> {
// 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<void> {
// 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}`,
@@ -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;
}
@@ -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<QueueName, boolean>;
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<QueueName, boolean>;
}
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<QueueStatus[]> {
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;
}
}
@@ -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<QueueName, Queue>;
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<QueueName, Queue>;
}
queue(name: QueueName): Queue {
const queue = this.queues[name];
if (!queue) {
throw new Error(`Queue ${name} is not registered`);
}
return queue;
}
}