feat(chatwoot): sync message status (delivered/read) from WhatsApp acks
Closes #2264 Co-authored-by: lucirlei <lucirleisantos@msn.com>
This commit is contained in:
1 parent
f1f25cf812
commit
1c6081a871
20 files changed
+410
-38
No files matched your search
@@ -1,6 +1,7 @@
|
||||
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 type { conversation_message_update } from '@figuro/chatwoot-sdk/dist/models/conversation_message_update';
|
||||
import { MessageStatus, MessageType } from '@waha/apps/chatwoot/client/types';
|
||||
|
||||
export class Conversation {
|
||||
public onError: (e: any) => void;
|
||||
@@ -43,6 +44,20 @@ export class Conversation {
|
||||
return this.send(data);
|
||||
}
|
||||
|
||||
/**
|
||||
* Update message status (sent/delivered/read/failed)
|
||||
*/
|
||||
public async updateMessageStatus(messageId: number, status: MessageStatus) {
|
||||
// The SDK model has no "status" field, but the endpoint accepts only status (and external_error)
|
||||
const data = { status: status } as unknown as conversation_message_update;
|
||||
return this.accountAPI.messages.update({
|
||||
accountId: this.accountId,
|
||||
conversationId: this.conversationId,
|
||||
messageId: messageId,
|
||||
data: data,
|
||||
});
|
||||
}
|
||||
|
||||
public async activity(text: string) {
|
||||
const data: conversation_message_create = {
|
||||
content: text,
|
||||
|
||||
@@ -36,6 +36,13 @@ export enum MessageType {
|
||||
ACTIVITY = 'activity',
|
||||
}
|
||||
|
||||
export enum MessageStatus {
|
||||
SENT = 'sent',
|
||||
DELIVERED = 'delivered',
|
||||
READ = 'read',
|
||||
FAILED = 'failed',
|
||||
}
|
||||
|
||||
export enum CustomAttributeType {
|
||||
TEXT = 0,
|
||||
NUMBER = 1,
|
||||
|
||||
@@ -38,6 +38,7 @@ import {
|
||||
} from '@waha/apps/chatwoot/dto/config.dto';
|
||||
import { Locale } from '@waha/apps/chatwoot/i18n/locale';
|
||||
import { isJidGroup } from '@waha/core/utils/jids';
|
||||
import { MessageStatusService } from '@waha/apps/chatwoot/services/MessageStatusService';
|
||||
|
||||
// eslint-disable-next-line @typescript-eslint/no-var-requires
|
||||
const mime = require('mime-types');
|
||||
@@ -65,6 +66,7 @@ export class ChatWootInboxMessageCreatedConsumer extends ChatWootInboxMessageCon
|
||||
session,
|
||||
container.ChatWootConfig(),
|
||||
container.Locale(),
|
||||
container.MessageStatusService(),
|
||||
);
|
||||
return await handler.handle(body);
|
||||
}
|
||||
@@ -77,6 +79,7 @@ export class MessageHandler {
|
||||
private session: WAHASessionAPI,
|
||||
private config: ChatWootConfig,
|
||||
private l: Locale,
|
||||
private statusService: MessageStatusService,
|
||||
) {}
|
||||
|
||||
async handle(body: any) {
|
||||
@@ -114,6 +117,7 @@ export class MessageHandler {
|
||||
// Send text (Part 1 if present)
|
||||
const attachments = message.attachments || [];
|
||||
const sendText = content && attachments.length !== 1;
|
||||
const parts = attachments.length + (sendText ? 1 : 0);
|
||||
if (sendText) {
|
||||
part += 1; // Text is the first possible part
|
||||
const exists = await this.getMapping(message, part);
|
||||
@@ -129,7 +133,7 @@ export class MessageHandler {
|
||||
});
|
||||
const msg = await this.sendTextMessage(chatId, text, replyTo, mentions);
|
||||
results.push(msg);
|
||||
await this.saveMapping(message, msg, part);
|
||||
await this.saveMapping(message, msg, part, parts);
|
||||
this.logger.info(`Text message sent: ${msg.id}`);
|
||||
}
|
||||
}
|
||||
@@ -155,20 +159,42 @@ export class MessageHandler {
|
||||
`File message sent: ${msg.id} - ${file.data_url} - ${file.file_type}`,
|
||||
);
|
||||
results.push(msg);
|
||||
await this.saveMapping(message, msg, part);
|
||||
await this.saveMapping(message, msg, part, parts);
|
||||
}
|
||||
if (this.config.conversations.syncMessageStatus) {
|
||||
await this.syncStatus(message, parts);
|
||||
}
|
||||
return results;
|
||||
}
|
||||
|
||||
/**
|
||||
* Push acks that arrived before the last part was mapped, best effort
|
||||
*/
|
||||
private async syncStatus(message: any, parts: number) {
|
||||
try {
|
||||
await this.statusService.sync({
|
||||
conversation_id: message.conversation.id,
|
||||
message_id: message.id,
|
||||
parts: parts,
|
||||
});
|
||||
} catch (err) {
|
||||
this.logger.warn(
|
||||
`ChatWoot => WhatsApp: error syncing status for Chatwoot message ${message.id}: ${err}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private async saveMapping(
|
||||
chatwootMessage: any,
|
||||
whatsappMessage: any,
|
||||
part: number,
|
||||
parts: number,
|
||||
): Promise<void> {
|
||||
const chatwoot: Omit<ChatwootMessage, 'id'> = {
|
||||
timestamp: new Date(chatwootMessage.created_at),
|
||||
conversation_id: chatwootMessage.conversation.id,
|
||||
message_id: chatwootMessage.id,
|
||||
parts: parts,
|
||||
};
|
||||
const whatsapp = EngineHelper.WhatsAppMessageKeys(whatsappMessage);
|
||||
await this.mappingService.map(chatwoot, whatsapp, part);
|
||||
|
||||
@@ -35,6 +35,7 @@ export class ChatWootInboxMessageUpdatedConsumer extends ChatWootInboxMessageCon
|
||||
session,
|
||||
container.ChatWootConfig(),
|
||||
container.Locale(),
|
||||
container.MessageStatusService(),
|
||||
);
|
||||
return await handler.handle(body);
|
||||
}
|
||||
|
||||
@@ -44,5 +44,7 @@ export class MessageCleanupConsumer extends ChatWootScheduledConsumer {
|
||||
.MessageMappingService()
|
||||
.cleanup(removeAfter);
|
||||
logger.info(`Removed ${removed} mappings for messages`);
|
||||
const acks = await container.MessageStatusService().cleanup(removeAfter);
|
||||
logger.info(`Removed ${acks} acks for messages`);
|
||||
}
|
||||
}
|
||||
@@ -53,7 +53,8 @@ export function ListenEventsForChatWoot(config: ChatWootConfig) {
|
||||
WAHAEvents.CALL_ACCEPTED,
|
||||
WAHAEvents.CALL_REJECTED,
|
||||
];
|
||||
if (config.conversations.markAsRead) {
|
||||
const conversations = config.conversations;
|
||||
if (conversations.markAsRead || conversations.syncMessageStatus) {
|
||||
events.push(WAHAEvents.MESSAGE_ACK);
|
||||
}
|
||||
return events;
|
||||
|
||||
@@ -14,12 +14,23 @@ 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';
|
||||
import { WAHAWebhookMessageAck } from '@waha/structures/webhooks.dto';
|
||||
import {
|
||||
WAHAWebhookMessageAck,
|
||||
WAMessageAckBody,
|
||||
} from '@waha/structures/webhooks.dto';
|
||||
import { Job } from 'bullmq';
|
||||
import { PinoLogger } from 'nestjs-pino';
|
||||
import { ShouldMarkAsReadInChatWoot } from '@waha/apps/chatwoot/consumers/waha/message.ack.utils';
|
||||
import { MultipleErrors } from '@waha/utils/errors';
|
||||
import {
|
||||
ShouldMarkAsReadInChatWoot,
|
||||
ShouldProcessAckInChatWoot,
|
||||
} from '@waha/apps/chatwoot/consumers/waha/message.ack.utils';
|
||||
import { parseMessageIdSerialized } from '@waha/core/utils/ids';
|
||||
import { MessageMappingService } from '@waha/apps/chatwoot/storage';
|
||||
import {
|
||||
ChatwootMessage,
|
||||
MessageMappingService,
|
||||
} from '@waha/apps/chatwoot/storage';
|
||||
import { MessageStatusService } from '@waha/apps/chatwoot/services/MessageStatusService';
|
||||
|
||||
@Processor(QueueName.WAHA_MESSAGE_ACK, { concurrency: JOB_CONCURRENCY })
|
||||
export class WAHAMessageAckConsumer extends ChatWootWAHABaseConsumer {
|
||||
@@ -32,7 +43,7 @@ export class WAHAMessageAckConsumer extends ChatWootWAHABaseConsumer {
|
||||
}
|
||||
|
||||
ShouldProcess(event: any): boolean {
|
||||
return ShouldMarkAsReadInChatWoot(event);
|
||||
return ShouldProcessAckInChatWoot(event);
|
||||
}
|
||||
|
||||
GetChatId(event: WAHAWebhookMessageAck): string {
|
||||
@@ -46,13 +57,20 @@ export class WAHAMessageAckConsumer extends ChatWootWAHABaseConsumer {
|
||||
const container = await this.DIContainer(job, job.data.app);
|
||||
const event = job.data.event as WAHAWebhookMessageAck;
|
||||
const session = new WAHASessionAPI(event.session, container.WAHASelf());
|
||||
const conversations = container.ChatWootConfig().conversations;
|
||||
const config: MessageAckHandlerConfig = {
|
||||
markAsRead: conversations.markAsRead,
|
||||
syncMessageStatus: conversations.syncMessageStatus,
|
||||
};
|
||||
const handler = new MessageAckHandler(
|
||||
config,
|
||||
container.ContactConversationService(),
|
||||
container.MessageMappingService(),
|
||||
container.Logger(),
|
||||
info,
|
||||
session,
|
||||
container.Locale(),
|
||||
container.MessageStatusService(),
|
||||
);
|
||||
try {
|
||||
await handler.handle(event);
|
||||
@@ -64,26 +82,72 @@ export class WAHAMessageAckConsumer extends ChatWootWAHABaseConsumer {
|
||||
}
|
||||
}
|
||||
|
||||
interface MessageAckHandlerConfig {
|
||||
markAsRead: boolean;
|
||||
syncMessageStatus: boolean;
|
||||
}
|
||||
|
||||
class MessageAckHandler {
|
||||
constructor(
|
||||
private readonly config: MessageAckHandlerConfig,
|
||||
private readonly contactConversationService: ContactConversationService,
|
||||
protected mappingService: MessageMappingService,
|
||||
private readonly logger: ILogger,
|
||||
private readonly info: IMessageInfo,
|
||||
private readonly session: WAHASessionAPI,
|
||||
private readonly locale: Locale,
|
||||
private readonly statusService: MessageStatusService,
|
||||
) {}
|
||||
|
||||
async handle(event: WAHAWebhookMessageAck): Promise<void> {
|
||||
const payload = event.payload;
|
||||
const promises: Promise<void>[] = [];
|
||||
if (this.config.syncMessageStatus) {
|
||||
promises.push(this.syncMessageStatus(event));
|
||||
}
|
||||
if (this.config.markAsRead && ShouldMarkAsReadInChatWoot(event)) {
|
||||
promises.push(this.markConversationAsRead(event));
|
||||
}
|
||||
const results = await Promise.allSettled(promises);
|
||||
const errors = results
|
||||
.filter((result) => result.status === 'rejected')
|
||||
.map((result: PromiseRejectedResult) => result.reason);
|
||||
if (errors.length > 0) {
|
||||
throw new MultipleErrors(errors);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* No chatwoot message - old message or some service ack, nothing to do
|
||||
*/
|
||||
private async getChatWootMessage(
|
||||
payload: WAMessageAckBody,
|
||||
): Promise<ChatwootMessage | null> {
|
||||
const key = parseMessageIdSerialized(payload.id);
|
||||
const chatwoot = await this.mappingService.getChatWootMessage({
|
||||
return this.mappingService.getChatWootMessage({
|
||||
chat_id: null,
|
||||
message_id: key.id,
|
||||
});
|
||||
// No chatwoot message found, so we don't mark it as read
|
||||
// Filters out old messages and some service ack messages
|
||||
}
|
||||
|
||||
private async syncMessageStatus(event: WAHAWebhookMessageAck) {
|
||||
const payload = event.payload;
|
||||
const key = parseMessageIdSerialized(payload.id);
|
||||
// Save the ack first, then look up the mapping.
|
||||
await this.statusService.record(
|
||||
key.id,
|
||||
payload.ack,
|
||||
new Date(event.timestamp),
|
||||
);
|
||||
const chatwoot = await this.getChatWootMessage(payload);
|
||||
if (!chatwoot) {
|
||||
return;
|
||||
}
|
||||
await this.statusService.sync(chatwoot);
|
||||
}
|
||||
|
||||
private async markConversationAsRead(event: WAHAWebhookMessageAck) {
|
||||
const payload = event.payload;
|
||||
const chatwoot = await this.getChatWootMessage(payload);
|
||||
if (!chatwoot) {
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -1,9 +1,13 @@
|
||||
import type { WAHAEngine, WAHAEvents } from '@waha/structures/enums.dto';
|
||||
import { ShouldMarkAsReadInChatWoot } from '@waha/apps/chatwoot/consumers/waha/message.ack.utils';
|
||||
import {
|
||||
ShouldMarkAsReadInChatWoot,
|
||||
ShouldProcessAckInChatWoot,
|
||||
} from '@waha/apps/chatwoot/consumers/waha/message.ack.utils';
|
||||
|
||||
interface TestParam {
|
||||
name: string;
|
||||
expected: boolean;
|
||||
ShouldMarkAsReadInChatWoot: boolean;
|
||||
ShouldProcessAckInChatWoot: boolean;
|
||||
GOWS: any;
|
||||
NOWEB: any;
|
||||
WEBJS: any;
|
||||
@@ -14,10 +18,11 @@ interface TestParam {
|
||||
* Participant - 11111111111 (lid - 011111111111111)
|
||||
* Group - 222222222222222222
|
||||
*/
|
||||
const ShouldMarkAsReadCases: TestParam[] = [
|
||||
const TestCases: TestParam[] = [
|
||||
{
|
||||
name: 'Participant read My message (READ ack)',
|
||||
expected: true,
|
||||
ShouldMarkAsReadInChatWoot: true,
|
||||
ShouldProcessAckInChatWoot: true,
|
||||
GOWS: {
|
||||
id: 'evt_00000000000000000000000000',
|
||||
session: 'default',
|
||||
@@ -131,7 +136,8 @@ const ShouldMarkAsReadCases: TestParam[] = [
|
||||
},
|
||||
{
|
||||
name: 'Participant received My message (DEVICE ack)',
|
||||
expected: false,
|
||||
ShouldMarkAsReadInChatWoot: false,
|
||||
ShouldProcessAckInChatWoot: true,
|
||||
GOWS: {
|
||||
id: 'evt_00000000000000000000000000',
|
||||
session: 'default',
|
||||
@@ -243,7 +249,8 @@ const ShouldMarkAsReadCases: TestParam[] = [
|
||||
},
|
||||
{
|
||||
name: 'Me read Participant message on another device (READ ack)',
|
||||
expected: false,
|
||||
ShouldMarkAsReadInChatWoot: false,
|
||||
ShouldProcessAckInChatWoot: false,
|
||||
GOWS: {
|
||||
id: 'evt_00000000000000000000000000',
|
||||
session: 'default',
|
||||
@@ -315,7 +322,8 @@ const ShouldMarkAsReadCases: TestParam[] = [
|
||||
},
|
||||
{
|
||||
name: 'Participant read My Group Message (READ ack)',
|
||||
expected: false,
|
||||
ShouldMarkAsReadInChatWoot: false,
|
||||
ShouldProcessAckInChatWoot: false,
|
||||
GOWS: {
|
||||
id: 'evt_00000000000000000000000000',
|
||||
session: 'default',
|
||||
@@ -446,15 +454,16 @@ const ShouldMarkAsReadCases: TestParam[] = [
|
||||
},
|
||||
];
|
||||
|
||||
const engines = ['GOWS', 'NOWEB', 'WEBJS'];
|
||||
|
||||
describe('ShouldMarkAsReadInChatWoot', () => {
|
||||
const engines = ['GOWS', 'NOWEB', 'WEBJS'];
|
||||
for (const engine of engines) {
|
||||
describe(engine, () => {
|
||||
for (const param of ShouldMarkAsReadCases) {
|
||||
const name = `[${param.expected}] ${param.name}`;
|
||||
for (const param of TestCases) {
|
||||
const name = `[${param.ShouldMarkAsReadInChatWoot}] ${param.name}`;
|
||||
const event = param[engine];
|
||||
const testfn = event ? test : test.skip;
|
||||
const expected = param.expected;
|
||||
const expected = param.ShouldMarkAsReadInChatWoot;
|
||||
testfn(name, () => {
|
||||
// Test
|
||||
const result = ShouldMarkAsReadInChatWoot(event);
|
||||
@@ -464,3 +473,21 @@ describe('ShouldMarkAsReadInChatWoot', () => {
|
||||
});
|
||||
}
|
||||
});
|
||||
|
||||
describe('ShouldProcessAckInChatWoot', () => {
|
||||
for (const engine of engines) {
|
||||
describe(engine, () => {
|
||||
for (const param of TestCases) {
|
||||
const name = `[${param.ShouldProcessAckInChatWoot}] ${param.name}`;
|
||||
const event = param[engine];
|
||||
const testfn = event ? test : test.skip;
|
||||
const expected = param.ShouldProcessAckInChatWoot;
|
||||
testfn(name, () => {
|
||||
// Test
|
||||
const result = ShouldProcessAckInChatWoot(event);
|
||||
expect(result).toBe(expected);
|
||||
});
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
@@ -1,13 +1,21 @@
|
||||
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 { isLidUser } from '@waha/core/utils/jids';
|
||||
import { EngineHelper } from '@waha/apps/chatwoot/waha';
|
||||
|
||||
export function ShouldMarkAsReadInChatWoot(
|
||||
const ACKS_TO_PROCESS = new Set([
|
||||
WAMessageAck.DEVICE,
|
||||
WAMessageAck.READ,
|
||||
WAMessageAck.PLAYED,
|
||||
]);
|
||||
|
||||
/**
|
||||
* Acks for OUR messages in DMs - the ones that can change Chatwoot message status
|
||||
*/
|
||||
export function ShouldProcessAckInChatWoot(
|
||||
event: WAHAWebhookMessageAck,
|
||||
): boolean {
|
||||
// Mark as seen only if it's DM
|
||||
// Ignore groups and other multiple participants chats
|
||||
const chatId = EngineHelper.ChatID(event.payload);
|
||||
if (!isJidCusFormat(chatId) && !isLidUser(chatId)) {
|
||||
@@ -15,18 +23,27 @@ export function ShouldMarkAsReadInChatWoot(
|
||||
}
|
||||
|
||||
const payload = event.payload;
|
||||
// Only READ and PLAYED
|
||||
const read = payload.ack === WAMessageAck.READ;
|
||||
const played = payload.ack == WAMessageAck.PLAYED;
|
||||
if (!read && !played) {
|
||||
if (!ACKS_TO_PROCESS.has(payload.ack)) {
|
||||
return false;
|
||||
}
|
||||
|
||||
// Only process when OUR message (fromMe: true) was read by recipient
|
||||
// Only process when OUR message (fromMe: true) was received or read by recipient
|
||||
if (!payload.fromMe) {
|
||||
return false;
|
||||
}
|
||||
|
||||
// Mark ChatWoot conversation as read when recipient reads our message
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Mark ChatWoot conversation as read when recipient reads our message
|
||||
*/
|
||||
export function ShouldMarkAsReadInChatWoot(
|
||||
event: WAHAWebhookMessageAck,
|
||||
): boolean {
|
||||
if (!ShouldProcessAckInChatWoot(event)) {
|
||||
return false;
|
||||
}
|
||||
// Only READ and PLAYED
|
||||
const ack = event.payload.ack;
|
||||
return ack === WAMessageAck.READ || ack === WAMessageAck.PLAYED;
|
||||
}
|
||||
@@ -21,8 +21,10 @@ import {
|
||||
ChatwootMessageRepository,
|
||||
MessageMappingRepository,
|
||||
MessageMappingService,
|
||||
WhatsAppAckRepository,
|
||||
WhatsAppMessageRepository,
|
||||
} from '@waha/apps/chatwoot/storage';
|
||||
import { MessageStatusService } from '@waha/apps/chatwoot/services/MessageStatusService';
|
||||
import { Job } from 'bullmq';
|
||||
import { Knex } from 'knex';
|
||||
import { i18n } from '@waha/apps/chatwoot/i18n';
|
||||
@@ -190,6 +192,20 @@ export class DIContainer {
|
||||
);
|
||||
}
|
||||
|
||||
@CacheSync()
|
||||
private WhatsAppAckRepository(): WhatsAppAckRepository {
|
||||
return new WhatsAppAckRepository(this.Knex(), this.AppPk());
|
||||
}
|
||||
|
||||
@CacheSync()
|
||||
public MessageStatusService(): MessageStatusService {
|
||||
return new MessageStatusService(
|
||||
this.MessageMappingService(),
|
||||
this.WhatsAppAckRepository(),
|
||||
this.ContactConversationService(),
|
||||
);
|
||||
}
|
||||
|
||||
public ChatWootErrorReporter(job: Job): ChatWootErrorReporter {
|
||||
return new ChatWootErrorReporter(this.Logger(), job, this.Locale());
|
||||
}
|
||||
@@ -234,6 +250,7 @@ export function ChatWootConfigDefaults(
|
||||
sort: ConversationSort.created_newest,
|
||||
status: null,
|
||||
markAsRead: true,
|
||||
syncMessageStatus: false,
|
||||
outgoing: ChatWootOutgoingMode.PRIVATE_NOTE,
|
||||
},
|
||||
};
|
||||
|
||||
@@ -52,6 +52,15 @@ export class ChatWootConversationsConfig {
|
||||
@IsBoolean()
|
||||
markAsRead?: boolean = true;
|
||||
|
||||
@ApiPropertyOptional({
|
||||
description:
|
||||
'Update Chatwoot message status (delivered/read) for messages sent from Chatwoot using WhatsApp acks. Disabled by default',
|
||||
default: false,
|
||||
})
|
||||
@IsOptional()
|
||||
@IsBoolean()
|
||||
syncMessageStatus?: boolean = false;
|
||||
|
||||
@ApiPropertyOptional({
|
||||
description:
|
||||
'How to show messages sent from WhatsApp (not from ChatWoot): ' +
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
import { Knex } from 'knex';
|
||||
|
||||
exports.up = async function (knex: Knex) {
|
||||
await knex.schema.alterTable('app_chatwoot_chatwoot_messages', (table) => {
|
||||
table.integer('parts').nullable();
|
||||
});
|
||||
// WhatsApp acks for our messages, kept apart from mappings because an ack can arrive before the mapping
|
||||
await knex.schema.createTable('app_chatwoot_whatsapp_acks', (table) => {
|
||||
table.increments('id');
|
||||
table.integer('app_pk');
|
||||
table
|
||||
.foreign('app_pk')
|
||||
.references('pk')
|
||||
.inTable('apps')
|
||||
.onDelete('CASCADE');
|
||||
table.string('message_id', 64);
|
||||
table.integer('ack');
|
||||
table.datetime('timestamp');
|
||||
table.unique(['app_pk', 'message_id'], {
|
||||
indexName: 'wa_ack_app_msg_unique',
|
||||
});
|
||||
table.index(['app_pk', 'timestamp'], 'wa_ack_app_ts_idx');
|
||||
});
|
||||
};
|
||||
|
||||
exports.down = async function (knex: Knex) {
|
||||
await knex.schema.dropTable('app_chatwoot_whatsapp_acks');
|
||||
await knex.schema.alterTable('app_chatwoot_chatwoot_messages', (table) => {
|
||||
table.dropColumn('parts');
|
||||
});
|
||||
};
|
||||
@@ -121,6 +121,7 @@ export class ChatWootAppService implements IAppService {
|
||||
const knex = manager.store.getWAHADatabase();
|
||||
const di = new DIContainer(appDb.pk, app.config, this.logger, knex);
|
||||
await di.MessageMappingService().purge();
|
||||
await di.MessageStatusService().purge();
|
||||
this.cleanCache(app);
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
import { ContactConversationService } from '@waha/apps/chatwoot/client/ContactConversationService';
|
||||
import { MessageStatus } from '@waha/apps/chatwoot/client/types';
|
||||
import { MessageMappingService } from '@waha/apps/chatwoot/storage/MessageMappingService';
|
||||
import { WhatsAppAckRepository } from '@waha/apps/chatwoot/storage/WhatsAppAckRepository';
|
||||
import { ChatwootMessage } from '@waha/apps/chatwoot/storage/types';
|
||||
import { WAMessageAck } from '@waha/structures/enums.dto';
|
||||
|
||||
export type ChatwootMessageStatusKey = Pick<
|
||||
ChatwootMessage,
|
||||
'conversation_id' | 'message_id' | 'parts'
|
||||
>;
|
||||
|
||||
/**
|
||||
* Mirrors WhatsApp acks for messages sent from Chatwoot into Chatwoot message status
|
||||
*/
|
||||
export class MessageStatusService {
|
||||
constructor(
|
||||
private readonly mappingService: MessageMappingService,
|
||||
private readonly ackRepository: WhatsAppAckRepository,
|
||||
private readonly contactConversationService: ContactConversationService,
|
||||
) {}
|
||||
|
||||
/**
|
||||
* Saves the ack for a WhatsApp message, keeps the highest one
|
||||
*/
|
||||
async record(
|
||||
messageId: string,
|
||||
ack: WAMessageAck,
|
||||
timestamp: Date,
|
||||
): Promise<void> {
|
||||
await this.ackRepository.record(messageId, ack, timestamp);
|
||||
}
|
||||
|
||||
/**
|
||||
* Pushes the status to Chatwoot once every part has an ack, uses the lowest one
|
||||
* @returns the status sent to Chatwoot, null if nothing to update yet
|
||||
*/
|
||||
async sync(message: ChatwootMessageStatusKey): Promise<MessageStatus | null> {
|
||||
// Messages mirrored from WhatsApp (and old rows) do not know how many parts were sent
|
||||
if (!message.parts) {
|
||||
return null;
|
||||
}
|
||||
const parts = await this.mappingService.getWhatsAppMessage(message);
|
||||
if (parts.length !== message.parts) {
|
||||
return null;
|
||||
}
|
||||
const ack = await this.ackRepository.minimumAck(
|
||||
parts.map((part) => part.message_id),
|
||||
);
|
||||
if (ack === null || ack < WAMessageAck.DEVICE) {
|
||||
return null;
|
||||
}
|
||||
let status = MessageStatus.DELIVERED;
|
||||
if (ack >= WAMessageAck.READ) {
|
||||
status = MessageStatus.READ;
|
||||
}
|
||||
// Chatwoot only moves status forward, so re-sending the same status is harmless
|
||||
await this.contactConversationService
|
||||
.ConversationById(message.conversation_id)
|
||||
.updateMessageStatus(message.message_id, status);
|
||||
return status;
|
||||
}
|
||||
|
||||
cleanup(removeAfter: Date): Promise<number> {
|
||||
return this.ackRepository.deleteOlderThan(removeAfter);
|
||||
}
|
||||
|
||||
purge(): Promise<number> {
|
||||
return this.ackRepository.deleteAll();
|
||||
}
|
||||
}
|
||||
@@ -77,6 +77,15 @@ export class MessageMappingRepository {
|
||||
.first();
|
||||
}
|
||||
|
||||
async getAllByChatwootMessageId(id: number): Promise<MessageMapping[]> {
|
||||
return this.knex(this.tableName)
|
||||
.where({
|
||||
app_pk: this.appPk,
|
||||
chatwoot_message_id: id,
|
||||
})
|
||||
.orderBy('part', 'asc');
|
||||
}
|
||||
|
||||
async getByChatwootMessageIdAndPart(
|
||||
id: number,
|
||||
part: number,
|
||||
|
||||
@@ -126,11 +126,11 @@ export class MessageMappingService {
|
||||
}
|
||||
const mappings = [];
|
||||
for (const message of messages) {
|
||||
const mapping =
|
||||
await this.messageMappingRepository.getByChatwootMessageId(message.id);
|
||||
if (mapping) {
|
||||
mappings.push(mapping);
|
||||
}
|
||||
const parts =
|
||||
await this.messageMappingRepository.getAllByChatwootMessageId(
|
||||
message.id,
|
||||
);
|
||||
mappings.push(...parts);
|
||||
}
|
||||
const whatsapp = [];
|
||||
for (const mapping of mappings) {
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
import { Knex } from 'knex';
|
||||
|
||||
export class WhatsAppAckRepository {
|
||||
static tableName = 'app_chatwoot_whatsapp_acks';
|
||||
|
||||
constructor(
|
||||
private readonly knex: Knex,
|
||||
private readonly appPk: number,
|
||||
) {}
|
||||
|
||||
get tableName() {
|
||||
return WhatsAppAckRepository.tableName;
|
||||
}
|
||||
|
||||
/**
|
||||
* Saves the ack for a WhatsApp message, keeps the highest one
|
||||
*/
|
||||
async record(messageId: string, ack: number, timestamp: Date): Promise<void> {
|
||||
const key = { app_pk: this.appPk, message_id: messageId };
|
||||
await this.knex(this.tableName)
|
||||
.insert({ ...key, ack: ack, timestamp: timestamp })
|
||||
.onConflict(['app_pk', 'message_id'])
|
||||
.ignore();
|
||||
await this.knex(this.tableName)
|
||||
.where(key)
|
||||
.where('ack', '<', ack)
|
||||
.update({ ack: ack });
|
||||
}
|
||||
|
||||
/**
|
||||
* The lowest ack across the messages, null if any message has no ack yet
|
||||
*/
|
||||
async minimumAck(messageIds: string[]): Promise<number | null> {
|
||||
if (messageIds.length === 0) {
|
||||
return null;
|
||||
}
|
||||
const rows = await this.knex(this.tableName)
|
||||
.where({ app_pk: this.appPk })
|
||||
.whereIn('message_id', messageIds)
|
||||
.select('ack');
|
||||
if (rows.length !== messageIds.length) {
|
||||
return null;
|
||||
}
|
||||
return Math.min(...rows.map((row) => row.ack));
|
||||
}
|
||||
|
||||
async deleteOlderThan(date: Date): Promise<number> {
|
||||
return this.knex(this.tableName)
|
||||
.where('app_pk', this.appPk)
|
||||
.andWhere('timestamp', '<', date)
|
||||
.del();
|
||||
}
|
||||
|
||||
/**
|
||||
* Deletes all rows for the app.
|
||||
*/
|
||||
async deleteAll(): Promise<number> {
|
||||
return this.knex(this.tableName).where({ app_pk: this.appPk }).delete();
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,7 @@ import { ChatwootMessageRepository } from './ChatwootMessageRepository';
|
||||
import { MessageMappingRepository } from './MessageMappingRepository';
|
||||
import { MessageMappingService } from './MessageMappingService';
|
||||
import { ChatwootMessage, MessageMapping, WhatsAppMessage } from './types';
|
||||
import { WhatsAppAckRepository } from './WhatsAppAckRepository';
|
||||
import { WhatsAppMessageRepository } from './WhatsAppMessageRepository';
|
||||
|
||||
// Export all types
|
||||
@@ -15,6 +16,7 @@ export {
|
||||
AppRepository,
|
||||
ChatwootMessageRepository,
|
||||
MessageMappingRepository,
|
||||
WhatsAppAckRepository,
|
||||
WhatsAppMessageRepository,
|
||||
};
|
||||
|
||||
|
||||
@@ -15,6 +15,9 @@ export interface ChatWootCombinedKey {
|
||||
export interface ChatwootMessage extends ChatWootCombinedKey {
|
||||
id?: number;
|
||||
timestamp: Date;
|
||||
// WhatsApp messages sent for it (text + attachments)
|
||||
// null - mirrored from WhatsApp
|
||||
parts?: number | null;
|
||||
}
|
||||
|
||||
export interface MessageMapping {
|
||||
|
||||
@@ -0,0 +1,9 @@
|
||||
/**
|
||||
* Several errors from parallel tasks as one, the message lists all of them
|
||||
*/
|
||||
export class MultipleErrors extends Error {
|
||||
constructor(public errors: any[]) {
|
||||
const messages = errors.map((error) => error?.message ?? String(error));
|
||||
super(`${errors.length} errors: ${messages.join('; ')}`);
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user