[core] ChatWoot - conversation selector

fix #1216
fix #1357
fix #1237
close #1343
close #1213
This commit is contained in:
devlikepro committed 2025-09-29 09:09:43 +07:00
1 parent de47d6e445
commit 61bb111768
15 files changed
+404 -93

No files matched your search

@@ -89,6 +89,7 @@ export class ChatwootWebhookController {
return { success: true };
default:
// Ignore other events
await this.chatWootQueueService.addJobToQueue(body.event, data);
return { success: true };
}
}
+13
View File
@@ -25,6 +25,8 @@ import { WAHASessionStatusConsumer } from './consumers/waha/session.status';
import { ChatWootQueueService } from './services/ChatWootQueueService';
import { ChatWootScheduleService } from './services/ChatWootScheduleService';
import { ChatWootWAHAQueueService } from './services/ChatWootWAHAQueueService';
import { ChatWootConversationCreatedConsumer } from './consumers/inbox/conversation_created';
import { ChatWootConversationStatusChangedConsumer } from '@waha/apps/chatwoot/consumers/inbox/conversation_status_changed';
const CONTROLLERS = [ChatwootWebhookController, ChatwootLocalesController];
@@ -65,6 +67,14 @@ const IMPORTS = lodash.flatten([
name: QueueName.INBOX_MESSAGE_UPDATED,
defaultJobOptions: merge(ExponentialRetriesJobOptions, JobRemoveOptions),
}),
RegisterAppQueue({
name: QueueName.INBOX_CONVERSATION_CREATED,
defaultJobOptions: merge(NoRetriesJobOptions, JobRemoveOptions),
}),
RegisterAppQueue({
name: QueueName.INBOX_CONVERSATION_STATUS_CHANGED,
defaultJobOptions: merge(NoRetriesJobOptions, JobRemoveOptions),
}),
RegisterAppQueue({
name: QueueName.INBOX_MESSAGE_DELETED,
defaultJobOptions: merge(ExponentialRetriesJobOptions, JobRemoveOptions),
@@ -79,6 +89,9 @@ const PROVIDERS = [
ChatWootInboxMessageCreatedConsumer,
ChatWootInboxMessageUpdatedConsumer,
ChatWootInboxMessageDeletedConsumer,
// Conversation events
ChatWootConversationCreatedConsumer,
ChatWootConversationStatusChangedConsumer,
ChatWootInboxCommandsConsumer,
WAHASessionStatusConsumer,
WAHAMessageAnyConsumer,
@@ -138,4 +138,26 @@ export class ContactConversationService {
public async InboxNotifications() {
return this.ConversationByContact(new InboxContactInfo(this.l));
}
public ResetCache(chatIds: Array<string>) {
this.logger.info(`Resetting cache chat ids: ${chatIds.join(', ')}`);
for (const chatId of chatIds) {
this.cache.delete(chatId);
}
}
public ResetMismatchedCache(chatIds: Array<string>, value: ConversationId) {
for (const chatId of chatIds) {
if (!this.cache.has(chatId)) {
continue;
}
const current = this.cache.get(chatId);
if (current !== value) {
this.logger.info(
`Resetting cache for chat id: ${chatId}, value changed from ${current} to ${value}`,
);
this.cache.delete(chatId);
}
}
}
}
@@ -5,7 +5,7 @@ import {
ChatWootInboxAPI,
} from '@waha/apps/chatwoot/client/interfaces';
import type { conversation } from '@figuro/chatwoot-sdk/dist/models/conversation';
import * as lodash from 'lodash';
import { ConversationSelector } from '@waha/apps/chatwoot/services/ConversationSelector';
export type ConversationResult = Pick<conversation, 'id' | 'account_id'>;
@@ -19,6 +19,7 @@ export class ConversationService {
private config: ChatWootAPIConfig,
private accountAPI: ChatwootClient,
private inboxAPI: ChatWootInboxAPI,
private selector: ConversationSelector,
private logger: ILogger,
) {}
@@ -28,13 +29,8 @@ export class ConversationService {
accountId: this.config.accountId,
id: contact.id,
})) as any;
const conversationsForInbox = lodash.filter(result.payload, {
inbox_id: this.config.inboxId,
}) as contact_conversations;
if (conversationsForInbox.length == 0) {
return null;
}
return conversationsForInbox[0];
const conversations = result.payload;
return this.selector.select(conversations);
}
private async create(contact: ContactIds): Promise<ConversationResult> {
+10
View File
@@ -15,6 +15,16 @@ export function GetChatID(contact: any): string | null {
return contact?.custom_attributes?.[AttributeKey.WA_CHAT_ID];
}
export function GetAllChatIDs(contact: any): Array<string> {
const attrs = contact?.custom_attributes || {};
const ids = [
attrs[AttributeKey.WA_CHAT_ID],
attrs[AttributeKey.WA_JID],
attrs[AttributeKey.WA_LID],
];
return ids.filter(Boolean);
}
export function FindChatID(contact: any): string | null {
if (GetJID(contact)) {
return GetJID(contact);
+7
View File
@@ -51,3 +51,10 @@ export enum CustomAttributeModel {
CONVERSATION = 0,
CONTACT = 1,
}
export enum ConversationStatus {
OPEN = 'open',
PENDING = 'pending',
SNOOZED = 'snoozed',
RESOLVED = 'resolved',
}
+2
View File
@@ -18,6 +18,8 @@ export enum QueueName {
//
INBOX_MESSAGE_CREATED = 'chatwoot.inbox | message_created',
INBOX_MESSAGE_UPDATED = 'chatwoot.inbox | message_updated',
INBOX_CONVERSATION_CREATED = 'chatwoot.inbox | conversation_created',
INBOX_CONVERSATION_STATUS_CHANGED = 'chatwoot.inbox | conversation_status_changed',
//
// ChatWoot Events - Artificial
//
+8 -1
View File
@@ -58,13 +58,20 @@ export abstract class ChatWootInboxMessageConsumer extends AppConsumer {
job: Job,
): Promise<any>;
protected GetConversationID(body) {
return body.conversation.id;
}
/**
* Process the job
* This method is called by the queue processor
*/
async processJob(job: Job<InboxData, any, EventName>): Promise<any> {
const body = job.data.body;
const key = ChatWootConversationKey(job.data.app, body.conversation.id);
const key = ChatWootConversationKey(
job.data.app,
this.GetConversationID(body),
);
return await this.withMutex(job, key, () =>
this.ProcessAndReportStatus(job),
);
@@ -0,0 +1,57 @@
import { Processor } from '@nestjs/bullmq';
import { JOB_CONCURRENCY } from '@waha/apps/app_sdk/constants';
import { QueueName } from '@waha/apps/chatwoot/consumers/QueueName';
import { ChatWootConversationKey } from '@waha/apps/chatwoot/consumers/mutex';
import { ChatWootInboxMessageConsumer } from '@waha/apps/chatwoot/consumers/inbox/base';
import { InboxData } from '@waha/apps/chatwoot/consumers/types';
import { Job } from 'bullmq';
import { GetAllChatIDs } from '@waha/apps/chatwoot/client/ids';
import { DIContainer } from '@waha/apps/chatwoot/di/DIContainer';
import { TKey } from '@waha/apps/chatwoot/i18n/templates';
import { ContactConversationService } from '@waha/apps/chatwoot/client/ContactConversationService';
import { SessionManager } from '@waha/core/abc/manager.abc';
import { PinoLogger } from 'nestjs-pino';
import { RMutexService } from '@waha/modules/rmutex';
@Processor(QueueName.INBOX_CONVERSATION_CREATED, {
concurrency: JOB_CONCURRENCY,
})
export class ChatWootConversationCreatedConsumer extends ChatWootInboxMessageConsumer {
constructor(
protected readonly manager: SessionManager,
log: PinoLogger,
rmutex: RMutexService,
) {
super(manager, log, rmutex, 'ChatWootConversationCreatedConsumer');
}
protected ErrorHeaderKey(): TKey | null {
return null;
}
protected GetConversationID(body) {
return body.id;
}
protected async Process(
container: DIContainer,
body: any,
job: Job,
): Promise<any> {
const handler = new ConversationCreatedHandler(
container.ContactConversationService(),
);
return handler.handle(body);
}
}
class ConversationCreatedHandler {
constructor(private service: ContactConversationService) {}
async handle(body: any) {
const ids = GetAllChatIDs(body?.meta?.sender);
if (!ids || ids.length === 0) {
return;
}
this.service.ResetMismatchedCache(ids, body.id);
}
}
@@ -0,0 +1,62 @@
import { Processor } from '@nestjs/bullmq';
import { JOB_CONCURRENCY } from '@waha/apps/app_sdk/constants';
import { QueueName } from '@waha/apps/chatwoot/consumers/QueueName';
import { ChatWootInboxMessageConsumer } from '@waha/apps/chatwoot/consumers/inbox/base';
import { Job } from 'bullmq';
import { DIContainer } from '../../di/DIContainer';
import { TKey } from '../../i18n/templates';
import { ContactConversationService } from '@waha/apps/chatwoot/client/ContactConversationService';
import { AttributeKey } from '@waha/apps/chatwoot/const';
import { ConversationSelector } from '@waha/apps/chatwoot/services/ConversationSelector';
import { GetAllChatIDs } from '@waha/apps/chatwoot/client/ids';
import { SessionManager } from '@waha/core/abc/manager.abc';
import { PinoLogger } from 'nestjs-pino';
import { RMutexService } from '@waha/modules/rmutex';
@Processor(QueueName.INBOX_CONVERSATION_STATUS_CHANGED, {
concurrency: JOB_CONCURRENCY,
})
export class ChatWootConversationStatusChangedConsumer extends ChatWootInboxMessageConsumer {
constructor(
protected readonly manager: SessionManager,
log: PinoLogger,
rmutex: RMutexService,
) {
super(manager, log, rmutex, 'ChatWootConversationStatusChangedConsumer');
}
protected ErrorHeaderKey(): TKey | null {
return null;
}
protected GetConversationID(body) {
return body.id;
}
protected async Process(
container: DIContainer,
body: any,
job: Job,
): Promise<any> {
const handler = new ConversationStatusChangedHandler(
container.ContactConversationService(),
container.ConversationSelector(),
);
return await handler.handle(body);
}
}
class ConversationStatusChangedHandler {
constructor(
private service: ContactConversationService,
private selector: ConversationSelector,
) {}
async handle(body) {
if (!this.selector.hasStatusFilter()) {
return;
}
const ids = GetAllChatIDs(body.meta?.sender);
this.service.ResetCache(ids);
}
}
+83 -81
View File
@@ -25,21 +25,17 @@ import {
import { Job } from 'bullmq';
import { Knex } from 'knex';
import { i18n } from '@waha/apps/chatwoot/i18n';
import { CacheSync } from '@waha/utils/Cache';
import {
ConversationSelector,
ConversationSort,
} from '@waha/apps/chatwoot/services/ConversationSelector';
/**
* Dependency Injection Container for ChatWoot
* Manages the creation and caching of various clients and repositories
*/
export class DIContainer {
private accountAPI: ChatwootClient;
private inboxAPI: ChatWootInboxAPI;
private contactService: ContactService;
private conversationService: ConversationService;
private contactConversationService: ContactConversationService;
private locale: Locale;
private messageMappingService: MessageMappingService;
private wahaSelf: WAHASelf;
/**
* Creates a new DIContainer with the given configuration
*/
@@ -68,123 +64,124 @@ export class DIContainer {
return this.knex;
}
@CacheSync()
public Locale(): Locale {
if (!this.locale) {
this.locale = i18n.locale(this.config.locale || DEFAULT_LOCALE);
this.locale = this.locale.override(this.ChatWootConfig().templates);
}
return this.locale;
let locale = i18n.locale(this.config.locale || DEFAULT_LOCALE);
locale = locale.override(this.ChatWootConfig().templates);
return locale;
}
/**
* Gets the AccountAPI client
* @returns ChatwootClient instance
*/
@CacheSync()
public AccountAPI(): ChatwootClient {
if (!this.accountAPI) {
this.accountAPI = new ChatwootClient({
config: {
basePath: this.config.url,
with_credentials: true,
credentials: 'include',
token: this.config.accountToken,
},
});
}
return this.accountAPI;
return new ChatwootClient({
config: {
basePath: this.config.url,
with_credentials: true,
credentials: 'include',
token: this.config.accountToken,
},
});
}
/**
* Gets the InboxAPI client
* @returns ChatWootInboxAPI instance
*/
@CacheSync()
public InboxAPI(): ChatWootInboxAPI {
if (!this.inboxAPI) {
const chatwootClientAPI = new ChatwootClient({
config: {
basePath: this.config.url,
with_credentials: true,
credentials: 'include',
token: this.config.inboxIdentifier,
},
});
this.inboxAPI = chatwootClientAPI.client as ChatWootInboxAPI;
}
return this.inboxAPI;
const chatwootClientAPI = new ChatwootClient({
config: {
basePath: this.config.url,
with_credentials: true,
credentials: 'include',
token: this.config.inboxIdentifier,
},
});
return chatwootClientAPI.client as ChatWootInboxAPI;
}
/**
* Gets the ContactService
* @returns ContactService instance
*/
@CacheSync()
private ContactService(): ContactService {
if (!this.contactService) {
this.contactService = new ContactService(
this.config,
this.AccountAPI(),
this.InboxAPI(),
this.logger,
);
}
return this.contactService;
return new ContactService(
this.config,
this.AccountAPI(),
this.InboxAPI(),
this.logger,
);
}
@CacheSync()
public ConversationSelector() {
const config = this.ChatWootConfig();
return new ConversationSelector({
sort: config.conversations.sort,
status: config.conversations.status,
inboxId: this.config.inboxId,
});
}
/**
* Gets the ConversationService
* @returns ConversationService instance
*/
@CacheSync()
private ConversationService(): ConversationService {
if (!this.conversationService) {
this.conversationService = new ConversationService(
this.config,
this.AccountAPI(),
this.InboxAPI(),
this.logger,
);
}
return this.conversationService;
return new ConversationService(
this.config,
this.AccountAPI(),
this.InboxAPI(),
this.ConversationSelector(),
this.logger,
);
}
/**
* Gets the ContactConversationService
* @returns ContactConversationService instance
*/
@CacheSync()
public ContactConversationService(): ContactConversationService {
if (!this.contactConversationService) {
this.contactConversationService = new ContactConversationService(
this.config,
this.ContactService(),
this.ConversationService(),
this.AccountAPI(),
this.logger,
this.Locale(),
);
}
return this.contactConversationService;
return new ContactConversationService(
this.config,
this.ContactService(),
this.ConversationService(),
this.AccountAPI(),
this.logger,
this.Locale(),
);
}
@CacheSync()
private ChatwootMessageRepository(): ChatwootMessageRepository {
return new ChatwootMessageRepository(this.Knex(), this.AppPk());
}
@CacheSync()
private WhatsAppMessageRepository(): WhatsAppMessageRepository {
return new WhatsAppMessageRepository(this.Knex(), this.AppPk());
}
@CacheSync()
private MessageMappingRepository(): MessageMappingRepository {
return new MessageMappingRepository(this.Knex(), this.AppPk());
}
@CacheSync()
public MessageMappingService(): MessageMappingService {
if (!this.messageMappingService) {
this.messageMappingService = new MessageMappingService(
this.Knex(),
this.WhatsAppMessageRepository(),
this.ChatwootMessageRepository(),
this.MessageMappingRepository(),
);
}
return this.messageMappingService;
return new MessageMappingService(
this.Knex(),
this.WhatsAppMessageRepository(),
this.ChatwootMessageRepository(),
this.MessageMappingRepository(),
);
}
public ChatWootErrorReporter(job: Job): ChatWootErrorReporter {
@@ -195,19 +192,20 @@ export class DIContainer {
* Gets the WAHASelf instance
* @returns WAHASelf instance
*/
@CacheSync()
public WAHASelf(): WAHASelf {
if (!this.wahaSelf) {
this.wahaSelf = new WAHASelf();
const logging = new AxiosLogging(this.Logger());
logging.applyTo(this.wahaSelf.client);
}
return this.wahaSelf;
const self = new WAHASelf();
const logging = new AxiosLogging(this.Logger());
logging.applyTo(self.client);
return self;
}
@CacheSync()
public CustomAttributesService() {
return new CustomAttributesService(this.config, this.AccountAPI());
}
@CacheSync()
public ChatWootConfig(): ChatWootConfig {
const defaults: ChatWootConfig = {
templates: {},
@@ -215,6 +213,10 @@ export class DIContainer {
commands: {
server: true,
},
conversations: {
sort: ConversationSort.created_newest,
status: null,
},
};
return lodash.defaults({}, this.config, defaults);
}
+21
View File
@@ -9,6 +9,11 @@ import {
} from 'class-validator';
import { Type } from 'class-transformer';
import { IsDynamicObject } from '@waha/nestjs/validation/IsDynamicObject';
import {
ConversationSelectorConfig,
ConversationSort,
} from '@waha/apps/chatwoot/services/ConversationSelector';
import { ConversationStatus } from '@waha/apps/chatwoot/client/types';
export const DEFAULT_LOCALE = 'en-US';
@@ -23,10 +28,20 @@ export enum LinkPreview {
HQ = 'HG',
}
export class ChatWootConversationsConfig {
@IsEnum(ConversationSort)
sort: ConversationSort;
@IsOptional()
@IsEnum(ConversationStatus, { each: true })
status: Array<ConversationStatus> | null;
}
export interface ChatWootConfig {
templates: Record<string, string>;
linkPreview: LinkPreview;
commands: ChatWootCommandsConfig;
conversations: ChatWootConversationsConfig;
}
export class ChatWootAppConfig implements ChatWootAPIConfig {
@@ -56,7 +71,13 @@ export class ChatWootAppConfig implements ChatWootAPIConfig {
@IsDynamicObject()
templates?: Record<string, string>;
@IsOptional()
@ValidateNested()
@Type(() => ChatWootCommandsConfig)
commands?: ChatWootCommandsConfig;
@IsOptional()
@ValidateNested()
@Type(() => ChatWootConversationsConfig)
conversations?: ChatWootConversationsConfig;
}
@@ -19,6 +19,10 @@ export class ChatWootQueueService {
private readonly messageUpdatedQueue: Queue,
@InjectQueue(QueueName.INBOX_MESSAGE_DELETED)
private readonly messageDeletedQueue: Queue,
@InjectQueue(QueueName.INBOX_CONVERSATION_CREATED)
private readonly conversationCreatedQueue: Queue,
@InjectQueue(QueueName.INBOX_CONVERSATION_STATUS_CHANGED)
private readonly conversationStatusChanged: Queue,
@InjectQueue(QueueName.INBOX_COMMANDS)
private readonly commandsQueue: Queue,
) {}
@@ -41,10 +45,14 @@ export class ChatWootQueueService {
*/
private getQueueForEvent(event: string): Queue | null {
switch (event) {
case EventName.CONVERSATION_CREATED:
return this.conversationCreatedQueue;
case EventName.MESSAGE_CREATED:
return this.messageCreatedQueue;
case EventName.MESSAGE_UPDATED:
return this.messageUpdatedQueue;
case EventName.CONVERSATION_STATUS_CHANGED:
return this.conversationStatusChanged;
case 'message_deleted':
return this.messageDeletedQueue;
case 'commands':
@@ -95,9 +103,9 @@ export class ChatWootQueueService {
*/
async addJobToQueue(event: string, data: InboxData): Promise<any> {
const queue = this.getQueueForEvent(event);
if (queue) {
return await this.add(queue, event, data);
if (!queue) {
return;
}
return { ignored: true, event };
await this.add(queue, event, data);
}
}
@@ -0,0 +1,72 @@
import { ConversationStatus } from '@waha/apps/chatwoot/client/types';
import { contact_conversations } from '@figuro/chatwoot-sdk/dist/models/contact_conversations';
import * as lodash from 'lodash';
import type { conversation } from '@figuro/chatwoot-sdk/dist/models/conversation';
export enum ConversationSort {
activity_newest = 'activity_newest',
created_newest = 'created_newest',
created_oldest = 'created_oldest',
activity_oldest = 'activity_oldest',
}
export type ConversationSelectorConfig = {
sort: ConversationSort;
status?: Array<ConversationStatus>;
inboxId: number;
};
export type ConversationResult = Pick<conversation, 'id' | 'account_id'>;
export class ConversationSelector {
constructor(private config: ConversationSelectorConfig) {}
hasStatusFilter() {
return this.config.status;
}
select(conversations: contact_conversations): ConversationResult | null {
conversations = this.filter(conversations);
conversations = this.sort(conversations);
return conversations[0] || null;
}
private filter(conversations: contact_conversations): contact_conversations {
// Filter by inbox id
conversations = lodash.filter(conversations, {
inbox_id: this.config.inboxId,
}) as contact_conversations;
// Filter by status
if (this.config.status && this.config.status.length > 0) {
conversations = lodash.filter(conversations, (conversation) => {
return this.config.status.includes(
conversation.status as ConversationStatus,
);
});
}
return conversations;
}
private sort(conversations: contact_conversations): contact_conversations {
let field = null;
let dir = null;
switch (this.config.sort) {
case ConversationSort.activity_newest:
[field, dir] = ['last_activity_at', 'desc'];
break;
case ConversationSort.created_newest:
[field, dir] = ['created_at', 'desc'];
break;
case ConversationSort.created_oldest:
[field, dir] = ['created_at', 'asc'];
break;
case ConversationSort.activity_oldest:
[field, dir] = ['last_activity_at', 'asc'];
}
if (!field || !dir) {
return conversations;
}
return lodash.orderBy(conversations, [field], [dir]);
}
}
+31
View File
@@ -21,3 +21,34 @@ export function CacheAsync() {
};
};
}
export function CacheSync() {
return function (
_target: any,
propertyKey: string,
descriptor: PropertyDescriptor,
) {
const original = descriptor.value as (...args: any[]) => any;
if (typeof original !== 'function') {
throw new Error('@CacheSync can only decorate methods');
}
const symbol = Symbol(`__cache_${propertyKey}`);
descriptor.value = function (...args: any[]) {
if (Object.prototype.hasOwnProperty.call(this, symbol)) {
return (this as any)[symbol];
}
const result = original.apply(this, args);
Object.defineProperty(this, symbol, {
value: result,
enumerable: false,
configurable: false,
writable: false,
});
return result;
};
return descriptor;
};
}