diff --git a/package.json b/package.json index 9af1ea55..9bfc483d 100644 --- a/package.json +++ b/package.json @@ -101,8 +101,8 @@ "rimraf": "^4.3.0", "rxjs": "^7.8.1", "sharp": "^0.33.4", + "shell-quote": "^1.8.3", "sqlite3": "^5.1.7", - "string-argv": "^0.3.2", "swagger-ui-express": "^4.1.4", "ulid": "^2.3.0", "undici": "^7.16.0", diff --git a/src/apps/app_sdk/JobUtils.ts b/src/apps/app_sdk/JobUtils.ts index 2c94ef8e..48ef207a 100644 --- a/src/apps/app_sdk/JobUtils.ts +++ b/src/apps/app_sdk/JobUtils.ts @@ -1,5 +1,5 @@ -import { Job } from 'bullmq'; -import { Backoffs } from 'bullmq'; +import { Backoffs, Job } from 'bullmq'; +import { FlowJob } from 'bullmq/dist/esm/interfaces'; /** * Calculates the delay until the next attempt for a job. @@ -65,3 +65,30 @@ export function JobLink(job: Job): { text: string; url: string } { const url = `${base}/jobs/queue/${queue}/${id}`; return { text: text, url: url }; } + +/** + * Chain jobs so that they run one at a time + * A | B | C + * to + * A -> B -> C + */ +export function ChainJobsOneAtATime(jobs: FlowJob[]): FlowJob { + if (jobs.length == 0) { + throw new Error('No jobs provided'); + } + + for (const job of jobs) { + if (job.children.length != 0) { + throw new Error('Jobs with children are not supported'); + } + } + const left = [...jobs]; + const root = left.pop(); + let parent = root; + while (left.length != 0) { + const job = left.pop(); + parent.children = [job]; + parent = job; + } + return root; +} diff --git a/src/apps/app_sdk/waha/Paginator.ts b/src/apps/app_sdk/waha/Paginator.ts new file mode 100644 index 00000000..c1d36942 --- /dev/null +++ b/src/apps/app_sdk/waha/Paginator.ts @@ -0,0 +1,43 @@ +import * as lodash from 'lodash'; + +type Call = (processed: number) => Promise; + +export interface PaginatorParams { + processed?: number; + max?: number; +} + +const DefaultParams: PaginatorParams = { + processed: 0, + max: Infinity, +}; + +/** + * Generic paginator for APIs that return arrays of records with limits and offsets. + */ +export class ArrayPaginator { + private readonly params: PaginatorParams; + + constructor(params: PaginatorParams = DefaultParams) { + this.params = lodash.merge({}, DefaultParams, params) as PaginatorParams; + } + + async *iterate(call: Call): AsyncGenerator { + let processed = this.params.processed; + let records: Record[] = []; + + while (true) { + records = await call(processed); + if (records.length === 0) { + return; + } + for (const record of records) { + yield record; + processed += 1; + if (processed >= this.params.max) { + return; + } + } + } + } +} diff --git a/src/apps/chatwoot/session/WAHASelf.ts b/src/apps/app_sdk/waha/WAHASelf.ts similarity index 85% rename from src/apps/chatwoot/session/WAHASelf.ts rename to src/apps/app_sdk/waha/WAHASelf.ts index d79f5b7f..45550c5e 100644 --- a/src/apps/chatwoot/session/WAHASelf.ts +++ b/src/apps/app_sdk/waha/WAHASelf.ts @@ -1,5 +1,8 @@ import { Channel } from '@waha/structures/channels.dto'; -import { ChatPictureResponse } from '@waha/structures/chats.dto'; +import { + GetChatMessagesFilter, + GetChatMessagesQuery, +} from '@waha/structures/chats.dto'; import { ChatRequest, MessageFileRequest, @@ -79,6 +82,20 @@ export class WAHASelf { .then((response) => response.data); } + async getChats( + session: string, + page: PaginationParams, + opts?: RequestOptions, + ) { + const url = `/api/${session}/chats`; + const params = { + ...page, + }; + return await this.client + .get(url, { params: params, signal: opts?.signal }) + .then((response) => response.data); + } + async getContacts( session: string, page: PaginationParams, @@ -232,6 +249,39 @@ export class WAHASelf { .then((response) => response.data); } + async getMessages( + session: string, + chatId: string, + query: GetChatMessagesQuery, + filter: GetChatMessagesFilter, + opts?: RequestOptions, + ) { + const url = `/api/${session}/chats/${chatId}/messages`; + const params = { + ...query, + ...filter, + }; + return await this.client + .get(url, { params: params, signal: opts?.signal }) + .then((response) => response.data); + } + + async getMessageById( + session: string, + chatId: string, + messageId: string, + media: boolean, + opts?: RequestOptions, + ) { + const url = `/api/${session}/chats/${chatId}/messages/${messageId}`; + const params = { + downloadMedia: media, + }; + return await this.client + .get(url, { params: params, signal: opts?.signal }) + .then((response) => response.data); + } + async findPNByLid( session: string, lid: string, @@ -298,6 +348,10 @@ export class WAHASessionAPI { private api: WAHASelf, ) {} + getChats(page: PaginationParams, opts?: RequestOptions) { + return this.api.getChats(this.session, page, opts); + } + getContacts(page: PaginationParams, opts?: RequestOptions): Promise { return this.api.getContacts(this.session, page, opts); } @@ -371,6 +425,30 @@ export class WAHASessionAPI { return this.api.readMessages(this.session, chatId, opts); } + async getMessages( + chatId: string, + query: GetChatMessagesQuery, + filter: GetChatMessagesFilter, + opts?: RequestOptions, + ) { + return this.api.getMessages(this.session, chatId, query, filter, opts); + } + + async getMessageById( + chatId: string, + messageId: string, + media: boolean, + opts?: RequestOptions, + ) { + return this.api.getMessageById( + this.session, + chatId, + messageId, + media, + opts, + ); + } + // // Lids // diff --git a/src/apps/chatwoot/SessionStatusEmoji.ts b/src/apps/chatwoot/SessionStatusEmoji.ts deleted file mode 100644 index 14cda02d..00000000 --- a/src/apps/chatwoot/SessionStatusEmoji.ts +++ /dev/null @@ -1,18 +0,0 @@ -import { WAHASessionStatus } from '@waha/structures/enums.dto'; - -export function SessionStatusEmoji(status: WAHASessionStatus): string { - switch (status) { - case WAHASessionStatus.STOPPED: - return '⚠️'; - case WAHASessionStatus.STARTING: - return '⏳'; - case WAHASessionStatus.SCAN_QR_CODE: - return '⚠️'; - case WAHASessionStatus.WORKING: - return '🟢'; - case WAHASessionStatus.FAILED: - return '🛑'; - default: - return '❓'; - } -} diff --git a/src/apps/chatwoot/api/chatwoot.webhook.controller.ts b/src/apps/chatwoot/api/chatwoot.webhook.controller.ts index bfc6ad02..9d608f7b 100644 --- a/src/apps/chatwoot/api/chatwoot.webhook.controller.ts +++ b/src/apps/chatwoot/api/chatwoot.webhook.controller.ts @@ -6,13 +6,13 @@ import { Post, } from '@nestjs/common'; import { ApiOperation, ApiTags } from '@nestjs/swagger'; -import { FindChatID } from '@waha/apps/chatwoot/client/ids'; +import { IsCommandsChat } from '@waha/apps/chatwoot/client/ids'; import { EventName, MessageType } from '@waha/apps/chatwoot/client/types'; -import { INBOX_CONTACT_CHAT_ID } from '@waha/apps/chatwoot/const'; import { InboxData } from '@waha/apps/chatwoot/consumers/types'; import { ChatWootQueueService } from '@waha/apps/chatwoot/services/ChatWootQueueService'; import { SessionManager } from '@waha/core/abc/manager.abc'; import { AppRepository } from '@waha/apps/app_sdk/storage/AppRepository'; +import { CommandPrefix } from '@waha/apps/chatwoot/cli'; @Controller('webhooks/chatwoot/') @ApiTags('🧩 Apps') @@ -36,16 +36,26 @@ export class ChatwootWebhookController { return { success: true }; } - // Ignore all private notes - if (body.private) { - return { success: true }; - } - // Ignore all incoming messages if (body.message_type == MessageType.INCOMING) { return { success: true }; } + const isCommandsChat = IsCommandsChat(body); + // Ignore private notes (most of them) + const deleted = body?.content_attributes?.deleted; + if (body.private) { + // Ignore any private note in commands chats + if (isCommandsChat) { + return { success: true }; + } + // keep "deleted" notes + // So Agent can delete "Sent From API/WhatsApp" messages in ChatWoot + if (!deleted) { + return { success: true }; + } + } + const data: InboxData = { session: session, app: id, @@ -61,12 +71,8 @@ export class ChatwootWebhookController { } // Check if it's a command message (sent to the special inbox contact) - const sender = body?.conversation?.meta?.sender; - const chatId = FindChatID(sender); - const isCommandsChat = chatId === INBOX_CONTACT_CHAT_ID; - // Check if it's a deleted message - if (body.content_attributes?.deleted && !isCommandsChat) { + if (deleted && !isCommandsChat) { await this.chatWootQueueService.addMessageDeletedJob(data); return { success: true }; } @@ -74,10 +80,10 @@ export class ChatwootWebhookController { // Route to specific queues based on an event type switch (body.event) { case EventName.MESSAGE_CREATED: - if (!isCommandsChat) { - await this.chatWootQueueService.addMessageCreatedJob(data); - } else { + if (isCommandsChat || body.content?.startsWith(CommandPrefix)) { await this.chatWootQueueService.addCommandsJob(body.event, data); + } else { + await this.chatWootQueueService.addMessageCreatedJob(data); } return { success: true }; case EventName.MESSAGE_UPDATED: @@ -94,10 +100,10 @@ export class ChatwootWebhookController { return { success: true }; } - if (!isCommandsChat) { - await this.chatWootQueueService.addMessageUpdatedJob(data); - } else { + if (isCommandsChat || body.content?.startsWith(CommandPrefix)) { await this.chatWootQueueService.addCommandsJob(body.event, data); + } else { + await this.chatWootQueueService.addMessageUpdatedJob(data); } return { success: true }; default: diff --git a/src/apps/chatwoot/chatwoot.module.ts b/src/apps/chatwoot/chatwoot.module.ts index 00ad4ce8..87a8d188 100644 --- a/src/apps/chatwoot/chatwoot.module.ts +++ b/src/apps/chatwoot/chatwoot.module.ts @@ -15,7 +15,7 @@ import { ChatWootInboxCommandsConsumer } from './consumers/inbox/commands'; import { ChatWootInboxMessageCreatedConsumer } from './consumers/inbox/message_created'; import { ChatWootInboxMessageDeletedConsumer } from './consumers/inbox/message_deleted'; import { ChatWootInboxMessageUpdatedConsumer } from './consumers/inbox/message_updated'; -import { QueueName } from './consumers/QueueName'; +import { FlowProducerName, QueueName } from './consumers/QueueName'; import { CheckVersionConsumer } from './consumers/scheduled/check.version'; import { WAHAMessageAnyConsumer } from './consumers/waha/message.any'; import { WAHAMessageEditedConsumer } from './consumers/waha/message.edited'; @@ -29,10 +29,15 @@ import { ChatWootWAHAQueueService } from './services/ChatWootWAHAQueueService'; 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'; const CONTROLLERS = [ChatwootWebhookController, ChatwootLocalesController]; const IMPORTS = lodash.flatten([ + BullModule.registerFlowProducer({ + name: FlowProducerName.MESSAGES_PULL_FLOW, + }), RegisterAppQueue({ name: QueueName.SCHEDULED_MESSAGE_CLEANUP, defaultJobOptions: merge(ExponentialRetriesJobOptions, JobRemoveOptions), @@ -45,6 +50,10 @@ const IMPORTS = lodash.flatten([ name: QueueName.TASK_CONTACTS_PULL, defaultJobOptions: merge(ExponentialRetriesJobOptions, JobRemoveOptions), }), + RegisterAppQueue({ + name: QueueName.TASK_MESSAGES_PULL, + defaultJobOptions: merge(ExponentialRetriesJobOptions, JobRemoveOptions), + }), RegisterAppQueue({ name: QueueName.WAHA_MESSAGE_ANY, defaultJobOptions: merge(ExponentialRetriesJobOptions, JobRemoveOptions), @@ -105,6 +114,7 @@ const PROVIDERS = [ ChatWootInboxCommandsConsumer, // Tasks TaskContactsPullConsumer, + TaskMessagesPullConsumer, // WAHA WAHASessionStatusConsumer, WAHAMessageAnyConsumer, diff --git a/src/apps/chatwoot/cli/cmd.contacts.ts b/src/apps/chatwoot/cli/cmd.contacts.ts index 5dd392f3..0876935b 100644 --- a/src/apps/chatwoot/cli/cmd.contacts.ts +++ b/src/apps/chatwoot/cli/cmd.contacts.ts @@ -8,14 +8,19 @@ import { Job, JobsOptions, Queue } from 'bullmq'; import { ChatWootScheduleService } from '@waha/apps/chatwoot/services/ChatWootScheduleService'; import { ILogger } from '@waha/apps/app_sdk/ILogger'; import { JobDataTimeout } from '@waha/apps/app_sdk/AppConsumer'; +import { + ExponentialRetriesJobOptions, + JobRemoveOptions, + merge, +} from '@waha/apps/app_sdk/constants'; export async function ContactsPullStart( ctx: CommandContext, options: ContactsPullOptions, jobOptions: JobsOptions & JobDataTimeout, ) { - const jobId = ChatWootScheduleService.SingleJobId(ctx.app); - const job: Job = await ctx.queues.importContacts.getJob(jobId); + const jobId = ChatWootScheduleService.SingleJobId(ctx.data.app); + const job: Job = await ctx.queues.contactsPull.getJob(jobId); if (job) { const state = await job.getState(); const done = state === 'completed' || state === 'failed'; @@ -26,32 +31,28 @@ export async function ContactsPullStart( } // Remove already done job - await ContactsPullRemove(ctx.queues.importContacts, ctx.app, ctx.logger); + await ContactsPullRemove(ctx.queues.contactsPull, ctx.data.app, ctx.logger); } - const opts: JobsOptions = { - ...jobOptions, - jobId: jobId, - }; + const opts: JobsOptions = merge( + ExponentialRetriesJobOptions, + JobRemoveOptions, + jobOptions, + ); const data = { - app: ctx.app, - session: ctx.session, + ...ctx.data, timeout: { job: jobOptions.timeout.job, }, options: options, }; - await ctx.queues.importContacts.add(QueueName.TASK_CONTACTS_PULL, data, opts); - - const msg = ctx.l.r('cli.cmd.contacts.pull.queued', { - batch: options.batch, - }); - await ctx.conversation.incoming(msg); + opts.jobId = jobId; + await ctx.queues.contactsPull.add(QueueName.TASK_CONTACTS_PULL, data, opts); } export async function ContactsPullStatus(ctx: CommandContext) { - const jobId = ChatWootScheduleService.SingleJobId(ctx.app); - const job: Job | null = await ctx.queues.importContacts.getJob(jobId); + const jobId = ChatWootScheduleService.SingleJobId(ctx.data.app); + const job: Job | null = await ctx.queues.contactsPull.getJob(jobId); if (!job) { const msg = ctx.l.r('cli.cmd.contacts.status.not-found'); diff --git a/src/apps/chatwoot/cli/cmd.messages.ts b/src/apps/chatwoot/cli/cmd.messages.ts new file mode 100644 index 00000000..db49ff21 --- /dev/null +++ b/src/apps/chatwoot/cli/cmd.messages.ts @@ -0,0 +1,182 @@ +import * as lodash from 'lodash'; +import { CommandContext } from '@waha/apps/chatwoot/cli/types'; +import { QueueName } from '@waha/apps/chatwoot/consumers/QueueName'; +import { + ChatID, + MessagesPullOptions, + MessagesPullStatusMessage, + TaskActivity, +} from '@waha/apps/chatwoot/consumers/task/messages.pull'; +import { ChatWootScheduleService } from '@waha/apps/chatwoot/services/ChatWootScheduleService'; +import { ILogger } from '@waha/apps/app_sdk/ILogger'; +import { JobDataTimeout } from '@waha/apps/app_sdk/AppConsumer'; +import { FlowProducer, Job, JobsOptions, Queue } from 'bullmq'; +import { FlowJob } from 'bullmq/dist/esm/interfaces'; +import { GetAllChatIDs, IsCommandsChat } from '@waha/apps/chatwoot/client/ids'; +import { ChainJobsOneAtATime } from '@waha/apps/app_sdk/JobUtils'; +import { + ExponentialRetriesJobOptions, + JobRemoveOptions, + merge, +} from '@waha/apps/app_sdk/constants'; +import { EngineHelper } from '@waha/apps/chatwoot/waha'; +import { WAHASelf, WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; +import { + IgnoreJidConfig, + isNullJid, + isPnUser, + JidFilter, +} from '@waha/core/utils/jids'; + +function oneJob(chat: string, data, opts: JobsOptions): FlowJob { + data = lodash.cloneDeep(data); + data.options.chat = chat; + opts = lodash.cloneDeep(opts); + return { + name: chat, + queueName: QueueName.TASK_MESSAGES_PULL, + data: data, + opts: merge(ExponentialRetriesJobOptions, JobRemoveOptions, opts), + children: [], + }; +} + +export async function MessagesPullStart( + ctx: CommandContext, + options: MessagesPullOptions, + jobOptions: JobsOptions & JobDataTimeout, +) { + const queue = ctx.queues.messagesPull; + const jobId = ChatWootScheduleService.SingleJobId(ctx.data.app); + const job: Job | null = await queue.getJob(jobId); + + if (job) { + const state = await job.getState(); + const done = state === 'completed' || state === 'failed'; + if (!done) { + const msg = ctx.l.r('cli.cmd.messages.pull.already-running', { + chat: job.data.options.chat, + }); + await ctx.conversation.incoming(msg); + return; + } + + await MessagesPullRemove(queue, ctx.data.app, ctx.logger); + } + const opts: JobsOptions = { + ...jobOptions, + ignoreDependencyOnFailure: true, + delay: 1_000, + }; + const data = { + ...ctx.data, + timeout: { + job: jobOptions.timeout.job, + }, + options: options, + }; + const producer: FlowProducer = ctx.flows.messagesPull; + let chats: string[] = [options.chat]; + + if (options.chat === ChatID.ALL) { + if (IsCommandsChat(data.body)) { + if (EngineHelper.SupportsAllChatForMessage()) { + // keep "all" chats - that's fine + } else { + // "all" chats for engine that doesn't support it + chats = await resolveAllToChats(ctx.waha, ctx.data.session, options); + if (chats.length == 0) { + throw new Error(ctx.l.r('cli.cmd.messages.pull.no-chats-found')); + } + } + } else { + // It's per-chat command, get the chat ids and run the pulling process + chats = GetAllChatIDs(data.body?.conversation?.meta?.sender); + chats = EngineHelper.FilterChatIdsForMessages(chats); + } + } + const children = chats.map((chat) => oneJob(chat, data, opts)); + + let root: FlowJob; + if (children.length == 1) { + // Schedule one job, no "summary" job required + root = children[0]; + } else { + // Add "summary" parent job so we get the overall progress at the end + root = oneJob(ChatID.SUMMARY, data, opts); + root.children = [ChainJobsOneAtATime(children)]; + } + root.opts.jobId = jobId; + await producer.add(root); + const activity = new TaskActivity(ctx.l, ctx.conversation); + await activity.details(data); +} + +export async function MessagesPullStatus(ctx: CommandContext) { + const queue = ctx.queues.messagesPull; + const jobId = ChatWootScheduleService.SingleJobId(ctx.data.app); + const job: Job | null = await queue.getJob(jobId); + + if (!job) { + const msg = ctx.l.r('cli.cmd.messages.status.not-found'); + await ctx.conversation.incoming(msg); + return; + } + + const state = await job.getState(); + const msg = MessagesPullStatusMessage(ctx.l, job, state); + await ctx.conversation.incoming(msg); +} + +export async function MessagesPullRemove( + queue: Queue, + app: string, + logger: ILogger, +): Promise { + const jobId = ChatWootScheduleService.SingleJobId(app); + const job: Job | null = await queue.getJob(jobId); + if (!job) { + logger.info('Pull Messages job has already been removed'); + return false; + } + + await job.remove(); + return true; +} + +/** + * WORKING ONLY FOR WEBJS + */ +async function resolveAllToChats( + waha: WAHASelf, + sessionName: string, + options: MessagesPullOptions, +): Promise { + const session = new WAHASessionAPI(sessionName, waha); + let chats = await session.getChats({ + limit: undefined, + offset: undefined, + }); + // new chat last + chats = lodash.sortBy(chats, (c) => c.timestamp || 0); + + const gte = Date.now() - options.period.start; + const result = []; + const jids = new JidFilter(options.ignore); + for (const chat of chats) { + const id = chat.id._serialized; + if (!jids.include(id)) { + continue; + } + if (isNullJid(id)) { + continue; + } + const timestamp = (chat.timestamp || Infinity) * 1000; + if (timestamp < gte) { + // Too old chat + continue; + } + result.push(id); + } + return result; +} diff --git a/src/apps/chatwoot/cli/cms.session.ts b/src/apps/chatwoot/cli/cms.session.ts index 7a3d7f68..e064db42 100644 --- a/src/apps/chatwoot/cli/cms.session.ts +++ b/src/apps/chatwoot/cli/cms.session.ts @@ -1,11 +1,11 @@ import { CommandContext } from '@waha/apps/chatwoot/cli/types'; import { TKey } from '@waha/apps/chatwoot/i18n/templates'; -import { SessionStatusEmoji } from '@waha/apps/chatwoot/SessionStatusEmoji'; +import { SessionStatusEmoji } from '@waha/apps/chatwoot/emoji'; import { AttachmentFromBuffer } from '@waha/apps/chatwoot/client/messages'; import { MessageType } from '@waha/apps/chatwoot/client/types'; export async function SessionRestart(ctx: CommandContext) { - await ctx.waha.restart(ctx.session); + await ctx.waha.restart(ctx.data.session); } export async function SessionStart(ctx: CommandContext) { @@ -13,20 +13,20 @@ export async function SessionStart(ctx: CommandContext) { } export async function SessionLogout(ctx: CommandContext) { - ctx.logger.info(`Logging out session ${ctx.session}`); - await ctx.waha.logout(ctx.session); + ctx.logger.info(`Logging out session ${ctx.data.session}`); + await ctx.waha.logout(ctx.data.session); const text = ctx.l.key(TKey.APP_LOGOUT_SUCCESS).render(); await ctx.conversation.incoming(text); } export async function SessionStop(ctx: CommandContext) { - ctx.logger.info(`Stopping session ${ctx.session}`); - await ctx.waha.stop(ctx.session); + ctx.logger.info(`Stopping session ${ctx.data.session}`); + await ctx.waha.stop(ctx.data.session); } export async function SessionStatus(ctx: CommandContext) { - ctx.logger.info(`Getting status for session ${ctx.session}`); - const session = await ctx.waha.get(ctx.session); + ctx.logger.info(`Getting status for session ${ctx.data.session}`); + const session = await ctx.waha.get(ctx.data.session); const emoji = SessionStatusEmoji(session.status); const text = ctx.l.key(TKey.APP_SESSION_CURRENT_STATUS).render({ emoji: emoji, @@ -39,15 +39,15 @@ export async function SessionStatus(ctx: CommandContext) { } export async function SessionQR(ctx: CommandContext) { - const content = await ctx.waha.qr(ctx.session); + const content = await ctx.waha.qr(ctx.data.session); const message = AttachmentFromBuffer(content, 'qr.jpg'); message.message_type = MessageType.INCOMING; await ctx.conversation.send(message); } export async function SessionScreenshot(ctx: CommandContext) { - ctx.logger.info(`Getting screenshot for session ${ctx.session}`); - const content = await ctx.waha.screenshot(ctx.session); + ctx.logger.info(`Getting screenshot for session ${ctx.data.session}`); + const content = await ctx.waha.screenshot(ctx.data.session); const message = AttachmentFromBuffer(content, 'screenshot.jpg'); message.message_type = MessageType.INCOMING; await ctx.conversation.send(message); diff --git a/src/apps/chatwoot/cli/index.ts b/src/apps/chatwoot/cli/index.ts index 654d765e..705e009c 100644 --- a/src/apps/chatwoot/cli/index.ts +++ b/src/apps/chatwoot/cli/index.ts @@ -1,253 +1,9 @@ -import { - Argument, - Command, - CommanderError, - Option, - OutputConfiguration, -} from 'commander'; -import argvSplit from 'string-argv'; +import { CommanderError } from 'commander'; +import { parse } from 'shell-quote'; import { BufferedOutput } from '@waha/apps/chatwoot/cli/utils/BufferedOutput'; -import { - buildFormatHelp, - fullCommandPath, -} from '@waha/apps/chatwoot/cli/utils/help'; -import { - JobAttemptsOption, - JobTimeoutOption, - ParseMS, -} from '@waha/apps/chatwoot/cli/utils/options'; import { CommandContext } from '@waha/apps/chatwoot/cli/types'; -import { - SessionLogout, - SessionQR, - SessionRestart, - SessionScreenshot, - SessionStart, - SessionStatus, - SessionStop, -} from '@waha/apps/chatwoot/cli/cms.session'; -import { ServerReboot, ServerStatus } from '@waha/apps/chatwoot/cli/cmd.server'; import { ChatWootCommandsConfig } from '@waha/apps/chatwoot/dto/config.dto'; -import { CommandDisabled } from '@waha/apps/chatwoot/cli/cmd.disabled'; -import { - ContactsPullStart, - ContactsPullStatus, -} from '@waha/apps/chatwoot/cli/cmd.contacts'; -import { ContactsPullOptions } from '@waha/apps/chatwoot/consumers/task/contacts.pull'; -import { JobsOptions } from 'bullmq'; -import { JobDataTimeout } from '@waha/apps/app_sdk/AppConsumer'; - -function BuildProgram( - commands: ChatWootCommandsConfig, - ctx: CommandContext, - output: OutputConfiguration, -) { - const l = ctx.l; - const SessionGroup = l.r('cli.cmd.root.sub.session'); - const ServerGroup = l.r('cli.cmd.root.sub.server'); - const SyncGroup = l.r('cli.cmd.root.sub.sync'); - - const program = new Command(); - program - .name('') - .description(l.r('cli.cmd.root.description')) - .exitOverride() - .configureOutput(output) - .showSuggestionAfterError(true) - .commandsGroup(SessionGroup) - .commandsGroup(ServerGroup) - .commandsGroup(SyncGroup) - .configureHelp({ - helpWidth: 200, - styleOptionTerm(str: string) { - return `- \`${str}\``; - }, - styleArgumentTerm(str: string) { - return `- **${str}**`; - }, - styleUsage(str) { - return `**${str}**`; - }, - styleSubcommandTerm(str: string) { - return `- **${str}**`; - }, - styleSubcommandDescription(str: string) { - return str ? `- ${str}` : str; - }, - subcommandTerm: fullCommandPath, - commandUsage(cmd) { - // For the root program - no usage - if (!cmd.parent) - return `${l.r('cli.usage.command')} ${l.r('cli.usage.options')}`; - // For subcommands, keep default behavior - if (cmd.commands.length > 0) { - return `${fullCommandPath(cmd)} ${l.r('cli.usage.command')} ${l.r( - 'cli.usage.options', - )}`; - } - return `${fullCommandPath(cmd)} ${l.r('cli.usage.options')}`; - }, - formatHelp: buildFormatHelp(l), - }); - - // - // Session - // - program - .command('status') - .alias('1') - .description(l.r('cli.cmd.session.status.description')) - .helpGroup(SessionGroup) - .action(() => SessionStatus(ctx)); - program - .command('restart') - .alias('2') - .description(l.r('cli.cmd.session.restart.description')) - .helpGroup(SessionGroup) - .action(() => SessionRestart(ctx)); - program - .command('start') - .alias('3') - .description(l.r('cli.cmd.session.start.description')) - .helpGroup(SessionGroup) - .action(() => SessionStart(ctx)); - program - .command('stop') - .alias('4') - .description(l.r('cli.cmd.session.stop.description')) - .helpGroup(SessionGroup) - .action(() => SessionStop(ctx)); - program - .command('logout') - .alias('5') - .description(l.r('cli.cmd.session.logout.description')) - .helpGroup(SessionGroup) - .action(() => SessionLogout(ctx)); - program - .command('qr') - .alias('6') - .description(l.r('cli.cmd.session.qr.description')) - .helpGroup(SessionGroup) - .action(() => SessionQR(ctx)); - program - .command('screenshot') - .alias('7') - .description(l.r('cli.cmd.session.screenshot.description')) - .helpGroup(SessionGroup) - .action(() => SessionScreenshot(ctx)); - - // - // Pull Contacts - // - program - .command('contacts') - .alias('/contacts') - .summary(l.r('cli.cmd.contacts.summary')) - .description(l.r('cli.cmd.contacts.description')) - .helpGroup(SyncGroup) - .addArgument( - new Argument( - '[action]', - l.r('cli.cmd.contacts.action.description'), - ).choices(['pull', 'status']), - ) - .addOption( - new Option('--avatar ', l.r('cli.cmd.contacts.pull.option.avatar')) - .choices(['skip', 'if-missing', 'update']) - .default('if-missing'), - ) - .option('--groups', l.r('cli.cmd.contacts.pull.option.groups')) - .option('--no-lids', l.r('cli.cmd.contacts.pull.option.no-lids')) - .option( - '--no-attributes', - l.r('cli.cmd.contacts.pull.option.no-attributes'), - ) - .option( - '--batch ', - l.r('cli.cmd.contacts.pull.option.batch'), - 100 as any, - ) - .addOption(new JobAttemptsOption(l, 6)) - .addOption(new JobTimeoutOption(l, '10m')) - .addOption( - new Option( - '--delay-contact ', - l.r('cli.cmd.contacts.pull.option.delay-contact'), - ) - .argParser(ParseMS) - .default(ParseMS('0.1s'), '0.1s'), - ) - .addOption( - new Option( - '--delay-batch ', - l.r('cli.cmd.contacts.pull.option.delay-batch'), - ) - .argParser(ParseMS) - .default(ParseMS('1s'), '1s'), - ) - .action(async (action, opts, cmd: Command) => { - if (!action) { - cmd.outputHelp(); - return; - } - - if (action === 'status') { - await ContactsPullStatus(ctx); - return; - } - - const options: ContactsPullOptions = { - batch: opts.batch, - avatar: opts.avatar, - attributes: opts.attributes, - contacts: { - lids: opts.lids, - groups: opts.groups, - }, - delay: { - contact: opts.delayContact, - batch: opts.delayBatch, - }, - }; - const jobOptions: JobsOptions & JobDataTimeout = { - attempts: opts.attempts, - timeout: { - job: opts.timeout, - }, - }; - await ContactsPullStart(ctx, options, jobOptions); - }); - - // - // Server - // - const server = program - .command('server', { hidden: !commands.server }) - .description(l.r('cli.cmd.server.description')) - .helpGroup(ServerGroup); - if (!commands.server) { - // Do not show help, show disabled ASAP - server.action(() => CommandDisabled(ctx, 'server')); - } - server - .command('status') - .description(l.r('cli.cmd.server.status.description')) - .action(() => - commands.server - ? ServerStatus(ctx) - : CommandDisabled(ctx, 'server status'), - ); - server - .command('reboot') - .description(l.r('cli.cmd.server.reboot.description')) - .option('-f, --force', l.r('cli.cmd.server.reboot.option.force'), false) - .action((opts) => - commands.server - ? ServerReboot(ctx, opts.force) - : CommandDisabled(ctx, 'server reboot'), - ); - return program; -} +import { BuildProgram } from '@waha/apps/chatwoot/cli/program.a'; export async function runText( commands: ChatWootCommandsConfig, @@ -262,7 +18,7 @@ export async function runText( .trim(); // final cleanup const output = new BufferedOutput(); const program = BuildProgram(commands, ctx, output); - const argv = argvSplit(text, '', ''); + const argv = parse(text); try { await program.parseAsync(argv, { from: 'user' }); } catch (err) { @@ -274,3 +30,5 @@ export async function runText( } return output; } + +export const CommandPrefix = process.env.WAHA_CHATWOOT_COMMAND_PREFIX || 'wa/'; diff --git a/src/apps/chatwoot/cli/program.a.ts b/src/apps/chatwoot/cli/program.a.ts new file mode 100644 index 00000000..f8200751 --- /dev/null +++ b/src/apps/chatwoot/cli/program.a.ts @@ -0,0 +1,68 @@ +import { ChatWootCommandsConfig } from '@waha/apps/chatwoot/dto/config.dto'; +import { CommandContext } from '@waha/apps/chatwoot/cli/types'; +import { Argument, Command, Option, OutputConfiguration } from 'commander'; +import { + buildFormatHelp, + fullCommandPath, +} from '@waha/apps/chatwoot/cli/utils/help'; +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'; + +function Program(ctx: CommandContext, output: OutputConfiguration) { + const l = ctx.l; + const program = new Command(); + program + .name('') + .description(l.r('cli.cmd.root.description')) + .exitOverride() + .configureOutput(output) + .showSuggestionAfterError(true) + .configureHelp({ + helpWidth: 200, + styleOptionTerm(str: string) { + return `- \`${str}\``; + }, + styleArgumentTerm(str: string) { + return `- **${str}**`; + }, + styleUsage(str) { + return `**${str}**`; + }, + styleSubcommandTerm(str: string) { + return `- **${str}**`; + }, + styleSubcommandDescription(str: string) { + return str ? `- ${str}` : str; + }, + subcommandTerm: fullCommandPath, + commandUsage(cmd) { + // For the root program - no usage + if (!cmd.parent) + return `${l.r('cli.usage.command')} ${l.r('cli.usage.options')}`; + // For subcommands, keep default behavior + if (cmd.commands.length > 0) { + return `${fullCommandPath(cmd)} ${l.r('cli.usage.command')} ${l.r( + 'cli.usage.options', + )}`; + } + return `${fullCommandPath(cmd)} ${l.r('cli.usage.options')}`; + }, + formatHelp: buildFormatHelp(l), + }); + return program; +} + +export function BuildProgram( + commands: ChatWootCommandsConfig, + ctx: CommandContext, + output: OutputConfiguration, +) { + const program = Program(ctx, output); + AddSessionCommand(program, ctx); + AddContactsCommand(program, ctx); + AddMessagesCommand(program, ctx); + AddServerCommand(program, ctx, commands.server); + return program; +} diff --git a/src/apps/chatwoot/cli/program.contacts.ts b/src/apps/chatwoot/cli/program.contacts.ts new file mode 100644 index 00000000..1e24a3d8 --- /dev/null +++ b/src/apps/chatwoot/cli/program.contacts.ts @@ -0,0 +1,113 @@ +import { Argument, Command, Option } from 'commander'; +import { CommandContext } from '@waha/apps/chatwoot/cli/types'; +import { + JobAttemptsOption, + JobTimeoutOption, + ParseMS, + NotNegativeNumber, + ProgressOption, +} from '@waha/apps/chatwoot/cli/utils/options'; +import { + ContactsPullStart, + ContactsPullStatus, +} from '@waha/apps/chatwoot/cli/cmd.contacts'; +import { ContactsPullOptions } from '@waha/apps/chatwoot/consumers/task/contacts.pull'; +import { JobsOptions } from 'bullmq'; +import { JobDataTimeout } from '@waha/apps/app_sdk/AppConsumer'; + +export function AddContactsCommand(program: Command, ctx: CommandContext) { + const l = ctx.l; + const SyncGroup = l.r('cli.cmd.root.sub.sync'); + program.commandsGroup(SyncGroup); + + program + .command('contacts') + .summary(l.r('cli.cmd.contacts.summary')) + .description(l.r('cli.cmd.contacts.description')) + .helpGroup(SyncGroup) + .addArgument( + new Argument( + '[action]', + l.r('cli.cmd.contacts.action.description'), + ).choices(['pull', 'status', 'help']), + ) + .addOption( + new Option( + '-a, --avatar ', + l.r('cli.cmd.contacts.pull.option.avatar'), + ) + .choices(['if-missing', 'update']) + .preset('if-missing'), + ) + .option('-g, --groups', l.r('cli.cmd.contacts.pull.option.groups')) + .option('-l, --lids', l.r('cli.cmd.contacts.pull.option.lids')) + .option( + '--na, --no-attributes', + l.r('cli.cmd.contacts.pull.option.no-attributes'), + ) + .option( + '-b, --batch ', + l.r('cli.cmd.contacts.pull.option.batch'), + NotNegativeNumber, + 100, + ) + .addOption( + ProgressOption(l.r('cli.cmd.contacts.pull.option.progress'), 100), + ) + .addOption(new JobAttemptsOption(l, 6)) + .addOption(new JobTimeoutOption(l, '10m')) + .addOption( + new Option( + '--dc, --delay-contact ', + l.r('cli.cmd.contacts.pull.option.delay-contact'), + ) + .argParser(ParseMS) + .default(ParseMS('0.1s'), '0.1s'), + ) + .addOption( + new Option( + '--db, --delay-batch ', + l.r('cli.cmd.contacts.pull.option.delay-batch'), + ) + .argParser(ParseMS) + .default(ParseMS('1s'), '1s'), + ) + .action(async (action, opts, cmd: Command) => { + if (!action) { + cmd.outputHelp(); + return; + } + + if (action === 'help') { + cmd.outputHelp(); + return; + } + + if (action === 'status') { + await ContactsPullStatus(ctx); + return; + } + + const options: ContactsPullOptions = { + batch: opts.batch, + progress: opts.progress, + avatar: opts.avatar, + attributes: opts.attributes, + contacts: { + lids: opts.lids, + groups: opts.groups, + }, + delay: { + contact: opts.delayContact, + batch: opts.delayBatch, + }, + }; + const jobOptions: JobsOptions & JobDataTimeout = { + attempts: opts.attempts, + timeout: { + job: opts.timeout, + }, + }; + await ContactsPullStart(ctx, options, jobOptions); + }); +} diff --git a/src/apps/chatwoot/cli/program.messages.ts b/src/apps/chatwoot/cli/program.messages.ts new file mode 100644 index 00000000..198142ae --- /dev/null +++ b/src/apps/chatwoot/cli/program.messages.ts @@ -0,0 +1,124 @@ +import { Argument, Command, Option } from 'commander'; +import { CommandContext } from '@waha/apps/chatwoot/cli/types'; +import { + JobAttemptsOption, + JobTimeoutOption, + ParseMS, + NotNegativeNumber, + ProgressOption, +} from '@waha/apps/chatwoot/cli/utils/options'; +import { + MessagesPullStart, + MessagesPullStatus, +} from '@waha/apps/chatwoot/cli/cmd.messages'; +import { MessagesPullOptions } from '@waha/apps/chatwoot/consumers/task/messages.pull'; +import { JobsOptions } from 'bullmq'; +import { JobDataTimeout } from '@waha/apps/app_sdk/AppConsumer'; + +export function AddMessagesCommand(program: Command, ctx: CommandContext) { + const l = ctx.l; + const SyncGroup = l.r('cli.cmd.root.sub.sync'); + program.commandsGroup(SyncGroup); + + program + .command('messages') + .summary(l.r('cli.cmd.messages.summary')) + .description(l.r('cli.cmd.messages.description')) + .helpGroup(SyncGroup) + .addArgument( + new Argument( + '[action]', + l.r('cli.cmd.messages.action.description'), + ).choices(['pull', 'status', 'help']), + ) + .addArgument( + new Argument('[end]', l.r('cli.cmd.messages.pull.argument.end')) + .argParser(ParseMS) + .default(ParseMS('1d'), '1d'), + ) + .addArgument( + new Argument('[start]', l.r('cli.cmd.messages.pull.argument.start')) + .argParser(ParseMS) + .default(ParseMS('0d'), '0d'), + ) + .option( + '-c, --chat ', + l.r('cli.cmd.messages.pull.option.chat'), + 'all', + ) + .option('-f, --force', l.r('cli.cmd.messages.pull.option.force')) + .option('--nd, --no-dm', l.r('cli.cmd.messages.pull.option.no-dm')) + .option('-g, --groups', l.r('cli.cmd.messages.pull.option.groups')) + .option('--ch, --channels', l.r('cli.cmd.messages.pull.option.channels')) + .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( + '-b, --batch ', + l.r('cli.cmd.messages.pull.option.batch'), + NotNegativeNumber, + 100 as any, + ) + .addOption(ProgressOption(l.r('cli.cmd.messages.pull.option.progress'))) + .addOption(new JobAttemptsOption(l, 6)) + .addOption(new JobTimeoutOption(l, '10m')) + .addOption( + new Option( + '--tm, --timeout-media ', + l.r('cli.cmd.messages.pull.option.timeout-media'), + ) + .argParser(ParseMS) + .default(ParseMS('30s'), '30s'), + ) + .action(async (action, end, start, opts, cmd: Command) => { + if (!action) { + cmd.outputHelp(); + return; + } + + if (action === 'help') { + cmd.outputHelp(); + return; + } + + if (action === 'status') { + await MessagesPullStatus(ctx); + return; + } + if (end > start) { + // Swap + const tmp = end; + end = start; + start = tmp; + } + + const options: MessagesPullOptions = { + chat: opts.chat, + progress: opts.progress, + period: { + end: end, + start: start, + }, + media: opts.media, + force: opts.force, + timeout: { + media: opts.timeoutMedia, + }, + batch: opts.batch, + ignore: { + dm: !opts.dm, + status: !opts.status, + groups: !opts.groups, + channels: !opts.channels, + broadcast: !opts.broadcast, + }, + }; + const jobOptions: JobsOptions & JobDataTimeout = { + attempts: opts.attempts, + timeout: { + job: opts.timeout, + }, + }; + await MessagesPullStart(ctx, options, jobOptions); + }); +} diff --git a/src/apps/chatwoot/cli/program.server.ts b/src/apps/chatwoot/cli/program.server.ts new file mode 100644 index 00000000..c32e2cf4 --- /dev/null +++ b/src/apps/chatwoot/cli/program.server.ts @@ -0,0 +1,38 @@ +import { Command } from 'commander'; +import { CommandContext } from '@waha/apps/chatwoot/cli/types'; +import { CommandDisabled } from '@waha/apps/chatwoot/cli/cmd.disabled'; +import { ServerReboot, ServerStatus } from '@waha/apps/chatwoot/cli/cmd.server'; + +export function AddServerCommand( + program: Command, + ctx: CommandContext, + enabled: boolean, +) { + const l = ctx.l; + const ServerGroup = l.r('cli.cmd.root.sub.server'); + program.commandsGroup(ServerGroup); + + const server = program + .command('server', { hidden: !enabled }) + .description(l.r('cli.cmd.server.description')) + .helpGroup(ServerGroup); + if (!enabled) { + // Do not show help, show disabled ASAP + server.action(() => CommandDisabled(ctx, 'server')); + } + server + .command('status') + .description(l.r('cli.cmd.server.status.description')) + .action(() => + enabled ? ServerStatus(ctx) : CommandDisabled(ctx, 'server status'), + ); + server + .command('reboot') + .description(l.r('cli.cmd.server.reboot.description')) + .option('-f, --force', l.r('cli.cmd.server.reboot.option.force'), false) + .action((opts) => + enabled + ? ServerReboot(ctx, opts.force) + : CommandDisabled(ctx, 'server reboot'), + ); +} diff --git a/src/apps/chatwoot/cli/program.session.ts b/src/apps/chatwoot/cli/program.session.ts new file mode 100644 index 00000000..35d80d1c --- /dev/null +++ b/src/apps/chatwoot/cli/program.session.ts @@ -0,0 +1,60 @@ +import { Command } from 'commander'; +import { CommandContext } from '@waha/apps/chatwoot/cli/types'; +import { + SessionLogout, + SessionQR, + SessionRestart, + SessionScreenshot, + SessionStart, + SessionStatus, + SessionStop, +} from '@waha/apps/chatwoot/cli/cms.session'; + +export function AddSessionCommand(program: Command, ctx: CommandContext) { + const l = ctx.l; + const SessionGroup = l.r('cli.cmd.root.sub.session'); + program.commandsGroup(SessionGroup); + + program + .command('status') + .alias('1') + .description(l.r('cli.cmd.session.status.description')) + .helpGroup(SessionGroup) + .action(() => SessionStatus(ctx)); + program + .command('restart') + .alias('2') + .description(l.r('cli.cmd.session.restart.description')) + .helpGroup(SessionGroup) + .action(() => SessionRestart(ctx)); + program + .command('start') + .alias('3') + .description(l.r('cli.cmd.session.start.description')) + .helpGroup(SessionGroup) + .action(() => SessionStart(ctx)); + program + .command('stop') + .alias('4') + .description(l.r('cli.cmd.session.stop.description')) + .helpGroup(SessionGroup) + .action(() => SessionStop(ctx)); + program + .command('logout') + .alias('5') + .description(l.r('cli.cmd.session.logout.description')) + .helpGroup(SessionGroup) + .action(() => SessionLogout(ctx)); + program + .command('qr') + .alias('6') + .description(l.r('cli.cmd.session.qr.description')) + .helpGroup(SessionGroup) + .action(() => SessionQR(ctx)); + program + .command('screenshot') + .alias('7') + .description(l.r('cli.cmd.session.screenshot.description')) + .helpGroup(SessionGroup) + .action(() => SessionScreenshot(ctx)); +} diff --git a/src/apps/chatwoot/cli/types.ts b/src/apps/chatwoot/cli/types.ts index c02b24b6..e0c3d93c 100644 --- a/src/apps/chatwoot/cli/types.ts +++ b/src/apps/chatwoot/cli/types.ts @@ -1,17 +1,21 @@ -import { WAHASelf } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASelf } from '@waha/apps/app_sdk/waha/WAHASelf'; import { Locale } from '@waha/apps/chatwoot/i18n/locale'; import { Conversation } from '../client/Conversation'; import { ILogger } from '@waha/apps/app_sdk/ILogger'; -import { Queue } from 'bullmq'; +import { FlowProducer, Queue } from 'bullmq'; +import { InboxData } from '@waha/apps/chatwoot/consumers/types'; export interface CommandContext { - app: string; - session: string; + data: InboxData; logger: ILogger; l: Locale; waha: WAHASelf; conversation: Conversation; queues: { - importContacts: Queue; + contactsPull: Queue; + messagesPull: Queue; + }; + flows: { + messagesPull: FlowProducer; }; } diff --git a/src/apps/chatwoot/cli/utils/options.ts b/src/apps/chatwoot/cli/utils/options.ts index 7a898515..3c895f70 100644 --- a/src/apps/chatwoot/cli/utils/options.ts +++ b/src/apps/chatwoot/cli/utils/options.ts @@ -19,15 +19,47 @@ export function ParseSeconds(value: string): number { export class JobAttemptsOption extends Option { constructor(l: Locale, def: number) { - super('--attempts ', l.r('cli.cmd.options.job.attempts')); + super('--at, --attempts ', l.r('cli.cmd.options.job.attempts')); + this.argParser(NotNegativeNumber); this.default(def); } } export class JobTimeoutOption extends Option { constructor(l: Locale, def: string) { - super('--timeout ', l.r('cli.cmd.options.job.timeout')); + super('-t, --timeout ', l.r('cli.cmd.options.job.timeout')); this.argParser(ParseMS); this.default(ParseMS(def)); } } + +export function ProgressOption( + description: string, + def: number = 1000, +): Option { + return new Option('-p, --progress [number]', description) + .argParser(NotNegativeNumber) + .default(def); +} + +export function NotNegativeNumber(value: string): number { + const n = parseInt(value, 10); + if (isNaN(n)) { + throw new Error(`Invalid number: "${value}"`); + } + if (n < 0) { + throw new Error(`Number must be positive or 0: "${value}"`); + } + return n; +} + +export function PositiveNumber(value: string): number { + const n = parseInt(value, 10); + if (isNaN(n)) { + throw new Error(`Invalid number: "${value}"`); + } + if (n <= 0) { + throw new Error(`Number must be positive: "${value}"`); + } + return n; +} diff --git a/src/apps/chatwoot/client/ContactConversationService.ts b/src/apps/chatwoot/client/ContactConversationService.ts index 9cb40da3..d8cd5502 100644 --- a/src/apps/chatwoot/client/ContactConversationService.ts +++ b/src/apps/chatwoot/client/ContactConversationService.ts @@ -53,10 +53,11 @@ export class ContactConversationService { return this.cache.get(chatId); } - let cwContact = await this.contactService.findOrCreateContact(contactInfo); + let [cwContact, created] = + await this.contactService.findOrCreateContact(contactInfo); // Update custom attributes - always - this.logger.info( + this.logger.debug( `Updating if required contact custom attributes for chat.id: ${chatId}, contact.id: ${cwContact.data.id}`, ); const attributes = await contactInfo.Attributes(); @@ -69,7 +70,7 @@ export class ContactConversationService { contactInfo, AvatarUpdateMode.IF_MISSING, ); - this.logger.info( + this.logger.debug( `Using contact for chat.id: ${chatId}, contact.id: ${cwContact.data.id}, contact.sourceId: ${cwContact.sourceId}`, ); @@ -80,7 +81,7 @@ export class ContactConversationService { id: cwContact.data.id, sourceId: cwContact.sourceId, }); - this.logger.info( + this.logger.debug( `Using conversation for chat.id: ${chatId}, conversation.id: ${conversation.id}, contact.id: ${cwContact.sourceId}`, ); @@ -134,7 +135,7 @@ export class ContactConversationService { } public ResetCache(chatIds: Array) { - this.logger.info(`Resetting cache chat ids: ${chatIds.join(', ')}`); + this.logger.debug(`Resetting cache chat ids: ${chatIds.join(', ')}`); for (const chatId of chatIds) { this.cache.delete(chatId); } @@ -147,7 +148,7 @@ export class ContactConversationService { } const current = this.cache.get(chatId); if (current.id !== contactId) { - this.logger.info( + this.logger.debug( `Resetting cache for chat id: ${chatId}, value changed from ${current} to ${contactId}`, ); this.cache.delete(chatId); diff --git a/src/apps/chatwoot/client/ContactService.ts b/src/apps/chatwoot/client/ContactService.ts index 4a5ba0f7..9f21778a 100644 --- a/src/apps/chatwoot/client/ContactService.ts +++ b/src/apps/chatwoot/client/ContactService.ts @@ -34,14 +34,18 @@ export class ContactService { private logger: ILogger, ) {} - async findOrCreateContact(contactInfo: ContactInfo) { + async findOrCreateContact( + contactInfo: ContactInfo, + ): Promise<[ContactResponse, boolean]> { const chatId = contactInfo.ChatId(); let contact = await this.searchByAnyID(chatId); - if (!contact) { - const request = await contactInfo.PublicContactCreate(); - contact = await this.create(chatId, request); + if (contact) { + return [contact, false]; } - return contact; + + const request = await contactInfo.PublicContactCreate(); + contact = await this.create(chatId, request); + return [contact, true]; } async searchByAnyID(chatId: string): Promise { @@ -156,10 +160,10 @@ export class ContactService { contact: ContactResponse, contactInfo: ContactInfo, mode: AvatarUpdateMode, - ) { + ): Promise { // Update Avatar if nothing, but keep the original one if any if (contact.data.thumbnail && mode == AvatarUpdateMode.IF_MISSING) { - return; + return false; } const chatId = contactInfo.ChatId(); const avatarUrl = await contactInfo.AvatarUrl().catch((err) => { @@ -169,13 +173,15 @@ export class ContactService { this.logger.warn(err); return null; }); - if (avatarUrl) { - this.updateAvatarUrlSafe(contact.data.id, avatarUrl); + if (!avatarUrl) { + return false; } + const success = await this.updateAvatarUrlSafe(contact.data.id, avatarUrl); + return success; } - public updateAvatarUrlSafe(contactId, avatarUrl: string) { - this.accountAPI.contacts + public updateAvatarUrlSafe(contactId, avatarUrl: string): Promise { + return this.accountAPI.contacts .update({ accountId: this.config.accountId, id: contactId, @@ -183,11 +189,15 @@ export class ContactService { avatar_url: avatarUrl, }, }) + .then(() => { + return true; + }) .catch((e) => { this.logger.warn( `Error updating avatar_url for contact.id: ${contactId}`, ); this.logger.warn(e); + return true; }); } } diff --git a/src/apps/chatwoot/client/Conversation.ts b/src/apps/chatwoot/client/Conversation.ts index 246e4eee..1e19385c 100644 --- a/src/apps/chatwoot/client/Conversation.ts +++ b/src/apps/chatwoot/client/Conversation.ts @@ -1,9 +1,11 @@ import ChatwootClient from '@figuro/chatwoot-sdk'; import type { conversation_message_create } from '@figuro/chatwoot-sdk/dist/models/conversation_message_create'; import { MessageType } from '@waha/apps/chatwoot/client/types'; +import { IsCommandsChat } from '@waha/apps/chatwoot/client/ids'; export class Conversation { public onError: (e: any) => void; + public overrideIncoming: any = null; constructor( private accountAPI: ChatwootClient, @@ -13,6 +15,9 @@ export class Conversation { ) {} public async send(data: conversation_message_create) { + if (data.message_type === MessageType.INCOMING && this.overrideIncoming) { + data = { ...data, ...this.overrideIncoming }; + } try { const message = await this.accountAPI.messages.create({ accountId: this.accountId, @@ -29,10 +34,13 @@ export class Conversation { } public async incoming(text: string) { - const data: conversation_message_create = { + let data: conversation_message_create = { content: text, message_type: MessageType.INCOMING as any, }; + if (this.overrideIncoming) { + data = { ...data, ...this.overrideIncoming }; + } return this.send(data); } @@ -44,11 +52,22 @@ export class Conversation { return this.send(data); } - public async outgoing(text: string) { + public async note(text: string) { const data: conversation_message_create = { content: text, + private: true, message_type: MessageType.OUTGOING as any, }; return this.send(data); } + + /** + * Force a note message (private outgoing) for all communications + */ + public forceNote() { + this.overrideIncoming = { + private: true, + message_type: MessageType.OUTGOING as any, + }; + } } diff --git a/src/apps/chatwoot/client/ConversationService.ts b/src/apps/chatwoot/client/ConversationService.ts index 1fca62bc..0c3824ad 100644 --- a/src/apps/chatwoot/client/ConversationService.ts +++ b/src/apps/chatwoot/client/ConversationService.ts @@ -40,7 +40,7 @@ export class ConversationService { inboxIdentifier: this.config.inboxIdentifier, contactIdentifier: contact.sourceId, }); - this.logger.info( + this.logger.debug( `Created conversation.id: ${conversation.id} for contact.id: ${contact.id}, contact.sourceId: ${contact.sourceId}`, ); return conversation; @@ -51,7 +51,7 @@ export class ConversationService { if (!conversation) { conversation = await this.create(contact); } - this.logger.info( + this.logger.debug( `Using conversation.id: ${conversation.id} for contact.id: ${contact.id}, contact.sourceId: ${contact.sourceId}`, ); return conversation; diff --git a/src/apps/chatwoot/client/ids.ts b/src/apps/chatwoot/client/ids.ts index c17a8706..c136ec03 100644 --- a/src/apps/chatwoot/client/ids.ts +++ b/src/apps/chatwoot/client/ids.ts @@ -1,4 +1,5 @@ -import { AttributeKey } from '@waha/apps/chatwoot/const'; +import { AttributeKey, INBOX_CONTACT_CHAT_ID } from '@waha/apps/chatwoot/const'; +import * as lodash from 'lodash'; import { WhatsAppMessage } from '@waha/apps/chatwoot/storage'; import { buildMessageId } from '@waha/core/engines/noweb/session.noweb.core'; import { isLidUser } from '@waha/core/utils/jids'; @@ -22,7 +23,7 @@ export function GetAllChatIDs(contact: any): Array { attrs[AttributeKey.WA_JID], attrs[AttributeKey.WA_LID], ]; - return ids.filter(Boolean); + return lodash.uniq(ids.filter(Boolean)); } export function FindChatID(contact: any): string | null { @@ -55,3 +56,9 @@ export function ContactAttr(chatId: string): AttributeKey { return AttributeKey.WA_JID; } } + +export function IsCommandsChat(body): boolean { + const sender = body?.conversation?.meta?.sender; + const chatId = FindChatID(sender); + return chatId === INBOX_CONTACT_CHAT_ID; +} diff --git a/src/apps/chatwoot/consumers/QueueName.ts b/src/apps/chatwoot/consumers/QueueName.ts index bc76ce23..677372e4 100644 --- a/src/apps/chatwoot/consumers/QueueName.ts +++ b/src/apps/chatwoot/consumers/QueueName.ts @@ -9,6 +9,7 @@ export enum QueueName { // Task // TASK_CONTACTS_PULL = 'chatwoot.task | contacts.pull', + TASK_MESSAGES_PULL = 'chatwoot.task | messages.pull', // // WAHA Events @@ -32,3 +33,7 @@ export enum QueueName { INBOX_MESSAGE_DELETED = 'chatwoot.inbox | message_deleted', INBOX_COMMANDS = 'chatwoot.inbox | commands', } + +export enum FlowProducerName { + MESSAGES_PULL_FLOW = 'messages.pull.flow', +} diff --git a/src/apps/chatwoot/consumers/inbox/base.ts b/src/apps/chatwoot/consumers/inbox/base.ts index faa8345a..a67de4cc 100644 --- a/src/apps/chatwoot/consumers/inbox/base.ts +++ b/src/apps/chatwoot/consumers/inbox/base.ts @@ -10,7 +10,7 @@ import { ChatIDNotFoundForContactError, PhoneNumberNotFoundInWhatsAppError, } from '@waha/apps/chatwoot/errors'; -import { WAHASessionAPI } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; import { SessionManager } from '@waha/core/abc/manager.abc'; import { RMutexService } from '@waha/modules/rmutex/rmutex.service'; import { Job } from 'bullmq'; @@ -18,6 +18,7 @@ import { PinoLogger } from 'nestjs-pino'; import { AppRepository } from '../../storage'; import { TKey } from '@waha/apps/chatwoot/i18n/templates'; +import { Conversation } from '@waha/apps/chatwoot/client/Conversation'; /** * Base class for ChatWoot inbox consumers @@ -90,6 +91,12 @@ export abstract class ChatWootInboxMessageConsumer extends AppConsumer { } } + protected conversationForReport(container, body): Conversation { + return container + .ContactConversationService() + .ConversationById(body.conversation.id); + } + /** * Report an error for a message */ @@ -98,9 +105,7 @@ export abstract class ChatWootInboxMessageConsumer extends AppConsumer { const header: string = this.ErrorHeaderKey() ? container.Locale().key(this.ErrorHeaderKey()).render() : err.message || `${err}`; - const conversation = container - .ContactConversationService() - .ConversationById(body.conversation.id); + const conversation = this.conversationForReport(container, body); const reporter = container.ChatWootErrorReporter(job); await reporter.ReportError( conversation, @@ -118,9 +123,7 @@ export abstract class ChatWootInboxMessageConsumer extends AppConsumer { } const container = await this.DIContainer(job, job.data.app); - const conversation = container - .ContactConversationService() - .ConversationById(body.conversation.id); + const conversation = this.conversationForReport(container, body); const reporter = container.ChatWootErrorReporter(job); await reporter.ReportSucceeded(conversation, body.message_type, body.id); } diff --git a/src/apps/chatwoot/consumers/inbox/commands.ts b/src/apps/chatwoot/consumers/inbox/commands.ts index 350ea8d8..31c71c8d 100644 --- a/src/apps/chatwoot/consumers/inbox/commands.ts +++ b/src/apps/chatwoot/consumers/inbox/commands.ts @@ -1,15 +1,19 @@ -import { InjectQueue, Processor } from '@nestjs/bullmq'; +import { InjectFlowProducer, InjectQueue, Processor } from '@nestjs/bullmq'; import { JOB_CONCURRENCY } from '@waha/apps/app_sdk/constants'; import { ChatWootInboxMessageConsumer } from '@waha/apps/chatwoot/consumers/inbox/base'; -import { QueueName } from '@waha/apps/chatwoot/consumers/QueueName'; +import { + FlowProducerName, + QueueName, +} from '@waha/apps/chatwoot/consumers/QueueName'; import { DIContainer } from '@waha/apps/chatwoot/di/DIContainer'; import { SessionManager } from '@waha/core/abc/manager.abc'; import { RMutexService } from '@waha/modules/rmutex/rmutex.service'; -import { Job, Queue } from 'bullmq'; +import { FlowProducer, Job, Queue } from 'bullmq'; import { PinoLogger } from 'nestjs-pino'; import { TKey } from '@waha/apps/chatwoot/i18n/templates'; -import { runText } from '@waha/apps/chatwoot/cli'; +import { CommandPrefix, runText } from '@waha/apps/chatwoot/cli'; import { CommandContext } from '@waha/apps/chatwoot/cli/types'; +import { IsCommandsChat } from '@waha/apps/chatwoot/client/ids'; @Processor(QueueName.INBOX_COMMANDS, { concurrency: JOB_CONCURRENCY }) export class ChatWootInboxCommandsConsumer extends ChatWootInboxMessageConsumer { @@ -18,7 +22,11 @@ export class ChatWootInboxCommandsConsumer extends ChatWootInboxMessageConsumer log: PinoLogger, rmutex: RMutexService, @InjectQueue(QueueName.TASK_CONTACTS_PULL) - private readonly importContactsQueue: Queue, + private readonly contactsPullQueue: Queue, + @InjectQueue(QueueName.TASK_MESSAGES_PULL) + private readonly messagesPullQueue: Queue, + @InjectFlowProducer(FlowProducerName.MESSAGES_PULL_FLOW) + private readonly messagesPullFlow: FlowProducer, ) { super(manager, log, rmutex, 'ChatWootInboxCommandsConsumer'); } @@ -27,23 +35,37 @@ export class ChatWootInboxCommandsConsumer extends ChatWootInboxMessageConsumer return null; } + protected conversationForReport(container, body) { + const conversation = super.conversationForReport(container, body); + if (!IsCommandsChat(body)) { + conversation.forceNote(); + } + return conversation; + } + protected async Process(container: DIContainer, body, job: Job) { - const cmd = body.content; + let cmd = body.content; + cmd = cmd.startsWith(CommandPrefix) ? cmd.slice(CommandPrefix.length) : cmd; this.logger.info( `Executing command '${cmd}' for session ${job.data.session}...`, ); const repo = container.ContactConversationService(); - const conversation = await repo.InboxNotifications(); - + let conversation = repo.ConversationById(body.conversation.id); + if (!IsCommandsChat(body)) { + conversation.forceNote(); + } const ctx: CommandContext = { - app: job.data.app, - session: job.data.session, + data: job.data, logger: this.logger, l: container.Locale(), waha: container.WAHASelf(), conversation: conversation, queues: { - importContacts: this.importContactsQueue, + contactsPull: this.contactsPullQueue, + messagesPull: this.messagesPullQueue, + }, + flows: { + messagesPull: this.messagesPullFlow, }, }; const commands = container.ChatWootConfig().commands; diff --git a/src/apps/chatwoot/consumers/inbox/message_created.ts b/src/apps/chatwoot/consumers/inbox/message_created.ts index 21aa296b..54eb2591 100644 --- a/src/apps/chatwoot/consumers/inbox/message_created.ts +++ b/src/apps/chatwoot/consumers/inbox/message_created.ts @@ -8,8 +8,8 @@ import { } from '@waha/apps/chatwoot/consumers/inbox/base'; import { QueueName } from '@waha/apps/chatwoot/consumers/QueueName'; import { DIContainer } from '@waha/apps/chatwoot/di/DIContainer'; -import { EngineHelper } from '@waha/apps/chatwoot/session'; -import { WAHASessionAPI } from '@waha/apps/chatwoot/session/WAHASelf'; +import { EngineHelper } from '@waha/apps/chatwoot/waha'; +import { WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; import { ChatwootMessage, MessageMapping, diff --git a/src/apps/chatwoot/consumers/inbox/message_deleted.ts b/src/apps/chatwoot/consumers/inbox/message_deleted.ts index faa628f1..bd57299e 100644 --- a/src/apps/chatwoot/consumers/inbox/message_deleted.ts +++ b/src/apps/chatwoot/consumers/inbox/message_deleted.ts @@ -8,7 +8,7 @@ import { } from '@waha/apps/chatwoot/consumers/inbox/base'; import { QueueName } from '@waha/apps/chatwoot/consumers/QueueName'; import { DIContainer } from '@waha/apps/chatwoot/di/DIContainer'; -import { WAHASessionAPI } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; import { MessageMappingService } from '@waha/apps/chatwoot/storage'; import { SessionManager } from '@waha/core/abc/manager.abc'; import { RMutexService } from '@waha/modules/rmutex/rmutex.service'; @@ -56,6 +56,13 @@ class MessageDeletedHandler { message_id: body.id, }); if (!messages || messages.length == 0) { + if (body.private) { + // Private notes can have no whatsapp messages + this.logger.info( + `No WhatsApp message to delete for private note '${body.id}'`, + ); + return; + } throw Error( `No WhatsApp message found for Chatwoot message '${body.id}'`, ); diff --git a/src/apps/chatwoot/consumers/inbox/message_updated.ts b/src/apps/chatwoot/consumers/inbox/message_updated.ts index ed006e08..b24dea4e 100644 --- a/src/apps/chatwoot/consumers/inbox/message_updated.ts +++ b/src/apps/chatwoot/consumers/inbox/message_updated.ts @@ -9,7 +9,7 @@ import { RMutexService } from '@waha/modules/rmutex/rmutex.service'; import { Job } from 'bullmq'; import { PinoLogger } from 'nestjs-pino'; -import { WAHASessionAPI } from '../../session/WAHASelf'; +import { WAHASessionAPI } from '../../../app_sdk/waha/WAHASelf'; import { TKey } from '@waha/apps/chatwoot/i18n/templates'; @Processor(QueueName.INBOX_MESSAGE_UPDATED, { concurrency: JOB_CONCURRENCY }) diff --git a/src/apps/chatwoot/consumers/task/base.ts b/src/apps/chatwoot/consumers/task/base.ts index 50761165..1ac48b48 100644 --- a/src/apps/chatwoot/consumers/task/base.ts +++ b/src/apps/chatwoot/consumers/task/base.ts @@ -12,6 +12,7 @@ import { PinoLogger } from 'nestjs-pino'; import { AppRepository } from '../../storage'; import { TKey } from '@waha/apps/chatwoot/i18n/templates'; import { SignalRace } from '@waha/utils/abortable'; +import { IsCommandsChat } from '@waha/apps/chatwoot/client/ids'; /** * Base class for ChatWoot background task consumers @@ -74,6 +75,16 @@ export abstract class ChatWootTaskConsumer extends AppConsumer { } } + protected conversationForReport(container, body) { + const conversation = container + .ContactConversationService() + .ConversationById(body.conversation.id); + if (!IsCommandsChat(body)) { + conversation.forceNote(); + } + return conversation; + } + /** * Report an error for a scheduled job */ @@ -84,9 +95,7 @@ export abstract class ChatWootTaskConsumer extends AppConsumer { : err.message || `${err}`; header = `${job.queueName}: ${header}`; - const conversation = await container - .ContactConversationService() - .InboxNotifications(); + const conversation = this.conversationForReport(container, job.data.body); const reporter = container.ChatWootErrorReporter(job); await reporter.ReportError(conversation, header, MessageType.INCOMING, err); @@ -99,9 +108,7 @@ export abstract class ChatWootTaskConsumer extends AppConsumer { } const container = await this.DIContainer(job, job.data.app); - const conversation = await container - .ContactConversationService() - .InboxNotifications(); + const conversation = this.conversationForReport(container, job.data.body); const reporter = container.ChatWootErrorReporter(job); await reporter.ReportSucceeded(conversation, MessageType.INCOMING); diff --git a/src/apps/chatwoot/consumers/task/contacts.pull.ts b/src/apps/chatwoot/consumers/task/contacts.pull.ts index 2e6b1151..7ef6023b 100644 --- a/src/apps/chatwoot/consumers/task/contacts.pull.ts +++ b/src/apps/chatwoot/consumers/task/contacts.pull.ts @@ -10,7 +10,7 @@ import { PinoLogger } from 'nestjs-pino'; import { ChatWootTaskConsumer } from '@waha/apps/chatwoot/consumers/task/base'; import { TKey } from '@waha/apps/chatwoot/i18n/templates'; -import { WAHASessionAPI } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; import { WhatsAppContactInfo } from '@waha/apps/chatwoot/contacts/WhatsAppContactInfo'; import { ContactSortField } from '@waha/structures/contacts.dto'; import { SortOrder } from '@waha/structures/pagination.dto'; @@ -23,11 +23,17 @@ import { Locale } from '@waha/apps/chatwoot/i18n/locale'; import { ILogger } from '@waha/apps/app_sdk/ILogger'; import { JobLink } from '@waha/apps/app_sdk/JobUtils'; import { sleep } from '@waha/utils/promiseTimeout'; -import { isJidGroup, isPnUser, isLidUser } from '@waha/core/utils/jids'; +import { isJidGroup, isLidUser, isPnUser } from '@waha/core/utils/jids'; +import { + ArrayPaginator, + PaginatorParams, +} from '@waha/apps/app_sdk/waha/Paginator'; +import { EngineHelper } from '@waha/apps/chatwoot/waha'; export interface ContactsPullOptions { batch: number; - avatar: 'skip' | 'if-missing' | 'update'; + progress: number | null; + avatar: null | 'if-missing' | 'update'; attributes: boolean; contacts: { groups: boolean; @@ -66,10 +72,7 @@ export class TaskContactsPullConsumer extends ChatWootTaskConsumer { const locale = container.Locale(); const waha = container.WAHASelf(); const session = new WAHASessionAPI(job.data.session, waha); - const conversation = await container - .ContactConversationService() - .InboxNotifications(); - + const conversation = this.conversationForReport(container, job.data.body); const handler = new ContactsPullHandler( signal, container.Logger(), @@ -83,18 +86,34 @@ export class TaskContactsPullConsumer extends ChatWootTaskConsumer { } interface Progress { - ok: number; - errors: number; + created: number; + updated: number; skipped: number; + errors: number; + avatar: { + updated: number; + }; +} + +function total(progress: Progress) { + return ( + progress.created + progress.updated + progress.skipped + progress.errors + ); } const NullProgress: Progress = { - ok: 0, - errors: 0, + created: 0, + updated: 0, skipped: 0, + errors: 0, + avatar: { + updated: 0, + }, }; class ContactsPullHandler { + private activity: TaskActivity; + constructor( private signal: AbortSignal, private logger: ILogger, @@ -102,99 +121,110 @@ class ContactsPullHandler { private session: WAHASessionAPI, private l: Locale, private contactService: ContactService, - ) {} + ) { + this.activity = new TaskActivity(this.l, this.conversation); + } async handle(options: ContactsPullOptions, job) { const batch = options.batch; let progress = lodash.merge({}, NullProgress, job.progress); - this.logger.info( + this.logger.debug( `Pulling contacts for session ${job.data.session} with batch size ${batch}...`, ); - if (progress.ok == 0 && progress.errors == 0 && progress.skipped == 0) { - await this.conversation.activity( - this.l.r('task.contacts.started', { progress: progress }), - ); - } else { - await this.conversation.activity( - this.l.r('task.contacts.progress', { progress: progress }), - ); + if (lodash.isEqual(progress, NullProgress)) { + await this.activity.started(progress); } - // fake 'null' first, so it's not empty list - // we overwrite it at the beginning of the loop - let contacts: any[] = [null]; - while (contacts.length != 0) { - this.signal.throwIfAborted(); - const processed = progress.ok + progress.errors + progress.skipped; - this.logger.info( - `Fetching contacts batch: offset=${processed}, limit=${batch}...`, - ); - contacts = await this.session.getContacts( - { - offset: processed, - limit: batch, - sortBy: ContactSortField.ID, - sortOrder: SortOrder.ASC, - }, - { signal: this.signal }, - ); + const query = { + sortBy: ContactSortField.ID, + sortOrder: SortOrder.ASC, + limit: batch, + offset: 0, + }; + const params: PaginatorParams = { + processed: total(progress), + }; + const paginator = new ArrayPaginator(params); + const contacts = paginator.iterate((processed: number) => { + query.offset = processed; + return this.session.getContacts(query, { + signal: this.signal, + }); + }); - for (const contact of contacts) { - try { - const chatId = contact.id; - if (isJidGroup(chatId)) { - if (!options.contacts.groups) { - this.logger.info(`Skipping group contact ${chatId}...`); - progress.skipped += 1; - await job.updateProgress(progress); - continue; - } - } else if (isLidUser(chatId)) { - if (!options.contacts.lids) { - this.logger.info(`Skipping PN contact ${chatId}...`); - progress.skipped += 1; - await job.updateProgress(progress); - continue; - } - this.logger.info(`Skipping LID contact ${chatId}...`); - } else if (!isPnUser(chatId)) { - this.logger.info(`Skipping non-phone-number contact ${chatId}...`); + for await (const contact of contacts) { + this.signal.throwIfAborted(); + + // + // Show progress + // + const thetotal = total(progress); + if (options.progress && thetotal) { + if (thetotal % options.progress == 0) { + await this.activity.progress(progress); + } + } + // Delay on batch + if (thetotal && options.delay.batch) { + if (thetotal % batch == 0) { + await sleep(options.delay.batch); + } + } + + // Process + try { + if (!EngineHelper.ContactIsMy(contact)) { + progress.skipped += 1; + await job.updateProgress(progress); + continue; + } + const chatId = contact.id; + if (isJidGroup(chatId)) { + if (!options.contacts.groups) { + this.logger.debug(`Skipping group contact ${chatId}...`); progress.skipped += 1; await job.updateProgress(progress); continue; } - - this.logger.info(`Pulling ${chatId}...`); - await this.pullOneContact(options, chatId); - progress.ok += 1; - } catch (e) { - this.logger.error(`Error pulling contact ${contact.id}: ${e}`); - progress.errors += 1; + } else if (isLidUser(chatId)) { + if (!options.contacts.lids) { + this.logger.debug(`Skipping LID contact ${chatId}...`); + progress.skipped += 1; + await job.updateProgress(progress); + continue; + } + } else if (!isPnUser(chatId)) { + this.logger.info(`Skipping non-phone-number contact ${chatId}...`); + progress.skipped += 1; + await job.updateProgress(progress); continue; } - // update progress before signal fails - await job.updateProgress(progress); - this.signal.throwIfAborted(); - await sleep(options.delay.contact); + + this.logger.debug(`Pulling ${chatId}...`); + const result = await this.pullOneContact(options, chatId); + progress.created += result.created; + progress.updated += result.updated; + progress.avatar.updated += result.avatar.updated; + this.logger.info( + `Contact ${chatId}: created=${result.created}, updated=${result.updated}, avatar.updated=${result.avatar.updated}`, + ); + } catch (e) { + this.logger.error(`Error pulling contact ${contact.id}: ${e}`); + progress.errors += 1; } - await this.conversation.activity( - this.l.r('task.contacts.progress', { progress: progress }), - ); - await sleep(options.delay.batch); + // update progress before signal fails + await job.updateProgress(progress); + await sleep(options.delay.contact); } - // - // Final report - // - job.progress = progress; - const msg = ContactsPullStatusMessage(this.l, job, 'completed'); - await this.conversation.incoming(msg); + await this.activity.completed(progress); } private async pullOneContact(options: ContactsPullOptions, chatId: string) { // Contact const contactInfo = WhatsAppContactInfo(this.session, chatId, this.l); - let cwContact = await this.contactService.findOrCreateContact(contactInfo); + let [cwContact, created] = + await this.contactService.findOrCreateContact(contactInfo); // Attributes if (options.attributes) { const attributes = await contactInfo.Attributes(); @@ -204,42 +234,73 @@ class ContactsPullHandler { ); } // Avatar + let avatarUpdated = false; switch (options.avatar) { case 'if-missing': - await this.contactService.updateAvatar( + avatarUpdated = await this.contactService.updateAvatar( cwContact, contactInfo, AvatarUpdateMode.IF_MISSING, ); break; case 'update': - await this.contactService.updateAvatar( + avatarUpdated = await this.contactService.updateAvatar( cwContact, contactInfo, AvatarUpdateMode.ALWAYS, ); break; } + return { + created: created ? 1 : 0, + updated: created ? 0 : 1, + avatar: { + updated: avatarUpdated ? 1 : 0, + }, + }; + } +} + +class TaskActivity { + constructor( + private l: Locale, + private conversation: Conversation, + ) {} + + public async started(progress: Progress) { + await this.conversation.activity( + this.l.r('task.contacts.started', { progress: progress }), + ); + } + + public async progress(progress: Progress) { + await this.conversation.activity( + this.l.r('task.contacts.progress', { progress: progress }), + ); + } + + /** + * Final report + */ + public async completed(progress: Progress) { + await this.conversation.activity( + this.l.r('task.contacts.completed', { progress: progress }), + ); } } export function ContactsPullStatusMessage( l: Locale, - job, + job: Job, state: JobState | 'unknown', ) { const details = JobLink(job); - job.progress = lodash.merge({}, NullProgress, job.progress); + const progress = lodash.merge({}, NullProgress, job.progress); const payload = { - error: state === 'failed' || job.progress.errors > 0, - state: lodash.capitalize(state), - progress: job.progress, + error: state === 'failed' || state === 'unknown' || progress.errors > 0, + state: lodash.capitalize(state ?? 'unknown'), + progress: progress, details: details, - job: { - timestamp: l.FormatTimestampSec(job.timestamp), - processedOn: l.FormatTimestampSec(job.processedOn), - finishedOn: l.FormatTimestampSec(job.finishedOn), - }, }; return l.r('task.contacts.status', payload); } diff --git a/src/apps/chatwoot/consumers/task/messages.pull.ts b/src/apps/chatwoot/consumers/task/messages.pull.ts new file mode 100644 index 00000000..496e18b3 --- /dev/null +++ b/src/apps/chatwoot/consumers/task/messages.pull.ts @@ -0,0 +1,410 @@ +import { Processor } from '@nestjs/bullmq'; +import * as ms from 'ms'; +import { JOB_CONCURRENCY } from '@waha/apps/app_sdk/constants'; +import { JobLink } from '@waha/apps/app_sdk/JobUtils'; +import { QueueName } from '@waha/apps/chatwoot/consumers/QueueName'; +import { ChatWootTaskConsumer } from '@waha/apps/chatwoot/consumers/task/base'; +import { DIContainer } from '@waha/apps/chatwoot/di/DIContainer'; +import { Locale } from '@waha/apps/chatwoot/i18n/locale'; +import { TKey } from '@waha/apps/chatwoot/i18n/templates'; +import { SessionManager } from '@waha/core/abc/manager.abc'; +import { RMutexService } from '@waha/modules/rmutex/rmutex.service'; +import { Job, JobState } from 'bullmq'; +import * as lodash from 'lodash'; +import { PinoLogger } from 'nestjs-pino'; +import { WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; +import { Conversation } from '@waha/apps/chatwoot/client/Conversation'; +import { ILogger } from '@waha/apps/app_sdk/ILogger'; +import { + GetChatMessagesFilter, + GetChatMessagesQuery, + MessageSortField, +} from '@waha/structures/chats.dto'; +import { WAMessage } from '@waha/structures/responses.dto'; +import { EnsureSeconds } from '@waha/utils/timehelper'; +import { MessageAnyHandler } from '@waha/apps/chatwoot/consumers/waha/message.any'; +import { MessageReportInfo } from '@waha/apps/chatwoot/consumers/waha/base'; +import { MessageAckEmoji } from '@waha/apps/chatwoot/emoji'; +import { SortOrder } from '@waha/structures/pagination.dto'; +import { + ArrayPaginator, + PaginatorParams, +} from '@waha/apps/app_sdk/waha/Paginator'; +import { EngineHelper } from '@waha/apps/chatwoot/waha/engines'; +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'; + +export enum ChatID { + ALL = 'all', + SUMMARY = 'summary', +} + +export type MessagesPullOptions = { + chat: string; + batch: number; + progress: number | null; + period: { + start: number; + end: number; + }; + force: boolean; + timeout: { + media: number; + }; + media: boolean; + ignore: IgnoreJidConfig; +}; + +interface Progress { + ok: number; + exists: number; + ignored: number; + errors: number; + chats: string[]; + last?: number; +} + +const NullProgress: Progress = { + ok: 0, + exists: 0, + ignored: 0, + errors: 0, + chats: [], +}; + +function total(progress: Progress) { + return progress.ok + progress.exists + progress.ignored + progress.errors; +} + +/** + * Placeholder task consumer for pulling WhatsApp messages. + * Throws until the implementation is completed. + */ +@Processor(QueueName.TASK_MESSAGES_PULL, { concurrency: JOB_CONCURRENCY }) +export class TaskMessagesPullConsumer extends ChatWootTaskConsumer { + constructor(manager: SessionManager, log: PinoLogger, rmutex: RMutexService) { + super(manager, log, rmutex, TaskMessagesPullConsumer.name); + } + + protected ErrorHeaderKey(): TKey | null { + return null; + } + + protected async Process( + container: DIContainer, + job: Job, + signal: AbortSignal, + ) { + const options: MessagesPullOptions = job.data.options; + const locale = container.Locale(); + const conversation = this.conversationForReport(container, job.data.body); + const waha = container.WAHASelf(); + const session = new WAHASessionAPI(job.data.session, waha); + const handler = new MessagesPullHandler( + container, + signal, + container.Logger(), + conversation, + session, + locale, + ); + if (options.chat === ChatID.SUMMARY) { + return await handler.summary(options, job); + } + + return await handler.handle(options, job); + } +} + +class MessagesPullHandler { + private activity: TaskActivity; + + constructor( + protected container: DIContainer, + protected signal: AbortSignal, + protected logger: ILogger, + conversation: Conversation, + protected session: WAHASessionAPI, + protected l: Locale, + ) { + this.activity = new TaskActivity(l, conversation); + } + + /** + * Calculate summary progress including current task and direct children + */ + protected async summaryProgress(job: Job): Promise { + const current = lodash.merge({}, NullProgress, job.progress); + const childrenValues = await job.getChildrenValues(); + const values: Progress[] = [current, ...Object.values(childrenValues)]; + let progress = lodash.merge({}, NullProgress); + + // Result + for (const p of values) { + progress.ok += p.ok; + progress.exists += p.exists; + progress.ignored += p.ignored; + progress.errors += p.errors; + } + // Chats + progress.chats = lodash.uniq(lodash.flatten(values.map((p) => p.chats))); + // Last timestamp + progress.last = lodash.max(values.map((p) => p.last || 0)); + if (progress.last === 0) { + progress.last = undefined; + } + return progress; + } + + async summary(options: MessagesPullOptions, job: Job): Promise { + const progress = await this.summaryProgress(job); + await job.updateProgress(progress); + await this.activity.completed(progress, options); + return progress; + } + + async handle(options: MessagesPullOptions, job: Job): Promise { + const jids = new JidFilter(options.ignore); + let progress = lodash.merge({}, NullProgress, job.progress); + const batch = options.batch ?? 100; + + this.logger.debug( + `Pulling messages for session ${job.data.session} with batch size ${batch}...`, + ); + if (lodash.isEqual(progress, NullProgress)) { + await this.activity.started(job, options); + } + + // Handler + const container = this.container; + const info = new MessageReportInfo(); + const handler = new MessageAnyHistoryHandler( + job, + container.MessageMappingService(), + container.ContactConversationService(), + container.Logger(), + info, + this.session, + container.Locale(), + container.WAHASelf(), + ); + handler.force = options.force; + + // Messages + const lte = job.timestamp - options.period.end; + const gte = job.timestamp - options.period.start; + const filters: GetChatMessagesFilter = { + 'filter.timestamp.lte': EnsureSeconds(lte), + 'filter.timestamp.gte': EnsureSeconds(gte), + }; + + const query: GetChatMessagesQuery = { + limit: batch, + offset: 0, + downloadMedia: false, + sortBy: MessageSortField.TIMESTAMP, + sortOrder: SortOrder.ASC, + }; + const params: PaginatorParams = { + processed: total(progress), + }; + const paginator = new ArrayPaginator(params); + let messages = paginator.iterate((processed: number) => { + query.offset = processed; + return this.session.getMessages(options.chat, query, filters, { + signal: this.signal, + }); + }); + + messages = EngineHelper.IterateMessages(messages); + const all = options.chat == ChatID.ALL; + + for await (let message of messages) { + this.signal.throwIfAborted(); + + // + // Show progress + // + const thetotal = total(progress); + if (options.progress && thetotal) { + if (thetotal % options.progress == 0) { + await this.activity.progress(progress, options); + } + } + + // Process + try { + progress.last = message.timestamp; + if (all && !jids.include(EngineHelper.ChatID(message))) { + progress.ignored += 1; + continue; + } + + if (isNullJid(EngineHelper.ChatID(message))) { + progress.ignored += 1; + continue; + } + + if ( + options.media && + message.hasMedia && + (await handler.ShouldProcessMessage(message)) + ) { + // Fetch media for the message + let signal = AbortSignal.timeout(options.timeout.media); + signal = AbortSignal.any([signal, this.signal]); + message = await this.session.getMessageById('all', message.id, true, { + signal: signal, + }); + } + + const chatwoot = await handler.handle(message); + progress.ok += chatwoot ? 1 : 0; + progress.exists += chatwoot ? 0 : 1; + if ( + chatwoot && + !progress.chats.includes(EngineHelper.ChatID(message)) + ) { + progress.chats.push(EngineHelper.ChatID(message)); + } + } catch (error) { + const renderer = new ErrorRenderer(); + const text = renderer.text(error); + this.logger.error(`Error: ${text}`); + this.logger.error(`Message:\n${JSON.stringify(message, null, 2)}`); + try { + const data = renderer.data(error); + this.logger.error(JSON.stringify(data, null, 2)); + } catch (err) { + this.logger.error( + `Error occurred while login details for error: ${err}`, + ); + } + progress.errors += 1; + } + // update progress before signal fails + await job.updateProgress(progress); + } + + if (job.parentKey) { + return await this.summaryProgress(job); + } + return await this.summary(options, job); + } +} + +export class TaskActivity { + constructor( + private l: Locale, + private conversation: Conversation, + ) {} + + public async details(data) { + await this.conversation.activity( + this.l.r('task.messages.details', { + prefix: IsCommandsChat(data.body) ? '' : CommandPrefix, + }), + ); + } + + public async started(job: Job, options: MessagesPullOptions) { + await this.conversation.activity( + this.l.r('task.messages.started', { + chat: options.chat, + period: period(options), + }), + ); + } + + public async progress(progress: Progress, options: MessagesPullOptions) { + const format: Intl.DateTimeFormatOptions = { + month: 'short', + day: 'numeric', + year: 'numeric', + }; + await this.conversation.activity( + this.l.r('task.messages.progress', { + progress: progress, + chat: options.chat, + last: this.l.FormatTimestampOpts(progress.last, format, false), + }), + ); + } + + /** + * Final report + */ + public async completed(progress, options) { + const msg = this.l.r('task.messages.completed', { + progress: progress, + period: period(options), + chat: options.chat, + }); + await this.conversation.activity(msg); + } +} + +class MessageAnyHistoryHandler extends MessageAnyHandler { + public force: boolean = false; + public shouldLogUnsupported = true; + + protected get delayFromMeAPI() { + return 0; + } + + async ShouldProcessMessage(payload: WAMessage) { + if (this.force) { + return true; + } + return await super.ShouldProcessMessage(payload); + } + + protected finalizeContent(content: string, payload: WAMessage): string { + let ack: any = null; + if (payload.fromMe) { + ack = { + emoji: MessageAckEmoji(payload.ack), + name: this.l.r(payload.ackName || 'UNKNOWN'), + }; + } + return this.l.r('whatsapp.history.message.wrapper', { + content: content, + payload: payload, + timestamp: this.l.FormatTimestamp(payload.timestamp, false), + ack: ack, + }); + } +} + +function period(options: MessagesPullOptions) { + const start = options.period.start; + const end = options.period.end; + if (end == 0) { + return ms(start); + } + return `${ms(start)}-${ms(end)}`; +} + +export function MessagesPullStatusMessage( + l: Locale, + job: Job, + state: JobState | 'unknown', +) { + const progress = lodash.merge({}, NullProgress, job.progress); + const details = JobLink(job); + const options = job.data.options as MessagesPullOptions; + const payload = { + chat: options.chat, + chats: progress.chats.length, + period: period(options), + error: state === 'failed' || state === 'unknown' || progress.errors > 0, + state: lodash.capitalize(state ?? 'unknown'), + progress: progress, + details: details, + last: l.FormatTimestamp(progress.last, false), + }; + + return l.r('task.messages.status', payload); +} diff --git a/src/apps/chatwoot/consumers/waha/base.ts b/src/apps/chatwoot/consumers/waha/base.ts index 8a7b333f..d4bcb94a 100644 --- a/src/apps/chatwoot/consumers/waha/base.ts +++ b/src/apps/chatwoot/consumers/waha/base.ts @@ -13,7 +13,7 @@ import { EventData } from '@waha/apps/chatwoot/consumers/types'; import { WhatsAppContactInfo } from '@waha/apps/chatwoot/contacts/WhatsAppContactInfo'; import { DIContainer } from '@waha/apps/chatwoot/di/DIContainer'; import { Locale } from '@waha/apps/chatwoot/i18n/locale'; -import { WAHASelf, WAHASessionAPI } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASelf, WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; import { AppRepository, ChatwootMessage, @@ -25,16 +25,13 @@ import { parseMessageIdSerialized } from '@waha/core/utils/ids'; import { RMutexService } from '@waha/modules/rmutex/rmutex.service'; import { WAHAEvents } from '@waha/structures/enums.dto'; import { MessageSource, WAMessageBase } from '@waha/structures/responses.dto'; -import { - WAHAWebhook, - WAHAWebhookMessageRevoked, -} from '@waha/structures/webhooks.dto'; import { sleep } from '@waha/utils/promiseTimeout'; import { Job } from 'bullmq'; import { PinoLogger } from 'nestjs-pino'; import { TKey } from '@waha/apps/chatwoot/i18n/templates'; import { isJidBroadcast, isJidGroup, toCusFormat } from '@waha/core/utils/jids'; +import { EngineHelper } from '@waha/apps/chatwoot/waha'; export function ListenEventsForChatWoot() { return [ @@ -177,7 +174,7 @@ export interface IMessageInfo { onMessageType(type: MessageType): void; } -class MessageReportInfo implements IMessageInfo { +export class MessageReportInfo implements IMessageInfo { public conversationId: number | null = null; public type: MessageType | null = null; @@ -212,10 +209,29 @@ export abstract class MessageBaseHandler { payload: Payload, ): Promise; + protected finalizeContent(content: string, payload: Payload): string { + return content; + } + abstract getReplyToWhatsAppID(payload: Payload): string | undefined; - async handle(data: WAHAWebhook) { - const payload = data.payload; + protected get delayFromMeAPI() { + return 3_000; + } + + async ShouldProcessMessage(payload: Payload): Promise { + const key = parseMessageIdSerialized(payload.id); + const chatwoot = await this.mappingService.getChatWootMessage({ + chat_id: toCusFormat(key.remoteJid), + message_id: key.id, + }); + if (chatwoot) { + return false; + } + return true; + } + + async handle(payload: Payload): Promise { // Find the type as soon as possible for error reporting const type = payload.fromMe ? MessageType.OUTGOING : MessageType.INCOMING; this.info.onMessageType(type); @@ -224,28 +240,25 @@ export abstract class MessageBaseHandler { // It also handles messages fromMe and ChatWoot (but does not if send via API) if (payload.fromMe && payload.source === MessageSource.API) { // Sleep a few seconds to save it to a database - await sleep(3_000); + await sleep(this.delayFromMeAPI); } - const key = parseMessageIdSerialized(payload.id); - const chatwoot = await this.mappingService.getChatWootMessage({ - chat_id: toCusFormat(key.remoteJid), - message_id: key.id, - }); - if (chatwoot) { - this.logger.info( - `Message '${payload.id}' already in ChatWoot: conversation.id=${chatwoot.conversation_id}, message.id=${chatwoot.message_id}`, - ); - return; + if (!(await this.ShouldProcessMessage(payload))) { + const log = `Skipping existing message '${payload.id}' from WhatsApp`; + this.logger.debug(log); + return null; } - - const contactInfo = WhatsAppContactInfo(this.session, payload.from, this.l); + const contactInfo = WhatsAppContactInfo( + this.session, + EngineHelper.ChatID(payload as any), + this.l, + ); const conversation = await this.repo.ConversationByContact(contactInfo); this.info.onConversationId(conversation.conversationId); const message = await this.buildChatWootMessage(payload); const response = await conversation.send(message); - this.logger.info( + this.logger.debug( `Created message as '${message.message_type}' from WhatsApp: ${response.id}`, ); await this.saveMapping(response, payload); @@ -288,14 +301,16 @@ export abstract class MessageBaseHandler { content = this.l.key(key).render({ text: content }); } - const chatId = payload.from; + const chatId = EngineHelper.ChatID(payload); // Add participant name to group messages const manyParticipants = isJidGroup(chatId) || isJidBroadcast(chatId); if (!payload.fromMe && manyParticipants) { const key = parseMessageIdSerialized(payload.id, true); - let participant = toCusFormat(key.participant); - const contact: any = await this.session.getContact(key.participant); + let participant = toCusFormat( + key.participant || (payload as any)._data.participant, + ); + const contact: any = await this.session.getContact(participant); const name = contact?.name || contact?.pushName || contact?.pushname; if (name) { participant = `${name} (${participant})`; @@ -315,6 +330,7 @@ export abstract class MessageBaseHandler { ); const type = payload.fromMe ? MessageType.OUTGOING : MessageType.INCOMING; + content = this.finalizeContent(content, payload); return { content: content, message_type: type, @@ -334,7 +350,7 @@ export abstract class MessageBaseHandler { return; } const chatwoot = await this.mappingService.getChatWootMessage({ - chat_id: payload.from, + chat_id: EngineHelper.ChatID(payload), message_id: replyToWhatsAppID, }); return chatwoot?.message_id; diff --git a/src/apps/chatwoot/consumers/waha/message.ack.ts b/src/apps/chatwoot/consumers/waha/message.ack.ts index c3e10092..e994acad 100644 --- a/src/apps/chatwoot/consumers/waha/message.ack.ts +++ b/src/apps/chatwoot/consumers/waha/message.ack.ts @@ -10,7 +10,7 @@ import { } from '@waha/apps/chatwoot/consumers/waha/base'; import { WhatsAppContactInfo } from '@waha/apps/chatwoot/contacts/WhatsAppContactInfo'; import { Locale } from '@waha/apps/chatwoot/i18n/locale'; -import { WAHASessionAPI } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; import { SessionManager } from '@waha/core/abc/manager.abc'; import { RMutexService } from '@waha/modules/rmutex/rmutex.service'; import { WAHAEvents } from '@waha/structures/enums.dto'; diff --git a/src/apps/chatwoot/consumers/waha/message.ack.utils.ts b/src/apps/chatwoot/consumers/waha/message.ack.utils.ts index d70a1a0c..6f4ea44c 100644 --- a/src/apps/chatwoot/consumers/waha/message.ack.utils.ts +++ b/src/apps/chatwoot/consumers/waha/message.ack.utils.ts @@ -2,13 +2,14 @@ import type { WAHAWebhookMessageAck } from '@waha/structures/webhooks.dto'; import { isJidCusFormat } from '@waha/utils/wa'; import { WAMessageAck } from '@waha/structures/enums.dto'; import { isLidUser, toCusFormat } from '@waha/core/utils/jids'; +import { EngineHelper } from '@waha/apps/chatwoot/waha'; export function ShouldMarkAsReadInChatWoot( event: WAHAWebhookMessageAck, ): boolean { // Mark as seen only if it's DM // Ignore groups and other multiple participants chats - const chatId = event.payload.from; + const chatId = EngineHelper.ChatID(event.payload); if (!isJidCusFormat(chatId) && !isLidUser(chatId)) { return false; } diff --git a/src/apps/chatwoot/consumers/waha/message.any.ts b/src/apps/chatwoot/consumers/waha/message.any.ts index c820300a..f6b4e026 100644 --- a/src/apps/chatwoot/consumers/waha/message.any.ts +++ b/src/apps/chatwoot/consumers/waha/message.any.ts @@ -8,7 +8,7 @@ import { IMessageInfo, MessageBaseHandler, } from '@waha/apps/chatwoot/consumers/waha/base'; -import { WAHASessionAPI } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; import { SessionManager } from '@waha/core/abc/manager.abc'; import { parseMessageIdSerialized } from '@waha/core/utils/ids'; import { RMutexService } from '@waha/modules/rmutex/rmutex.service'; @@ -30,6 +30,8 @@ import { UnsupportedMessage, resolveProtoMessage, } from '@waha/apps/chatwoot/messages/to/chatwoot'; +import { EngineHelper } from '@waha/apps/chatwoot/waha'; +import { HasMediaWithNoMediaMessage } from '@waha/apps/chatwoot/messages/to/chatwoot/HasMediaWithNoMediaMessage'; @Processor(QueueName.WAHA_MESSAGE_ANY, { concurrency: JOB_CONCURRENCY }) export class WAHAMessageAnyConsumer extends ChatWootWAHABaseConsumer { @@ -42,7 +44,7 @@ export class WAHAMessageAnyConsumer extends ChatWootWAHABaseConsumer { } GetChatId(event: WAHAWebhookMessageAny): string { - return event.payload.from; + return EngineHelper.ChatID(event.payload); } async Process( @@ -62,11 +64,13 @@ export class WAHAMessageAnyConsumer extends ChatWootWAHABaseConsumer { container.Locale(), container.WAHASelf(), ); - return await handler.handle(event); + return await handler.handle(event.payload); } } -class MessageAnyHandler extends MessageBaseHandler { +export class MessageAnyHandler extends MessageBaseHandler { + public shouldLogUnsupported: boolean = false; + protected async getMessage( payload: WAMessage, ): Promise { @@ -81,7 +85,7 @@ class MessageAnyHandler extends MessageBaseHandler { return msg; } - converter = new TextMessage(this.l, this.logger, this.waha); + converter = new TextMessage(this.l, this.logger, this.waha, this.job); msg = await converter.convert(payload, null); if (msg) { return msg; @@ -123,8 +127,19 @@ class MessageAnyHandler extends MessageBaseHandler { return msg; } + converter = new HasMediaWithNoMediaMessage(this.l, this.job); + msg = await converter.convert(payload, protoMessage); + if (msg) { + return msg; + } + converter = new UnsupportedMessage(this.l, this.job); msg = await converter.convert(payload, protoMessage); + if (this.shouldLogUnsupported) { + this.logger.warn( + `UnsupportedMessage:\n${JSON.stringify(payload, null, 2)}`, + ); + } return msg; } @@ -133,6 +148,9 @@ class MessageAnyHandler extends MessageBaseHandler { if (!replyTo) { return undefined; } + if (!replyTo.id) { + return undefined; + } const key = parseMessageIdSerialized(replyTo.id, true); return key.id; } diff --git a/src/apps/chatwoot/consumers/waha/message.edited.ts b/src/apps/chatwoot/consumers/waha/message.edited.ts index 224d911f..2bcd10db 100644 --- a/src/apps/chatwoot/consumers/waha/message.edited.ts +++ b/src/apps/chatwoot/consumers/waha/message.edited.ts @@ -18,12 +18,13 @@ import { import { Job } from 'bullmq'; import { PinoLogger } from 'nestjs-pino'; -import { WAHASessionAPI } from '../../session/WAHASelf'; +import { WAHASessionAPI } from '../../../app_sdk/waha/WAHASelf'; import { MessageEdited, MessageToChatWootConverter, resolveProtoMessage, } from '@waha/apps/chatwoot/messages/to/chatwoot'; +import { EngineHelper } from '@waha/apps/chatwoot/waha'; @Processor(QueueName.WAHA_MESSAGE_EDITED, { concurrency: JOB_CONCURRENCY }) export class WAHAMessageEditedConsumer extends ChatWootWAHABaseConsumer { @@ -36,7 +37,7 @@ export class WAHAMessageEditedConsumer extends ChatWootWAHABaseConsumer { } GetChatId(event: WAHAWebhookMessageEdited): string { - return event.payload.from; + return EngineHelper.ChatID(event.payload); } async Process( @@ -56,7 +57,7 @@ export class WAHAMessageEditedConsumer extends ChatWootWAHABaseConsumer { container.Locale(), container.WAHASelf(), ); - return await handler.handle(event); + return await handler.handle(event.payload); } } diff --git a/src/apps/chatwoot/consumers/waha/message.reaction.ts b/src/apps/chatwoot/consumers/waha/message.reaction.ts index 8d1826eb..44d2b0fe 100644 --- a/src/apps/chatwoot/consumers/waha/message.reaction.ts +++ b/src/apps/chatwoot/consumers/waha/message.reaction.ts @@ -9,7 +9,7 @@ import { IMessageInfo, } from '@waha/apps/chatwoot/consumers/waha/base'; import { MessageBaseHandler } from '@waha/apps/chatwoot/consumers/waha/base'; -import { WAHASessionAPI } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; import { SessionManager } from '@waha/core/abc/manager.abc'; import { parseMessageIdSerialized } from '@waha/core/utils/ids'; import { RMutexService } from '@waha/modules/rmutex/rmutex.service'; @@ -19,6 +19,7 @@ import { WAHAWebhookMessageReaction } from '@waha/structures/webhooks.dto'; import { Job } from 'bullmq'; import { PinoLogger } from 'nestjs-pino'; import { TKey } from '@waha/apps/chatwoot/i18n/templates'; +import { EngineHelper } from '@waha/apps/chatwoot/waha'; @Processor(QueueName.WAHA_MESSAGE_REACTION, { concurrency: JOB_CONCURRENCY }) export class WAHAMessageReactionConsumer extends ChatWootWAHABaseConsumer { @@ -37,7 +38,7 @@ export class WAHAMessageReactionConsumer extends ChatWootWAHABaseConsumer { // but for backward compatability we do this return event.payload.to; } - return event.payload.from; + return EngineHelper.ChatID(event.payload); } async Process( @@ -57,7 +58,7 @@ export class WAHAMessageReactionConsumer extends ChatWootWAHABaseConsumer { container.Locale(), container.WAHASelf(), ); - return await handler.handle(event); + return await handler.handle(event.payload); } } diff --git a/src/apps/chatwoot/consumers/waha/session.status.ts b/src/apps/chatwoot/consumers/waha/session.status.ts index 05743db3..ccb18006 100644 --- a/src/apps/chatwoot/consumers/waha/session.status.ts +++ b/src/apps/chatwoot/consumers/waha/session.status.ts @@ -11,8 +11,8 @@ import { IMessageInfo, } from '@waha/apps/chatwoot/consumers/waha/base'; import { Locale } from '@waha/apps/chatwoot/i18n/locale'; -import { WAHASelf } from '@waha/apps/chatwoot/session/WAHASelf'; -import { SessionStatusEmoji } from '@waha/apps/chatwoot/SessionStatusEmoji'; +import { WAHASelf } from '@waha/apps/app_sdk/waha/WAHASelf'; +import { SessionStatusEmoji } from '@waha/apps/chatwoot/emoji'; import { SessionManager } from '@waha/core/abc/manager.abc'; import { RMutexService } from '@waha/modules/rmutex/rmutex.service'; import { WAHAEvents, WAHASessionStatus } from '@waha/structures/enums.dto'; diff --git a/src/apps/chatwoot/contacts/WhatsAppContactInfo.ts b/src/apps/chatwoot/contacts/WhatsAppContactInfo.ts index deaade1e..9dc31ba0 100644 --- a/src/apps/chatwoot/contacts/WhatsAppContactInfo.ts +++ b/src/apps/chatwoot/contacts/WhatsAppContactInfo.ts @@ -2,7 +2,7 @@ import { public_contact_create_update_payload as Contact } from '@figuro/chatwoo import { ContactInfo } from '@waha/apps/chatwoot/client/ContactConversationService'; import { AttributeKey } from '@waha/apps/chatwoot/const'; import { Locale } from '@waha/apps/chatwoot/i18n/locale'; -import { WAHASessionAPI } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASessionAPI } from '@waha/apps/app_sdk/waha/WAHASelf'; import { Channel } from '@waha/structures/channels.dto'; import { CacheAsync } from '@waha/utils/Cache'; import { TKey } from '@waha/apps/chatwoot/i18n/templates'; @@ -91,6 +91,9 @@ class LidContactInfo extends ChatContactInfo { @CacheAsync() async jid() { const pn = await this.session.findPNByLid(this.chatId); + if (!pn) { + return null; + } return new JidContactInfo(this.session, pn, this.locale); } @@ -99,7 +102,7 @@ class LidContactInfo extends ChatContactInfo { if (jid) { return await jid.AvatarUrl(); } - return null; + return await this.session.getChatPicture(this.chatId); } @CacheAsync() @@ -109,7 +112,6 @@ class LidContactInfo extends ChatContactInfo { if (jid) { attributes = await jid.Attributes(); } - attributes[AttributeKey.WA_CHAT_ID] = this.chatId; attributes[AttributeKey.WA_LID] = this.chatId; return attributes; } diff --git a/src/apps/chatwoot/di/DIContainer.ts b/src/apps/chatwoot/di/DIContainer.ts index 4c663802..01b8b135 100644 --- a/src/apps/chatwoot/di/DIContainer.ts +++ b/src/apps/chatwoot/di/DIContainer.ts @@ -15,7 +15,7 @@ import { } from '@waha/apps/chatwoot/dto/config.dto'; import { ChatWootErrorReporter } from '@waha/apps/chatwoot/error/ChatWootErrorReporter'; import { Locale } from '@waha/apps/chatwoot/i18n/locale'; -import { WAHASelf } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASelf } from '@waha/apps/app_sdk/waha/WAHASelf'; import { ChatwootMessageRepository, MessageMappingRepository, diff --git a/src/apps/chatwoot/emoji.ts b/src/apps/chatwoot/emoji.ts new file mode 100644 index 00000000..0a419305 --- /dev/null +++ b/src/apps/chatwoot/emoji.ts @@ -0,0 +1,37 @@ +import { WAHASessionStatus, WAMessageAck } from '@waha/structures/enums.dto'; + +export function SessionStatusEmoji(status: WAHASessionStatus): string { + switch (status) { + case WAHASessionStatus.STOPPED: + return '⚠️'; + case WAHASessionStatus.STARTING: + return '⏳'; + case WAHASessionStatus.SCAN_QR_CODE: + return '⚠️'; + case WAHASessionStatus.WORKING: + return '🟢'; + case WAHASessionStatus.FAILED: + return '🛑'; + default: + return '❓'; + } +} + +export function MessageAckEmoji(ack: WAMessageAck) { + switch (ack) { + case WAMessageAck.ERROR: + return '❌'; + case WAMessageAck.PENDING: + return '⏳'; + case WAMessageAck.SERVER: + return '✔️'; + case WAMessageAck.DEVICE: + return '✔️'; + case WAMessageAck.READ: + return '✅'; + case WAMessageAck.PLAYED: + return '✅'; + default: + return '❔'; + } +} diff --git a/src/apps/chatwoot/i18n/locale.ts b/src/apps/chatwoot/i18n/locale.ts index fb886913..78749e07 100644 --- a/src/apps/chatwoot/i18n/locale.ts +++ b/src/apps/chatwoot/i18n/locale.ts @@ -1,7 +1,9 @@ import * as Mustache from 'mustache'; +import * as lodash from 'lodash'; import { TemplatePayloads, TKey } from '@waha/apps/chatwoot/i18n/templates'; import Long from 'long'; import { ensureNumber } from '@waha/core/engines/noweb/utils'; +import { EnsureMilliseconds } from '@waha/utils/timehelper'; export class Locale { constructor(private readonly strings: Record) {} @@ -55,28 +57,9 @@ export class Locale { } } - FormatDatetime(date: Date | null): string | null { - if (!date) { - return null; - } - const options: any = this.strings['datetime'] || {}; - options.timeZone = options.timeZone || options.timezone || process.env.TZ; - return date.toLocaleString(this.locale, options); - } - - FormatDatetimeSec(date: Date | null) { - if (!date) { - return null; - } - const options: any = this.strings['datetime'] || {}; - options.second = '2-digit'; - options.timeZone = options.timeZone || options.timezone || process.env.TZ; - return date.toLocaleString(this.locale, options); - } - - FormatTimestampSec(timestamp: any) { - const date = this.ParseTimestamp(timestamp); - return this.FormatDatetimeSec(date); + FormatDatetime(date: Date | null, year): string | null { + const options: any = lodash.clone(this.strings['datetime'] || {}); + return this.FormatDatetimeOpts(date, options, year); } ParseTimestamp(timestamp: Long | string | number | null): Date | null { @@ -87,7 +70,7 @@ export class Locale { if (!Number.isFinite(value)) { return undefined; } - const milliseconds = value >= 1e12 ? value : value * 1000; + const milliseconds = EnsureMilliseconds(value); const date = new Date(milliseconds); if (Number.isNaN(date.getTime())) { return undefined; @@ -95,9 +78,44 @@ export class Locale { return date; } - FormatTimestamp(timestamp: Long | string | number | null): string | null { + FormatTimestamp( + timestamp: Long | string | number | null, + year: boolean = true, + ): string | null { const date = this.ParseTimestamp(timestamp); - return this.FormatDatetime(date); + return this.FormatDatetime(date, year); + } + + /** + * Format date using custom options + */ + + FormatDatetimeOpts( + date: Date, + options: Intl.DateTimeFormatOptions, + year: boolean, + ) { + options = lodash.cloneDeep(options); + if (!date) { + return null; + } + const opts: any = this.strings['datetime'] || {}; + // Copy timezone if any + options.timeZone = opts.timeZone || opts.timezone || process.env.TZ; + if (!year && date.getFullYear() === new Date().getFullYear()) { + // Hide year if current year + options.year = undefined; + } + return date.toLocaleDateString(this.locale, options); + } + + FormatTimestampOpts( + timestamp: Long | string | number | null, + options: Intl.DateTimeFormatOptions, + year: boolean = true, + ): string | null { + const date = this.ParseTimestamp(timestamp); + return this.FormatDatetimeOpts(date, options, year); } } diff --git a/src/apps/chatwoot/i18n/locales/ar-AE.yaml b/src/apps/chatwoot/i18n/locales/ar-AE.yaml index 639d8704..e8a4aa4e 100644 --- a/src/apps/chatwoot/i18n/locales/ar-AE.yaml +++ b/src/apps/chatwoot/i18n/locales/ar-AE.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # "options" من date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: ar-AE day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/bn-BD.yaml b/src/apps/chatwoot/i18n/locales/bn-BD.yaml index 776d1a0c..5000fd7d 100644 --- a/src/apps/chatwoot/i18n/locales/bn-BD.yaml +++ b/src/apps/chatwoot/i18n/locales/bn-BD.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # date.toLocaleString এর "options" # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: bn-BD day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/de-DE.yaml b/src/apps/chatwoot/i18n/locales/de-DE.yaml index f8601176..b9af4906 100644 --- a/src/apps/chatwoot/i18n/locales/de-DE.yaml +++ b/src/apps/chatwoot/i18n/locales/de-DE.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # "Optionen" von date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: de-DE day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/en-US.yaml b/src/apps/chatwoot/i18n/locales/en-US.yaml index 5a790dbf..cb619cc7 100644 --- a/src/apps/chatwoot/i18n/locales/en-US.yaml +++ b/src/apps/chatwoot/i18n/locales/en-US.yaml @@ -12,8 +12,9 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # "options" from date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit @@ -52,12 +53,12 @@ job.scheduled.error.header: Scheduled job failed to execute message.from.whatsapp: |- {{{text}}} - 📱 *Sent from WhatsApp* + *📱 Sent from WhatsApp* message.from.api: |- {{{text}}} - 📱 *Sent via API* + *📱 Sent via API* message.removed.in.whatsapp: |- ❌ *This message was deleted in WhatsApp* @@ -72,6 +73,35 @@ whatsapp.group.message: |- {{{text}}} +whatsapp.history.message.wrapper: |- + {{{content}}} + {{^payload.fromMe}} + + {{/payload.fromMe}} + `🗓️ {{{timestamp}}}` + +# WhatsApp Messages Ack +UNKNOWN: |- + Unknown + +ERROR: |- + Error + +PENDING: |- + Sending + +SERVER: |- + Server + +DEVICE: |- + Device + +READ: |- + Read + +PLAYED: |- + Played + whatsapp.reaction.added: |- *Reacted* {{{emoji}}} @@ -237,6 +267,18 @@ whatsapp.to.chatwoot.message.pix: |- 💳 **PIX Copy and Paste sent** {{/pixData}} +whatsapp.to.chatwoot.message.has.media.no.media: |- + {{#content}} + {{{content}}} + + {{/content}} + 🖼️⚠️ **The message has media, but we couldn't download it** + 📱 Please open **WhatsApp** to view it. + + {{^content}} + Details: [{{{details.text}}}]({{{details.url}}}) + {{/content}} + whatsapp.to.chatwoot.message.unsupported: |- ⚠️ **This message type is not supported in this inbox.** 📱 Please open **WhatsApp** to view it. @@ -338,6 +380,65 @@ cli.cmd.session.qr.description: |- cli.cmd.session.screenshot.description: |- Capture a session screenshot (available only with the **WEBJS** engine). +cli.cmd.messages.summary: |- + Pull WhatsApp messages. + +cli.cmd.messages.description: |- + 💬 Run `messages pull` to request message sync or `messages status` to review progress. + +cli.cmd.messages.action.description: |- + `pull` to start a message pull or `status` to inspect the active job. + +cli.cmd.messages.pull.option.chat: |- + Limit the pull to a specific WhatsApp Chat ID. + +cli.cmd.messages.pull.option.force: |- + Force ChatWoot to create messages even when the WhatsApp message is already present in the inbox mapping. + +cli.cmd.messages.pull.option.batch: |- + Number of messages to pull in each batch. + +cli.cmd.messages.pull.option.progress: |- + Report progress every N messages. Set to 0 to disable progress reports. + +cli.cmd.messages.pull.option.no-dm: |- + Skip direct messages while pulling history. + +cli.cmd.messages.pull.option.groups: |- + Include WhatsApp group chats when pulling history. + +cli.cmd.messages.pull.option.channels: |- + Include WhatsApp Channels when pulling history. + +cli.cmd.messages.pull.option.status: |- + Include WhatsApp status updates when pulling history. + +cli.cmd.messages.pull.option.broadcast: |- + Include broadcast lists when pulling history. + +cli.cmd.messages.pull.option.media: |- + Fetch media attachments (images, videos, files) when pulling history. + +cli.cmd.messages.pull.option.timeout-media: |- + Set how long to try downloading media for each message (defaults to `30s`). + +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`). + +cli.cmd.messages.pull.argument.start: |- + Set where the pull window starts, e.g., `1d` begins one day ago (defaults to `0d`, meaning now). + +cli.cmd.messages.pull.already-running: |- + ⏳ **Pull Messages** job is already running for **{{{chat}}}**. Wait for it to finish or stop it before starting a new one. + + Send `messages status` for details. + +cli.cmd.messages.status.not-found: |- + ℹ️ No active **Pull Messages** job found. Start one with `messages pull`. + +cli.cmd.messages.pull.no-chats-found: |- + No WhatsApp chats found to pull messages from for the period. + cli.cmd.contacts.summary: |- Pull WhatsApp contacts. @@ -353,6 +454,9 @@ cli.cmd.contacts.pull.description: |- cli.cmd.contacts.pull.option.batch: |- Number of contacts to pull in each batch. +cli.cmd.contacts.pull.option.progress: |- + Report progress every N contacts. Set to 0 to disable progress reports. + cli.cmd.contacts.pull.option.delay-contact: |- Delay between syncing each contact to avoid rate limiting in pulling avatars. @@ -360,13 +464,13 @@ cli.cmd.contacts.pull.option.delay-batch: |- Delay between batches of contacts. cli.cmd.contacts.pull.option.avatar: |- - Choose how to pull avatars for contacts: `skip`, `if-missing`, or `update`. + Choose how to pull avatars for contacts. cli.cmd.contacts.pull.option.groups: |- Include WhatsApp group contacts in the sync. -cli.cmd.contacts.pull.option.no-lids: |- - Skip WhatsApp LID contacts (linked-device JIDs) in the sync. +cli.cmd.contacts.pull.option.lids: |- + Include WhatsApp LID contacts (anonymous contacts) in the sync. cli.cmd.contacts.pull.option.no-attributes: |- Do not pull contact attributes if contact already exists (it speeds up the process). @@ -376,12 +480,6 @@ cli.cmd.contacts.pull.already-running: |- Send `contacts status` for details. -cli.cmd.contacts.pull.queued: |- - 👥 **Pull Contacts** job queued. - - We'll share a progress update every **{{batch}}** contacts. - Pulling contacts may take a while - check the latest status anytime with `contacts status`. - cli.cmd.options.job.attempts: |- Number of times to retry the job before marking it as failed. @@ -421,33 +519,52 @@ cli.help.options.defaultGroup: |- cli.help.globalOptions.title: |- **Global Options:** -task.contacts.status: |- - 👤 Pull Contacts: **{{state}}** - - ✅ Successful: {{{progress.ok}}} contacts - - ⚪ Skipped: {{{progress.skipped}}} contacts - - ⚠️ Errors: {{{progress.errors}}} contacts - - {{#job}} - ⌛ Timeline: - - Queued: {{{timestamp}}} - {{#processedOn}} - - Start At: {{{processedOn}}} - {{/processedOn}} - {{#finishedOn}} - - Finished At: {{{finishedOn}}} - {{/finishedOn}} - {{/job}} - - {{#error}} - Details: - [{{{details.text}}}]({{{details.url}}}) - {{/error}} - task.contacts.started: |- - Pulling contacts... + 👤⏳ Pulling contacts... Send "contacts status" for more details. task.contacts.progress: |- - Pull Contacts: ✅ {{progress.ok}}, ⚪ {{progress.skipped}}, ⚠️ {{progress.errors}} + 👤⏳ Contacts: 🆕 {{progress.created}}, 🔄 {{progress.updated}}, ⚪ {{progress.skipped}}, 🔴 {{progress.errors}}, 🖼️ {{progress.avatar.updated}} + +task.contacts.completed: |- + 👤✅ Contacts Completed → 🆕 {{progress.created}}, 🔄 {{progress.updated}}, ⚪ {{progress.skipped}}, 🔴 {{progress.errors}}, 🖼️ {{progress.avatar.updated}} + +task.contacts.status: |- + 👤 Contacts: **{{state}}** + - 🆕 Created: {{{progress.created}}} contacts + - 🔄 Updated: {{{progress.updated}}} contacts + - ⚪ Skipped: {{{progress.skipped}}} contacts + - 🔴 Errors: {{{progress.errors}}} contacts + - 🖼️ Avatar Updated: {{{progress.avatar.updated}}} contacts + + Details: + [{{{details.text}}}]({{{details.url}}}) + +task.messages.started: |- + ⏳ Pulling messages 💬 {{{chat}}} ⏱️ {{{period}}}... + +task.messages.details: |- + Send "{{{prefix}}}messages status" to see more details. + +task.messages.progress: |- + ⏳ Messages {{{chat}}}: 🆕 {{progress.ok}}, 🟢 {{progress.exists}}, ⚪ {{progress.ignored}}, 🔴 {{progress.errors}}{{#last}} | 📅 {{{last}}}{{/last}} + +task.messages.completed: |- + ✅ Messages {{{chat}}} in {{{period}}} → 🆕 {{progress.ok}}, 🟢 {{progress.exists}}, ⚪ {{progress.ignored}}, 🔴 {{progress.errors}} + +task.messages.status: |- + Messages **{{{chat}}}** in ⏱️ **{{{period}}}** period: + - ⚙️ Job State: **{{state}}** + {{#last}} + - 📅 Last Processed: **{{{last}}}** + {{/last}} + - 💬 Created in: **{{{chats}}}** chats + - 🆕 Created: **{{{progress.ok}}}** messages + - 🟢 Already exists: {{{progress.exists}}} messages + - ⚪ Ignored: {{{progress.ignored}}} messages + - 🔴 Errors: {{{progress.errors}}} messages + + Details: + [{{{details.text}}}]({{{details.url}}}) app.help.reminder: |- Send **help** to get the list of available commands. diff --git a/src/apps/chatwoot/i18n/locales/es-ES.yaml b/src/apps/chatwoot/i18n/locales/es-ES.yaml index 34bd4b49..2f13e9b0 100644 --- a/src/apps/chatwoot/i18n/locales/es-ES.yaml +++ b/src/apps/chatwoot/i18n/locales/es-ES.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # "opciones" de date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: es-ES day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/fa-IR.yaml b/src/apps/chatwoot/i18n/locales/fa-IR.yaml index cb3064f3..03c74c99 100644 --- a/src/apps/chatwoot/i18n/locales/fa-IR.yaml +++ b/src/apps/chatwoot/i18n/locales/fa-IR.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # گزینه‌های "options" از date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: fa-IR day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/fr-FR.yaml b/src/apps/chatwoot/i18n/locales/fr-FR.yaml index 9e745976..a58fa34b 100644 --- a/src/apps/chatwoot/i18n/locales/fr-FR.yaml +++ b/src/apps/chatwoot/i18n/locales/fr-FR.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # "options" de date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: fr-FR day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/he-IL.yaml b/src/apps/chatwoot/i18n/locales/he-IL.yaml index b981a223..8bf8ab25 100644 --- a/src/apps/chatwoot/i18n/locales/he-IL.yaml +++ b/src/apps/chatwoot/i18n/locales/he-IL.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # "אפשרויות" של date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: he-IL day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/hi-IN.yaml b/src/apps/chatwoot/i18n/locales/hi-IN.yaml index 211968de..04ecde16 100644 --- a/src/apps/chatwoot/i18n/locales/hi-IN.yaml +++ b/src/apps/chatwoot/i18n/locales/hi-IN.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # date.toLocaleString की "options" # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: hi-IN day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/id-ID.yaml b/src/apps/chatwoot/i18n/locales/id-ID.yaml index 8402ea9c..3206bcb9 100644 --- a/src/apps/chatwoot/i18n/locales/id-ID.yaml +++ b/src/apps/chatwoot/i18n/locales/id-ID.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # "opsi" dari date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: id-ID day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/pa-PK.yaml b/src/apps/chatwoot/i18n/locales/pa-PK.yaml index 26a21c7a..8617af2f 100644 --- a/src/apps/chatwoot/i18n/locales/pa-PK.yaml +++ b/src/apps/chatwoot/i18n/locales/pa-PK.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # date.toLocaleString ਦੀਆਂ "options" # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: pa-PK day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/pt-BR.yaml b/src/apps/chatwoot/i18n/locales/pt-BR.yaml index d5ca2d80..4f294860 100644 --- a/src/apps/chatwoot/i18n/locales/pt-BR.yaml +++ b/src/apps/chatwoot/i18n/locales/pt-BR.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # "opções" de date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: pt-BR day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/ru-RU.yaml b/src/apps/chatwoot/i18n/locales/ru-RU.yaml index c2f19b7b..274992d1 100644 --- a/src/apps/chatwoot/i18n/locales/ru-RU.yaml +++ b/src/apps/chatwoot/i18n/locales/ru-RU.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # "параметры" из date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: ru-RU day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/tr-TR.yaml b/src/apps/chatwoot/i18n/locales/tr-TR.yaml index 40a4df6e..49b44fc8 100644 --- a/src/apps/chatwoot/i18n/locales/tr-TR.yaml +++ b/src/apps/chatwoot/i18n/locales/tr-TR.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # date.toLocaleString için "options" # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: tr-TR day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/uk-UA.yaml b/src/apps/chatwoot/i18n/locales/uk-UA.yaml index b7035b78..260e106a 100644 --- a/src/apps/chatwoot/i18n/locales/uk-UA.yaml +++ b/src/apps/chatwoot/i18n/locales/uk-UA.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # "параметри" з date.toLocaleString # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: uk-UA day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/ur-PK.yaml b/src/apps/chatwoot/i18n/locales/ur-PK.yaml index 4f8625c8..4dae16c1 100644 --- a/src/apps/chatwoot/i18n/locales/ur-PK.yaml +++ b/src/apps/chatwoot/i18n/locales/ur-PK.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # date.toLocaleString کی "options" # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: ur-PK day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/i18n/locales/zh-CN.yaml b/src/apps/chatwoot/i18n/locales/zh-CN.yaml index c96b3add..b6441cc3 100644 --- a/src/apps/chatwoot/i18n/locales/zh-CN.yaml +++ b/src/apps/chatwoot/i18n/locales/zh-CN.yaml @@ -12,9 +12,10 @@ app.inbox.contact.avatar.url: 'https://raw.githubusercontent.com/devlikeapro/wah # date.toLocaleString 的 "options" # https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Global_Objects/Intl/DateTimeFormat/DateTimeFormat#options datetime: + weekday: short locales: zh-CN day: 2-digit - month: 2-digit + month: short year: numeric hour: 2-digit minute: 2-digit diff --git a/src/apps/chatwoot/messages/to/chatwoot/HasMediaWithNoMediaMessage.ts b/src/apps/chatwoot/messages/to/chatwoot/HasMediaWithNoMediaMessage.ts new file mode 100644 index 00000000..538b7dca --- /dev/null +++ b/src/apps/chatwoot/messages/to/chatwoot/HasMediaWithNoMediaMessage.ts @@ -0,0 +1,39 @@ +import { JobLink } from '@waha/apps/app_sdk/JobUtils'; +import { Locale } from '@waha/apps/chatwoot/i18n/locale'; +import { TKey } from '@waha/apps/chatwoot/i18n/templates'; +import { Job } from 'bullmq'; +import { ChatWootMessagePartial } from '@waha/apps/chatwoot/consumers/waha/base'; +import type { proto } from '@adiwajshing/baileys'; +import { WAMessage } from '@waha/structures/responses.dto'; +import { MessageToChatWootConverter } from '@waha/apps/chatwoot/messages/to/chatwoot'; + +export class HasMediaWithNoMediaMessage implements MessageToChatWootConverter { + constructor( + private readonly locale: Locale, + private readonly job: Job, + ) {} + + convert( + payload: WAMessage, + protoMessage: proto.Message | null, + ): ChatWootMessagePartial { + void protoMessage; + if (!payload.hasMedia) { + return null; + } + + const content = this.locale.r( + 'whatsapp.to.chatwoot.message.has.media.no.media', + { + content: null, + details: JobLink(this.job), + }, + ); + + return { + content, + attachments: [], + private: true, + }; + } +} diff --git a/src/apps/chatwoot/messages/to/chatwoot/TextMessage.ts b/src/apps/chatwoot/messages/to/chatwoot/TextMessage.ts index c471d147..5056f142 100644 --- a/src/apps/chatwoot/messages/to/chatwoot/TextMessage.ts +++ b/src/apps/chatwoot/messages/to/chatwoot/TextMessage.ts @@ -3,12 +3,14 @@ import { SendAttachment } from '@waha/apps/chatwoot/client/types'; import { ChatWootMessagePartial } from '@waha/apps/chatwoot/consumers/waha/base'; import { Locale } from '@waha/apps/chatwoot/i18n/locale'; import { TKey } from '@waha/apps/chatwoot/i18n/templates'; -import { WAHASelf } from '@waha/apps/chatwoot/session/WAHASelf'; +import { WAHASelf } from '@waha/apps/app_sdk/waha/WAHASelf'; import { isEmptyString } from './utils/proto'; import type { proto } from '@adiwajshing/baileys'; import { WAMessage } from '@waha/structures/responses.dto'; import { MessageToChatWootConverter } from '@waha/apps/chatwoot/messages/to/chatwoot'; import { WhatsappToMarkdown } from '@waha/apps/chatwoot/messages/to/chatwoot/utils/markdown'; +import { JobLink } from '@waha/apps/app_sdk/JobUtils'; +import { Job } from 'bullmq'; // eslint-disable-next-line @typescript-eslint/no-var-requires const mime = require('mime-types'); @@ -18,6 +20,7 @@ export class TextMessage implements MessageToChatWootConverter { private readonly locale: Locale, private readonly logger: ILogger, private readonly waha: WAHASelf, + private readonly job: Job, ) {} async convert( @@ -28,11 +31,25 @@ export class TextMessage implements MessageToChatWootConverter { const attachments = await this.getAttachments(payload); let content = this.locale.key(TKey.WA_TO_CW_MESSAGE).render({ payload }); if (isEmptyString(content) && attachments.length === 0) { + // No media, no content - return null so we can process it later return null; } if (isEmptyString(content)) { + // There's some media, but no content + // force content to be null for nice UI in ChatWoot content = null; } + if (attachments.length == 0 && payload.hasMedia) { + // Has media flag, but we couldn't find any media + // Add a warning at the end of content + content = this.locale.r( + 'whatsapp.to.chatwoot.message.has.media.no.media', + { + content: content, + details: JobLink(this.job), + }, + ); + } return { content: WhatsappToMarkdown(content), attachments, @@ -47,7 +64,7 @@ export class TextMessage implements MessageToChatWootConverter { } const media = payload.media!; - this.logger.info(`Downloading media from '${media.url}'...`); + this.logger.debug(`Downloading media from '${media.url}'...`); const buffer = await this.waha.fetch(media.url); const fileContent = buffer.toString('base64'); let filename = media.filename; diff --git a/src/apps/chatwoot/services/ChatWootScheduleService.ts b/src/apps/chatwoot/services/ChatWootScheduleService.ts index 525a7b91..138dbe91 100644 --- a/src/apps/chatwoot/services/ChatWootScheduleService.ts +++ b/src/apps/chatwoot/services/ChatWootScheduleService.ts @@ -5,6 +5,7 @@ import { Queue } from 'bullmq'; 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'; /** * Service for scheduling ChatWoot tasks @@ -20,7 +21,9 @@ export class ChatWootScheduleService { @InjectQueue(QueueName.SCHEDULED_CHECK_VERSION) private readonly checkVersionQueue: Queue, @InjectQueue(QueueName.TASK_CONTACTS_PULL) - private readonly importContactsQueue: Queue, + private readonly contactsPullQueue: Queue, + @InjectQueue(QueueName.TASK_MESSAGES_PULL) + private readonly messagesPullQueue: Queue, ) {} /** @@ -66,7 +69,7 @@ export class ChatWootScheduleService { ); // contacts - ContactsPullRemove(this.importContactsQueue, appId, this.logger).catch( + ContactsPullRemove(this.contactsPullQueue, appId, this.logger).catch( (reason) => { // Ignore errors this.logger.warn( @@ -74,6 +77,14 @@ export class ChatWootScheduleService { ); }, ); + + MessagesPullRemove(this.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/session/engines.ts b/src/apps/chatwoot/session/engines.ts deleted file mode 100644 index 100c89ce..00000000 --- a/src/apps/chatwoot/session/engines.ts +++ /dev/null @@ -1,74 +0,0 @@ -import { WhatsAppMessage } from '@waha/apps/chatwoot/storage'; -import { WAHAEngine } from '@waha/structures/enums.dto'; -import { getEngineName } from '@waha/version'; -import { Message as MessageInstance } from 'whatsapp-web.js/src/structures'; -import { toCusFormat } from '@waha/core/utils/jids'; - -interface IEngineHelper { - WhatsAppMessageKeys(message: any): WhatsAppMessage; -} - -class NOWEBHelper implements IEngineHelper { - WhatsAppMessageKeys(message: any): WhatsAppMessage { - const timestamp = parseInt(message.messageTimestamp) * 1000; - return { - timestamp: new Date(timestamp), - from_me: message.key.fromMe, - chat_id: toCusFormat(message.key.remoteJid), - message_id: message.key.id, - participant: message.key.participant, - }; - } -} - -class GOWSHelper implements IEngineHelper { - /** - * Parse API response and get the data - * API Response depends on engine right now - */ - WhatsAppMessageKeys(message: any): WhatsAppMessage { - const Info = message._data.Info; - const timestamp = new Date(Info.Timestamp).getTime(); - return { - timestamp: new Date(timestamp), - from_me: Info.IsFromMe, - chat_id: toCusFormat(Info.Chat), - message_id: Info.ID, - participant: Info.Sender ? toCusFormat(Info.Sender) : null, - }; - } -} - -class WEBJSHelper implements IEngineHelper { - /** - * Parse API response and get the data for WEBJS engine - */ - WhatsAppMessageKeys(message: MessageInstance): WhatsAppMessage { - return { - timestamp: new Date(message.timestamp * 1000), - from_me: message.fromMe, - chat_id: message.from, - message_id: message.id.id, - participant: message.author || null, - }; - } -} - -// Choose the right EngineHelper based on getEngineName() function -let engineHelper: IEngineHelper; - -switch (getEngineName()) { - case WAHAEngine.NOWEB: - engineHelper = new NOWEBHelper(); - break; - case WAHAEngine.GOWS: - engineHelper = new GOWSHelper(); - break; - case WAHAEngine.WEBJS: - engineHelper = new WEBJSHelper(); - break; - default: - engineHelper = new WEBJSHelper(); // Default to WEBJS as it's the default engine -} - -export const EngineHelper = engineHelper; diff --git a/src/apps/chatwoot/storage/ChatwootMessageRepository.ts b/src/apps/chatwoot/storage/ChatwootMessageRepository.ts index b4318dd3..17a76a6f 100644 --- a/src/apps/chatwoot/storage/ChatwootMessageRepository.ts +++ b/src/apps/chatwoot/storage/ChatwootMessageRepository.ts @@ -39,6 +39,7 @@ export class ChatwootMessageRepository { app_pk: this.appPk, id: id, }) + .orderBy('id', 'desc') .first(); } diff --git a/src/apps/chatwoot/storage/MessageMappingRepository.ts b/src/apps/chatwoot/storage/MessageMappingRepository.ts index 5e3d1205..28c363b7 100644 --- a/src/apps/chatwoot/storage/MessageMappingRepository.ts +++ b/src/apps/chatwoot/storage/MessageMappingRepository.ts @@ -63,6 +63,7 @@ export class MessageMappingRepository { app_pk: this.appPk, whatsapp_message_id: id, }) + .orderBy('id', 'desc') .first(); } @@ -72,6 +73,7 @@ export class MessageMappingRepository { app_pk: this.appPk, chatwoot_message_id: id, }) + .orderBy('id', 'desc') .first(); } @@ -85,6 +87,7 @@ export class MessageMappingRepository { chatwoot_message_id: id, part: part, }) + .orderBy('id', 'desc') .first(); } } diff --git a/src/apps/chatwoot/storage/WhatsAppMessageRepository.ts b/src/apps/chatwoot/storage/WhatsAppMessageRepository.ts index dab9cf09..ba24da54 100644 --- a/src/apps/chatwoot/storage/WhatsAppMessageRepository.ts +++ b/src/apps/chatwoot/storage/WhatsAppMessageRepository.ts @@ -39,6 +39,7 @@ export class WhatsAppMessageRepository { app_pk: this.appPk, id: id, }) + .orderBy('id', 'desc') .first(); } @@ -48,6 +49,7 @@ export class WhatsAppMessageRepository { app_pk: this.appPk, message_id: messageId, }) + .orderBy('id', 'desc') .first(); } diff --git a/src/apps/chatwoot/waha/engines.ts b/src/apps/chatwoot/waha/engines.ts new file mode 100644 index 00000000..724aa915 --- /dev/null +++ b/src/apps/chatwoot/waha/engines.ts @@ -0,0 +1,177 @@ +import * as lodash from 'lodash'; +import { WhatsAppMessage } from '@waha/apps/chatwoot/storage'; +import { WAHAEngine } from '@waha/structures/enums.dto'; +import { getEngineName } from '@waha/version'; +import { Message as MessageInstance } from 'whatsapp-web.js/src/structures'; +import { isLidUser, isPnUser, toCusFormat } from '@waha/core/utils/jids'; +import { WAMessage } from '@waha/structures/responses.dto'; + +interface IEngineHelper { + ChatID(message: WAMessage | any): string; + + WhatsAppMessageKeys(message: any): WhatsAppMessage; + + IterateMessages( + messages: AsyncGenerator, + ): AsyncGenerator; + + ContactIsMy(contact); + + FilterChatIdsForMessages(chats: string[]): string[]; + + SupportsAllChatForMessage(): boolean; +} + +class NOWEBHelper implements IEngineHelper { + ChatID(message: WAMessage): string { + return message.from; + } + + WhatsAppMessageKeys(message: any): WhatsAppMessage { + const timestamp = parseInt(message.messageTimestamp) * 1000; + return { + timestamp: new Date(timestamp), + from_me: message.key.fromMe, + chat_id: toCusFormat(message.key.remoteJid), + message_id: message.key.id, + participant: message.key.participant, + }; + } + + IterateMessages( + messages: AsyncGenerator, + ): AsyncGenerator { + return messages; + } + + FilterChatIdsForMessages(chats: string[]): string[] { + return chats; + } + + ContactIsMy(contact) { + return true; + } + + SupportsAllChatForMessage(): boolean { + return true; + } +} + +class GOWSHelper implements IEngineHelper { + ChatID(message: WAMessage): string { + return message.from; + } + + /** + * Parse API response and get the data + * API Response depends on engine right now + */ + WhatsAppMessageKeys(message: any): WhatsAppMessage { + const Info = message._data.Info; + const timestamp = new Date(Info.Timestamp).getTime(); + return { + timestamp: new Date(timestamp), + from_me: Info.IsFromMe, + chat_id: toCusFormat(Info.Chat), + message_id: Info.ID, + participant: Info.Sender ? toCusFormat(Info.Sender) : null, + }; + } + + IterateMessages( + messages: AsyncGenerator, + ): AsyncGenerator { + return messages; + } + + FilterChatIdsForMessages(chats: string[]): string[] { + return chats; + } + + SupportsAllChatForMessage(): boolean { + return true; + } + + ContactIsMy(contact) { + return true; + } +} + +class WEBJSHelper implements IEngineHelper { + ChatID(message: WAMessage): string { + return message._data?.id?.remote || message.from; + } + + /** + * Parse API response and get the data for WEBJS engine + */ + WhatsAppMessageKeys(message: MessageInstance): WhatsAppMessage { + return { + timestamp: new Date(message.timestamp * 1000), + from_me: message.fromMe, + chat_id: message.from, + message_id: message.id.id, + participant: message.author || null, + }; + } + + /** + * WEBJS API lacks server-side sorting hooks, so we buffer and sort by the unix timestamp in memory. + */ + async *IterateMessages( + messages: AsyncGenerator, + ): AsyncGenerator { + const buffer: T[] = []; + + for await (const message of messages) { + buffer.push(message); + } + + const sorted = lodash.sortBy(buffer, (item) => item.timestamp); + + for (const message of sorted) { + yield message; + } + } + + FilterChatIdsForMessages(chats: string[]): string[] { + if (chats.length == 2) { + const lidChat = chats.find(isLidUser); + const cusChat = chats.find(isPnUser); + if (lidChat && cusChat) { + return [lidChat]; + } + // WEBJS engine merges messages for @lid and @c.us + // into single chat, so it's fine to pull only from one + } + // Otherwise - return the original + return chats; + } + + SupportsAllChatForMessage(): boolean { + return false; + } + + ContactIsMy(contact) { + return contact.isMyContact; + } +} + +// Choose the right EngineHelper based on getEngineName() function +let engineHelper: IEngineHelper; + +switch (getEngineName()) { + case WAHAEngine.NOWEB: + engineHelper = new NOWEBHelper(); + break; + case WAHAEngine.GOWS: + engineHelper = new GOWSHelper(); + break; + case WAHAEngine.WEBJS: + engineHelper = new WEBJSHelper(); + break; + default: + engineHelper = new WEBJSHelper(); // Default to WEBJS as it's the default engine +} + +export const EngineHelper = engineHelper; diff --git a/src/apps/chatwoot/session/index.ts b/src/apps/chatwoot/waha/index.ts similarity index 100% rename from src/apps/chatwoot/session/index.ts rename to src/apps/chatwoot/waha/index.ts diff --git a/src/core/utils/jids.ts b/src/core/utils/jids.ts index 75cfc31c..24084368 100644 --- a/src/core/utils/jids.ts +++ b/src/core/utils/jids.ts @@ -78,6 +78,7 @@ export function toJID(chatId) { } export interface IgnoreJidConfig { + dm?: boolean; status: boolean; groups: boolean; channels: boolean; @@ -100,6 +101,10 @@ export class JidFilter { return false; } else if (this.ignore.channels && isJidNewsletter(jid)) { return false; + } else if (this.ignore.dm && isLidUser(jid)) { + return false; + } else if (this.ignore.dm && isPnUser(jid)) { + return false; } return true; } diff --git a/src/utils/promiseTimeout.ts b/src/utils/promiseTimeout.ts index 2bcd14ae..c0d8e6ee 100644 --- a/src/utils/promiseTimeout.ts +++ b/src/utils/promiseTimeout.ts @@ -34,6 +34,9 @@ export const promiseTimeout = function ( }; export async function sleep(ms: number) { + if (ms == 0) { + return; + } return new Promise((resolve) => setTimeout(resolve, ms)); } diff --git a/src/utils/timehelper.ts b/src/utils/timehelper.ts new file mode 100644 index 00000000..132a3c81 --- /dev/null +++ b/src/utils/timehelper.ts @@ -0,0 +1,21 @@ +/** + * If value is ms - we convert it to seconds + */ +export function EnsureSeconds(ms: number) { + if (!ms) { + return ms; + } + if (ms >= 1e12) { + return Math.floor(ms / 1000); + } + return ms; +} +export function EnsureMilliseconds(seconds: number) { + if (!seconds) { + return seconds; + } + if (seconds < 1e12) { + return seconds * 1000; + } + return seconds; +} diff --git a/yarn.lock b/yarn.lock index 99fb8600..0751441d 100644 --- a/yarn.lock +++ b/yarn.lock @@ -11221,6 +11221,13 @@ __metadata: languageName: node linkType: hard +"shell-quote@npm:^1.8.3": + version: 1.8.3 + resolution: "shell-quote@npm:1.8.3" + checksum: 550dd84e677f8915eb013d43689c80bb114860649ec5298eb978f40b8f3d4bc4ccb072b82c094eb3548dc587144bb3965a8676f0d685c1cf4c40b5dc27166242 + languageName: node + linkType: hard + "side-channel-list@npm:^1.0.0": version: 1.0.0 resolution: "side-channel-list@npm:1.0.0" @@ -11550,13 +11557,6 @@ __metadata: languageName: node linkType: hard -"string-argv@npm:^0.3.2": - version: 0.3.2 - resolution: "string-argv@npm:0.3.2" - checksum: 8703ad3f3db0b2641ed2adbb15cf24d3945070d9a751f9e74a924966db9f325ac755169007233e8985a39a6a292f14d4fee20482989b89b96e473c4221508a0f - languageName: node - linkType: hard - "string-length@npm:^4.0.1": version: 4.0.2 resolution: "string-length@npm:4.0.2" @@ -12617,8 +12617,8 @@ __metadata: rimraf: ^4.3.0 rxjs: ^7.8.1 sharp: ^0.33.4 + shell-quote: ^1.8.3 sqlite3: ^5.1.7 - string-argv: ^0.3.2 supertest: ^4.0.2 swagger-ui-express: ^4.1.4 ts-jest: ^29.1.3