From ab781d915a0607e2f2bee45c476df7103aaf5c82 Mon Sep 17 00:00:00 2001 From: christianini-debug <338734608+christianini-debug@users.noreply.github.com> Date: Tue, 6 Oct 2026 16:21:23 -0300 Subject: [PATCH 1/2] fix(meta): preserve Coexistence webhook batches --- .gitignore | 3 +- package-lock.json | 1 + package.json | 4 +- .../channel/evohub/evohub.controller.ts | 2 + .../channel/meta/meta.controller.ts | 61 ++- .../channel/meta/whatsapp.business.service.ts | 138 ++---- test/meta-coexistence.test.ts | 417 ++++++++++++++++++ 7 files changed, 502 insertions(+), 124 deletions(-) create mode 100644 test/meta-coexistence.test.ts diff --git a/.gitignore b/.gitignore index 74b5ad7e5b..3b6b372378 100644 --- a/.gitignore +++ b/.gitignore @@ -34,7 +34,8 @@ lerna-debug.log* # Project related /instances/* !/instances/.gitkeep -/test/ +/test/* +!/test/meta-coexistence.test.ts /src/env.yml /store *.env diff --git a/package-lock.json b/package-lock.json index 55bf83fe04..84060d5274 100644 --- a/package-lock.json +++ b/package-lock.json @@ -88,6 +88,7 @@ "@types/qrcode-terminal": "^0.12.2", "commitizen": "^4.3.1", "cz-conventional-changelog": "^3.3.0", + "esbuild": "0.27.7", "eslint": "^10.4.1", "eslint-config-prettier": "^10.1.8", "eslint-plugin-import-x": "^4.16.2", diff --git a/package.json b/package.json index 4a5cade827..e1378ea824 100644 --- a/package.json +++ b/package.json @@ -21,7 +21,8 @@ "db:migrate:dev": "node runWithProvider.js \"rm -rf ./prisma/migrations && cp -r ./prisma/DATABASE_PROVIDER-migrations ./prisma/migrations && npx prisma migrate dev --schema ./prisma/DATABASE_PROVIDER-schema.prisma && cp -r ./prisma/migrations/* ./prisma/DATABASE_PROVIDER-migrations\"", "db:migrate:dev:win": "node runWithProvider.js \"xcopy /E /I prisma\\DATABASE_PROVIDER-migrations prisma\\migrations && npx prisma migrate dev --schema prisma\\DATABASE_PROVIDER-schema.prisma\"", "postinstall": "patch-package", - "prepare": "husky" + "prepare": "husky", + "test:coexistence": "tsx --test ./test/meta-coexistence.test.ts" }, "repository": { "type": "git", @@ -147,6 +148,7 @@ "@types/qrcode-terminal": "^0.12.2", "commitizen": "^4.3.1", "cz-conventional-changelog": "^3.3.0", + "esbuild": "0.27.7", "eslint": "^10.4.1", "eslint-config-prettier": "^10.1.8", "eslint-plugin-import-x": "^4.16.2", diff --git a/src/api/integrations/channel/evohub/evohub.controller.ts b/src/api/integrations/channel/evohub/evohub.controller.ts index 3958078edc..91e688ffed 100644 --- a/src/api/integrations/channel/evohub/evohub.controller.ts +++ b/src/api/integrations/channel/evohub/evohub.controller.ts @@ -1,6 +1,7 @@ import { MetaController } from '@api/integrations/channel/meta/meta.controller'; import { PrismaRepository } from '@api/repository/repository.service'; import { WAMonitoringService } from '@api/services/monitor.service'; +import { Integration } from '@api/types/wa.types'; import { ConfigService, EvolutionHub } from '@config/env.config'; import { Logger } from '@config/logger.config'; import * as crypto from 'crypto'; @@ -17,6 +18,7 @@ import * as crypto from 'crypto'; */ export class EvoHubController extends MetaController { private readonly hubLogger = new Logger('EvoHubController'); + protected readonly channelIntegration: string = Integration.EVOHUB; constructor( prismaRepository: PrismaRepository, diff --git a/src/api/integrations/channel/meta/meta.controller.ts b/src/api/integrations/channel/meta/meta.controller.ts index 558a22e98b..6148a95341 100644 --- a/src/api/integrations/channel/meta/meta.controller.ts +++ b/src/api/integrations/channel/meta/meta.controller.ts @@ -1,5 +1,6 @@ import { PrismaRepository } from '@api/repository/repository.service'; import { WAMonitoringService } from '@api/services/monitor.service'; +import { Integration } from '@api/types/wa.types'; import { Logger } from '@config/logger.config'; import axios from 'axios'; @@ -7,6 +8,7 @@ import { ChannelController, ChannelControllerInterface } from '../channel.contro export class MetaController extends ChannelController implements ChannelControllerInterface { private readonly logger = new Logger('MetaController'); + protected readonly channelIntegration: string = Integration.WHATSAPP_BUSINESS; constructor(prismaRepository: PrismaRepository, waMonitor: WAMonitoringService) { super(prismaRepository, waMonitor); @@ -15,54 +17,51 @@ export class MetaController extends ChannelController implements ChannelControll integrationEnabled: boolean; public async receiveWebhook(data: any) { - if (data.object === 'whatsapp_business_account') { - if (data.entry[0]?.changes[0]?.field === 'message_template_status_update') { - const template = await this.prismaRepository.template.findFirst({ - where: { templateId: `${data.entry[0].changes[0].value.message_template_id}` }, - }); + if (data?.object !== 'whatsapp_business_account') return { status: 'success' }; - if (!template) { - console.log('template not found'); - return; - } + for (const entry of data.entry ?? []) { + for (const change of entry.changes ?? []) { + if (change.field === 'message_template_status_update') { + const template = await this.prismaRepository.template.findFirst({ + where: { templateId: `${change.value?.message_template_id}` }, + }); - const { webhookUrl } = template; + if (!template) { + console.log('template not found'); + continue; + } - await axios.post(webhookUrl, data.entry[0].changes[0].value, { - headers: { - 'Content-Type': 'application/json', - }, - }); - return; - } + const { webhookUrl } = template; - data.entry?.forEach(async (entry: any) => { - const numberId = entry.changes[0].value.metadata.phone_number_id; + await axios.post(webhookUrl, change.value, { + headers: { + 'Content-Type': 'application/json', + }, + }); + continue; + } + + const numberId = change.value?.metadata?.phone_number_id; if (!numberId) { this.logger.error('WebhookService -> receiveWebhookMeta -> numberId not found'); - return { - status: 'success', - }; + continue; } const instance = await this.prismaRepository.instance.findFirst({ - where: { number: numberId }, + where: { number: numberId, integration: this.channelIntegration }, }); if (!instance) { this.logger.error('WebhookService -> receiveWebhookMeta -> instance not found'); - return { - status: 'success', - }; + continue; } - await this.waMonitor.waInstances[instance.name].connectToWhatsapp(data); + const channel = this.waMonitor.waInstances[instance.name]; + if (!channel) throw new Error('Instância da Cloud API indisponível para processar o webhook'); - return { - status: 'success', - }; - }); + await channel.connectToWhatsapp({ ...data, entry: [{ ...entry, changes: [change] }] }); + } } return { diff --git a/src/api/integrations/channel/meta/whatsapp.business.service.ts b/src/api/integrations/channel/meta/whatsapp.business.service.ts index 098654c2d0..ea3775ecc2 100644 --- a/src/api/integrations/channel/meta/whatsapp.business.service.ts +++ b/src/api/integrations/channel/meta/whatsapp.business.service.ts @@ -126,17 +126,31 @@ export class BusinessStartupService extends ChannelStartupService { public async connectToWhatsapp(data?: any): Promise { if (!data) return; - const content = data.entry[0].changes[0].value; - const normalizedContent = this.normalizeWebhookContent(content); - const remoteId = this.resolveRemoteId(normalizedContent); - try { - this.loadChatwoot(); - - await this.eventHandler(normalizedContent); - - if (remoteId) { - this.phoneNumber = createJid(remoteId); + await this.loadChatwoot(); + + for (const entry of data.entry ?? []) { + for (const change of entry.changes ?? []) { + const content = change.value; + if (!content) continue; + + const messages = Array.isArray(content.messages) ? content.messages : []; + // Meta uses the smb_message_echoes webhook field with a message_echoes array in its value. + const echoes = + [content.message_echoes, content.smb_message_echoes].find( + (value) => Array.isArray(value) && value.length > 0, + ) ?? []; + + if (messages.length) { + await this.eventHandler({ ...content, messages, statuses: undefined, isEcho: false }); + } + if (Array.isArray(echoes) && echoes.length) { + await this.eventHandler({ ...content, messages: echoes, statuses: undefined, isEcho: true }); + } + if (Array.isArray(content.statuses) && content.statuses.length) { + await this.eventHandler({ ...content, messages: undefined }); + } + } } } catch (error) { this.logger.error(error); @@ -144,70 +158,20 @@ export class BusinessStartupService extends ChannelStartupService { } } - private normalizeWebhookContent(content: any) { - if (!content || typeof content !== 'object') return content; - - const normalized = { ...content }; - const messageEchoes = Array.isArray(normalized?.message_echoes) ? normalized.message_echoes : undefined; - const smbMessageEchoes = Array.isArray(normalized?.smb_message_echoes) ? normalized.smb_message_echoes : undefined; - const echoes = messageEchoes?.length ? messageEchoes : smbMessageEchoes?.length ? smbMessageEchoes : undefined; - - if (!Array.isArray(normalized.messages) && Array.isArray(echoes) && echoes.length > 0) { - normalized.messages = echoes; - } - - return normalized; - } - private normalizePhoneNumber(value?: string) { return typeof value === 'string' ? value.replace(/\D/g, '') : ''; } - private resolveRemoteId(content: any) { - const firstMessage = content?.messages?.[0]; - const recipient = content?.statuses?.[0]?.recipient_id; - - const candidates = [firstMessage?.from, firstMessage?.to, recipient].filter(Boolean) as string[]; - if (candidates.length === 0) return undefined; - - const businessNumbers = [ - this.normalizePhoneNumber(content?.metadata?.display_phone_number), - this.normalizePhoneNumber(content?.metadata?.phone_number_id), - ].filter(Boolean); - - const externalCounterpart = candidates.find((candidate) => { - const normalizedCandidate = this.normalizePhoneNumber(candidate); - return normalizedCandidate && !businessNumbers.includes(normalizedCandidate); - }); - - return externalCounterpart ?? candidates[0]; - } - - private isCloudApiEchoPayload(received: any) { - return ( - (Array.isArray(received?.message_echoes) && received.message_echoes.length > 0) || - (Array.isArray(received?.smb_message_echoes) && received.smb_message_echoes.length > 0) - ); - } - private resolveMessageRemoteId(message: any, received: any) { - if (this.isCloudApiEchoPayload(received)) { - return message?.to ?? message?.from; - } - - return message?.from ?? message?.to; + return this.isCloudApiFromMe(message, received) ? message?.to : message?.from; } private isCloudApiFromMe(message: any, received: any) { - if (this.isCloudApiEchoPayload(received)) return true; + if (received.isEcho) return true; const from = this.normalizePhoneNumber(message?.from); const displayPhone = this.normalizePhoneNumber(received?.metadata?.display_phone_number); - const phoneNumberId = this.normalizePhoneNumber(received?.metadata?.phone_number_id); - - if (!from) return false; - - return from === displayPhone || from === phoneNumberId; + return !!from && from === displayPhone; } private isCloudApiStatusFromMe(item: any, received: any) { @@ -484,6 +448,7 @@ export class BusinessStartupService extends ChannelStartupService { try { let messageRaw: any; let pushName: any; + let persistedMedia = false; const incomingContact = received?.contacts?.[0]; if (incomingContact) { @@ -605,6 +570,7 @@ export class BusinessStartupService extends ChannelStartupService { const createdMessage = await this.prismaRepository.message.create({ data: messageRaw, }); + persistedMedia = true; await this.prismaRepository.media.create({ data: { @@ -786,7 +752,7 @@ export class BusinessStartupService extends ChannelStartupService { } } - this.sendDataWebhook(Events.MESSAGES_UPSERT, messageRaw); + await this.sendDataWebhook(Events.MESSAGES_UPSERT, messageRaw); await chatbotController.emit({ instance: { instanceName: this.instance.name, instanceId: this.instanceId }, @@ -795,17 +761,14 @@ export class BusinessStartupService extends ChannelStartupService { pushName: messageRaw.pushName, }); - if (!this.isMediaMessage(message) && message.type !== 'sticker') { + if (!persistedMedia) { await this.prismaRepository.message.create({ data: messageRaw, }); } - const contactPhone = incomingContact?.profile?.phone ?? incomingContact?.wa_id ?? remoteId; - if (!contactPhone) return; - const contactRaw: any = { - remoteJid: createJid(contactPhone), + remoteJid, pushName, // profilePicUrl: '', instanceId: this.instanceId, @@ -816,14 +779,7 @@ export class BusinessStartupService extends ChannelStartupService { } if (contact) { - const contactRaw: any = { - remoteJid: createJid(contactPhone), - pushName, - // profilePicUrl: '', - instanceId: this.instanceId, - }; - - this.sendDataWebhook(Events.CONTACTS_UPDATE, contactRaw); + await this.sendDataWebhook(Events.CONTACTS_UPDATE, contactRaw); if (this.configService.get('CHATWOOT').ENABLED && this.localChatwoot?.enabled) { await this.chatwootService.eventWhatsapp( @@ -834,21 +790,21 @@ export class BusinessStartupService extends ChannelStartupService { } await this.prismaRepository.contact.updateMany({ - where: { remoteJid: contact.remoteJid }, + where: { instanceId: this.instanceId, remoteJid: contact.remoteJid }, data: contactRaw, }); return; } - this.sendDataWebhook(Events.CONTACTS_UPSERT, contactRaw); + await this.sendDataWebhook(Events.CONTACTS_UPSERT, contactRaw); - this.prismaRepository.contact.create({ + await this.prismaRepository.contact.create({ data: contactRaw, }); } if (received.statuses) { for await (const item of received.statuses) { - const remoteId = item?.recipient_id ?? this.phoneNumber; + const remoteId = item?.recipient_id; if (!remoteId) continue; const key: any = { @@ -857,7 +813,7 @@ export class BusinessStartupService extends ChannelStartupService { fromMe: this.isCloudApiStatusFromMe(item, received), }; if (settings?.groups_ignore && key.remoteJid.includes('@g.us')) { - return; + continue; } if (key.remoteJid !== 'status@broadcast' && !key?.remoteJid?.match(/(:\d+)/)) { const findMessage = await this.prismaRepository.message.findFirst({ @@ -871,7 +827,7 @@ export class BusinessStartupService extends ChannelStartupService { }); if (!findMessage) { - return; + continue; } const findMessageKey: any = findMessage?.key ?? {}; @@ -907,7 +863,7 @@ export class BusinessStartupService extends ChannelStartupService { ); } - return; + continue; } const message: any = { @@ -934,6 +890,7 @@ export class BusinessStartupService extends ChannelStartupService { } } catch (error) { this.logger.error(error); + throw error; } } @@ -1015,8 +972,7 @@ export class BusinessStartupService extends ChannelStartupService { const database = this.configService.get('DATABASE'); const settings = await this.findSettings(); - if (content.messages && content.messages.length > 0) { - const message = content.messages[0]; + for (const message of content.messages ?? []) { this.logger.log(`Tipo de mensaje recibido: ${message.type}`); if ( @@ -1032,18 +988,18 @@ export class BusinessStartupService extends ChannelStartupService { message.type === 'button' || message.type === 'reaction' ) { - await this.messageHandle(content, database, settings); + await this.messageHandle({ ...content, messages: [message], statuses: undefined }, database, settings); } else { this.logger.warn(`Tipo de mensaje no reconocido: ${message.type}`); } - } else if (content.statuses) { - await this.messageHandle(content, database, settings); - } else { - this.logger.warn('No se encontraron mensajes ni estados en el contenido recibido'); + } + if (content.statuses?.length) { + await this.messageHandle({ ...content, messages: undefined }, database, settings); } } catch (error) { this.logger.error('Error en eventHandler:'); this.logger.error(error); + throw error; } } diff --git a/test/meta-coexistence.test.ts b/test/meta-coexistence.test.ts new file mode 100644 index 0000000000..24bd7cc8e3 --- /dev/null +++ b/test/meta-coexistence.test.ts @@ -0,0 +1,417 @@ +import assert from 'node:assert/strict'; +import { createRequire } from 'node:module'; +import { resolve } from 'node:path'; +import { test } from 'node:test'; + +import { build } from 'esbuild'; + +// Bundle the production controller and service, replacing only infrastructure boundaries. +const root = resolve(__dirname, '..'); +const productionClasses = build({ + stdin: { + contents: `export { BusinessStartupService } from ${JSON.stringify(process.env.COEX_SERVICE_SOURCE ?? './src/api/integrations/channel/meta/whatsapp.business.service')}; +export { MetaController } from ${JSON.stringify(process.env.COEX_CONTROLLER_SOURCE ?? './src/api/integrations/channel/meta/meta.controller')}; +export { EvoHubController } from './src/api/integrations/channel/evohub/evohub.controller'; +export { BaseChatbotController } from './src/api/integrations/chatbot/base-chatbot.controller';`, + resolveDir: root, + loader: 'ts', + }, + bundle: true, + write: false, + platform: 'node', + format: 'cjs', + packages: 'external', + tsconfig: resolve(root, 'tsconfig.json'), + plugins: [ + { + name: 'mock-infrastructure', + setup(builder) { + const modules = { + '@api/services/channel.service': `export class ChannelStartupService { + get instanceId() { return this.instance.id; } + }`, + '@api/server.module': `export const chatbotController = { + emit: async (data) => globalThis.__metaTestEmit(data), + };`, + '@utils/sendTelemetry': 'export const sendTelemetry = () => {};', + '@api/integrations/storage/s3/libs/minio.server': + 'export const uploadFile = () => {}; export const getObjectUrl = () => {};', + '@exceptions': + 'export class InternalServerErrorException extends Error {} export class BadRequestException extends Error {}', + '@config/logger.config': 'export class Logger { error() {} warn() {} log() {} }', + '../channel.controller': `export class ChannelController { + constructor(prismaRepository, waMonitor) { Object.assign(this, { prismaRepository, waMonitor }); } + }`, + './chatbot.controller': `export class ChatbotController { + constructor(prismaRepository, waMonitor) { Object.assign(this, { prismaRepository, waMonitor }); } + }`, + '@config/env.config': 'export const configService = { get: () => ({ ENABLE: false }) };', + }; + builder.onResolve({ filter: /.*/ }, (args) => + Object.hasOwn(modules, args.path) ? { path: args.path, namespace: 'test-boundary' } : undefined, + ); + builder.onLoad({ filter: /.*/, namespace: 'test-boundary' }, (args) => ({ + contents: modules[args.path], + loader: 'js', + })); + }, + }, + ], +}).then((result) => { + const compiled = { exports: {} as any }; + new Function('require', 'module', 'exports', result.outputFiles[0].text)( + createRequire(resolve(root, 'package.json')), + compiled, + compiled.exports, + ); + return compiled.exports; +}); + +const metadata = { display_phone_number: '15550783881', phone_number_id: '106540352242922' }; +const customer = '16505551234'; +const customerJid = `${customer}@s.whatsapp.net`; +const message = (id = 'wamid.phone-1', to = customer) => ({ + id, + from: metadata.display_phone_number, + to, + timestamp: '1739321024', + type: 'text', + text: { body: 'A colleague replied from the Business app' }, +}); +const envelope = (value: any, field = 'smb_message_echoes') => ({ + object: 'whatsapp_business_account', + entry: [{ id: 'test-waba', changes: [{ field, value }] }], +}); + +async function harness(instanceId = 'tenant-a') { + const { BusinessStartupService, MetaController } = await productionClasses; + const events: any[] = []; + const emitted: any[] = []; + const stored: any[] = []; + const updates: any[] = []; + const errors: any[] = []; + const service = Object.create(BusinessStartupService.prototype); + Object.assign(service, { + instance: { id: instanceId, name: instanceId }, + localSettings: {}, + localWebhook: {}, + localChatwoot: {}, + logger: { log() {}, warn() {}, error: (error: any) => errors.push(error) }, + configService: { get: () => ({ ENABLED: false, ENABLE: false, SAVE_DATA: {} }) }, + findSettings: async () => ({}), + loadChatwoot: async () => {}, + sendDataWebhook: async (event: string, data: any) => events.push({ event, data }), + prismaRepository: { + message: { + create: async ({ data }: any) => { + stored.push(data); + return { ...data, id: 'db-id' }; + }, + findFirst: async () => null, + }, + messageUpdate: { create: async ({ data }: any) => updates.push(data) }, + contact: { + findFirst: async () => null, + create: async () => {}, + updateMany: async (query: any) => updates.push(query), + }, + }, + }); + (globalThis as any).__metaTestEmit = async (data: any) => emitted.push(data); + const controller = new MetaController( + { + instance: { + findFirst: async ({ where }: any) => + where.number === metadata.phone_number_id ? { id: instanceId, name: instanceId } : null, + }, + }, + { waInstances: { [instanceId]: service } }, + ); + return { service, controller, events, emitted, stored, updates, errors }; +} + +test('a Business app echo reaches the agent and webhook as an outgoing message in the customer conversation', async () => { + const h = await harness(); + await h.service.connectToWhatsapp(envelope({ metadata, message_echoes: [message()] })); + const upserts = h.events.filter((item) => item.event === 'messages.upsert'); + assert.equal(upserts.length, 1, 'the phone message must reach MESSAGES_UPSERT before the handler resolves'); + assert.deepEqual(upserts[0].data.key, { id: 'wamid.phone-1', remoteJid: customerJid, fromMe: true }); + assert.equal(h.emitted.length, 1, 'the agent must receive the colleague message'); + assert.equal(h.emitted[0].msg.key.fromMe, true); + assert.equal(h.stored.length, 1); + assert.deepEqual(h.errors, []); +}); + +test('the HTTP controller waits for echo delivery before acknowledging it', async () => { + const h = await harness(); + await h.controller.receiveWebhook(envelope({ metadata, message_echoes: [message()] })); + assert.equal(h.emitted.length, 1); + assert.equal(h.stored.length, 1); +}); + +test('Chatwoot runs before webhook/bot and preserves distinct Chatwoot identifiers', async () => { + const h = await harness(); + const order: string[] = []; + h.service.configService.get = (key: string) => ({ ENABLED: key === 'CHATWOOT', ENABLE: false, SAVE_DATA: {} }); + h.service.localChatwoot.enabled = true; + h.service.chatwootService = { + eventWhatsapp: async (event: string) => { + if (event === 'messages.upsert') { + order.push('chatwoot'); + return { id: 10, inbox_id: 20, conversation_id: 30 }; + } + }, + }; + h.service.sendDataWebhook = async (event: string, data: any) => { + if (event === 'messages.upsert') { + order.push('webhook'); + assert.equal(data.chatwootMessageId, 10); + assert.equal(data.chatwootInboxId, 20); + assert.equal(data.chatwootConversationId, 30); + } + }; + (globalThis as any).__metaTestEmit = async () => order.push('bot'); + await h.service.connectToWhatsapp(envelope({ metadata, message_echoes: [message()] })); + assert.deepEqual(order, ['chatwoot', 'webhook', 'bot']); + assert.deepEqual(h.errors, []); +}); + +test('classic inbound messages still reach the customer conversation with fromMe false and no contacts', async () => { + const h = await harness(); + const inbound = { ...message('wamid.customer'), from: customer, to: metadata.display_phone_number }; + await h.service.connectToWhatsapp(envelope({ metadata, messages: [inbound] }, 'messages')); + assert.equal(h.emitted.length, 1); + assert.deepEqual(h.emitted[0].msg.key, { id: inbound.id, remoteJid: customerJid, fromMe: false }); + assert.deepEqual(h.errors, []); +}); + +test('empty messages and message_echoes do not hide the alternate smb_message_echoes array', async () => { + const h = await harness(); + await h.service.connectToWhatsapp( + envelope({ metadata, messages: [], message_echoes: [], smb_message_echoes: [message()] }), + ); + assert.equal(h.emitted.length, 1); + assert.equal(h.emitted[0].msg.key.fromMe, true); +}); + +test('every message in an echo batch is delivered with its own recipient', async () => { + const h = await harness(); + await h.service.connectToWhatsapp( + envelope({ metadata, message_echoes: [message(), message('wamid.phone-2', '16505559999')] }), + ); + assert.deepEqual( + h.emitted.map((item) => item.remoteJid), + [customerJid, '16505559999@s.whatsapp.net'], + ); +}); + +test('the controller routes every change to the correct tenant', async () => { + const { MetaController } = await productionClasses; + const calls: any[] = []; + const controller = new MetaController( + { instance: { findFirst: async ({ where }: any) => ({ name: where.number }) } }, + { + waInstances: Object.fromEntries( + ['number-a', 'number-b', 'number-c'].map((name) => [ + name, + { + connectToWhatsapp: async (data: any) => calls.push({ name, data }), + }, + ]), + ), + }, + ); + const values = ['number-a', 'number-b', 'number-c'].map((number) => ({ + metadata: { phone_number_id: number }, + message_echoes: [message()], + })); + const data = envelope(values[0]); + data.entry[0].changes.push({ field: 'smb_message_echoes', value: values[1] }); + data.entry.push({ id: 'second-waba', changes: [{ field: 'smb_message_echoes', value: values[2] }] }); + await controller.receiveWebhook(data); + assert.deepEqual( + calls.map((item) => item.name), + ['number-a', 'number-b', 'number-c'], + ); + assert.equal(calls[1].data.entry[0].changes.length, 1); + assert.equal(calls[1].data.entry[0].changes[0].value.metadata.phone_number_id, 'number-b'); + assert.equal(calls[2].data.entry.length, 1); + assert.equal(calls[2].data.entry[0].id, 'second-waba'); +}); + +test('Meta and EvoHub route the same phone number ID only to their own integration', async () => { + const { MetaController, EvoHubController } = await productionClasses; + const calls: string[] = []; + const integrations = ['WHATSAPP-BUSINESS', 'EVOHUB']; + const repository = { + instance: { + findFirst: async ({ where }: any) => { + assert.equal(where.number, metadata.phone_number_id); + assert.ok(integrations.includes(where.integration)); + return { name: where.integration }; + }, + }, + }; + const monitor = { + waInstances: Object.fromEntries( + integrations.map((name) => [ + name, + { + connectToWhatsapp: async () => calls.push(name), + }, + ]), + ), + }; + await new MetaController(repository, monitor).receiveWebhook(envelope({ metadata, message_echoes: [message()] })); + await new EvoHubController(repository, monitor, {}).receiveWebhook( + envelope({ metadata, message_echoes: [message()] }), + ); + assert.deepEqual(calls, integrations); +}); + +test('the controller rejects unavailable channels and failed processing', async () => { + const h = await harness(); + h.service.connectToWhatsapp = async () => { + throw new Error('processing failed'); + }; + const payload = envelope({ metadata, message_echoes: [message()] }); + await assert.rejects(() => h.controller.receiveWebhook(payload), /processing failed/); + delete h.controller.waMonitor.waInstances['tenant-a']; + await assert.rejects(() => h.controller.receiveWebhook(payload), /indisponível/); +}); + +test('contact updates stay within the originating tenant', async () => { + const h = await harness(); + h.service.prismaRepository.contact.findFirst = async () => ({ remoteJid: customerJid, pushName: 'Customer' }); + await h.service.connectToWhatsapp(envelope({ metadata, message_echoes: [message()] })); + assert.equal(h.updates.length, 1); + assert.deepEqual(h.updates[0].where, { instanceId: 'tenant-a', remoteJid: customerJid }); +}); + +test('persistence failure rejects processing so the upstream can retry', async () => { + const h = await harness(); + h.service.prismaRepository.message.create = async () => { + throw new Error('database unavailable'); + }; + await assert.rejects( + () => h.service.connectToWhatsapp(envelope({ metadata, message_echoes: [message()] })), + /database unavailable/, + ); +}); + +test('media and stickers sent from the app are emitted and persisted without S3', async () => { + for (const type of ['audio', 'image', 'video', 'document', 'sticker']) { + const h = await harness(); + const media = { ...message(`wamid.${type}`), type, [type]: { id: 'synthetic-media', mime_type: `${type}/test` } }; + await h.service.connectToWhatsapp(envelope({ metadata, message_echoes: [media] })); + assert.equal(h.emitted.length, 1, type); + assert.equal(h.emitted[0].msg.key.fromMe, true, type); + assert.equal(h.stored.length, 1, `${type} must be persisted`); + } +}); + +test('media already persisted by the S3 path is not inserted again', async () => { + const h = await harness(); + const mediaRows: any[] = []; + h.service.configService.get = (key: string) => ({ ENABLE: key === 'S3', ENABLED: false, SAVE_DATA: {} }); + h.service.hasValidMediaContent = () => true; + h.service.fetchMediaFromGraph = async (id: string) => { + assert.equal(id, 'synthetic-image'); + return { + result: { data: { mime_type: 'image/png' }, headers: {} }, + buffer: { data: Buffer.from('synthetic image bytes') }, + }; + }; + h.service.prismaRepository.media = { create: async ({ data }: any) => mediaRows.push(data) }; + const image = { ...message('wamid.image'), type: 'image', image: { id: 'synthetic-image' } }; + await h.service.connectToWhatsapp(envelope({ metadata, message_echoes: [image] })); + assert.equal(h.stored.length, 1); + assert.equal(mediaRows.length, 1); + assert.equal(h.emitted.length, 1); + assert.deepEqual(h.errors, []); +}); + +test('concurrent echoes never reuse another customer recipient', async () => { + const h = await harness(); + h.service.findSettings = async () => { + await new Promise((done) => setTimeout(done, 5)); + return {}; + }; + await Promise.all([ + h.service.connectToWhatsapp(envelope({ metadata, message_echoes: [message('wamid.a')] })), + h.service.connectToWhatsapp(envelope({ metadata, message_echoes: [message('wamid.b', '16505559999')] })), + ]); + assert.deepEqual( + h.emitted.map((item) => [item.msg.key.id, item.remoteJid]), + [ + ['wamid.a', customerJid], + ['wamid.b', '16505559999@s.whatsapp.net'], + ], + ); +}); + +test('co-located inbound and echo arrays preserve distinct message directions', async () => { + const h = await harness(); + await h.service.connectToWhatsapp( + envelope({ metadata, messages: [{ ...message('wamid.in'), from: customer }], message_echoes: [message()] }), + ); + assert.deepEqual( + h.emitted.map((item) => item.msg.key.fromMe), + [false, true], + ); +}); + +test('unknown and deleted status items do not discard later updates and saved direction is retained', async () => { + const h = await harness(); + h.service.prismaRepository.message.findFirst = async ({ where }: any) => + ['wamid.deleted', 'wamid.known'].includes(where.key.equals) + ? { id: where.key.equals, key: { fromMe: false, remoteJid: customerJid } } + : null; + await h.service.connectToWhatsapp( + envelope( + { + metadata, + statuses: [ + { id: 'wamid.unknown', recipient_id: customer, status: 'delivered' }, + { id: 'wamid.deleted', recipient_id: customer, message: null }, + { id: 'wamid.known', recipient_id: customer, status: 'read' }, + ], + }, + 'messages', + ), + ); + assert.deepEqual( + h.updates.map((item) => item.status), + ['DELETED', 'READ'], + ); + assert.equal(h.updates[1].fromMe, false); + assert.equal(h.updates[1].remoteJid, customerJid); +}); + +for (const pause of [false, true]) { + test(`the native bot ${pause ? 'pauses an active session' : 'ignores the app echo'} without generating a reply`, async () => { + const h = await harness(); + const { BaseChatbotController } = await productionClasses; + const bot = Object.create(BaseChatbotController.prototype); + const responses: any[] = []; + const pauses: any[] = []; + Object.assign(bot, { + integrationEnabled: true, + integrationName: 'N8n', + settingsRepository: { findFirst: async () => ({ listeningFromMe: false, stopBotFromMe: pause }) }, + getSession: async () => (pause ? { id: 'session-a', status: 'opened' } : null), + checkIgnoreJids: () => false, + findBotTrigger: async () => ({ listeningFromMe: false, stopBotFromMe: pause }), + processBot: async (...args: any[]) => responses.push(args), + prismaRepository: { integrationSession: { update: async (query: any) => pauses.push(query) } }, + waMonitor: { waInstances: {} }, + logger: h.service.logger, + }); + (globalThis as any).__metaTestEmit = (data: any) => bot.emit(data); + await h.service.connectToWhatsapp(envelope({ metadata, message_echoes: [message()] })); + assert.equal(responses.length, 0); + assert.equal(pauses.length, pause ? 1 : 0); + if (pause) assert.deepEqual(pauses[0], { where: { id: 'session-a' }, data: { status: 'paused' } }); + assert.deepEqual(h.errors, []); + }); +} From 01ef5965e1d0faf828e9f67ac7efe36dbcd9e2b0 Mon Sep 17 00:00:00 2001 From: christianini-debug <338734608+christianini-debug@users.noreply.github.com> Date: Tue, 6 Oct 2026 16:43:38 -0300 Subject: [PATCH 2/2] fix(meta)!: address webhook review findings Validate Meta POST signatures before dispatch, preserve business ID senders and contact names, and await status delivery. BREAKING CHANGE: Meta POST webhooks require WA_BUSINESS_APP_SECRET. Missing configuration returns 503; invalid signatures return 401. --- .env.example | 4 + src/api/guards/meta-webhook.guard.ts | 26 ++ .../integrations/channel/meta/meta.router.ts | 3 +- .../channel/meta/whatsapp.business.service.ts | 22 +- src/config/env.config.ts | 2 + src/main.ts | 3 +- test/meta-coexistence.test.ts | 255 +++++++++++++++++- 7 files changed, 301 insertions(+), 14 deletions(-) create mode 100644 src/api/guards/meta-webhook.guard.ts diff --git a/.env.example b/.env.example index fa7d9e7bd3..c63c57acf0 100644 --- a/.env.example +++ b/.env.example @@ -253,6 +253,10 @@ KAFKA_SSL_CERT= # WhatsApp Business API - Environment variables # Token used to validate the webhook on the Facebook APP WA_BUSINESS_TOKEN_WEBHOOK=evolution +# Required for Meta POST webhooks: App Secret from the Meta app that signs the notifications. +# This is different from the verification token above and the Cloud API access token. +# Without this secret, /webhook/meta returns 503; invalid or missing signatures return 401. +WA_BUSINESS_APP_SECRET= WA_BUSINESS_URL=https://graph.facebook.com WA_BUSINESS_VERSION=v20.0 WA_BUSINESS_LANGUAGE=en_US diff --git a/src/api/guards/meta-webhook.guard.ts b/src/api/guards/meta-webhook.guard.ts new file mode 100644 index 0000000000..183ff35993 --- /dev/null +++ b/src/api/guards/meta-webhook.guard.ts @@ -0,0 +1,26 @@ +import { ConfigService, WaBusiness } from '@config/env.config'; +import { createHmac, timingSafeEqual } from 'crypto'; +import { Request, RequestHandler } from 'express'; + +export function metaWebhookGuard(configService: ConfigService): RequestHandler { + return (req: Request & { rawBody?: Buffer }, res, next) => { + const secret = configService.get('WA_BUSINESS').APP_SECRET; + if (!secret) { + return res.status(503).json({ error: 'WA_BUSINESS_APP_SECRET não configurado' }); + } + + const signature = req.headers['x-hub-signature-256']; + if (typeof signature !== 'string' || !/^sha256=[a-f\d]{64}$/i.test(signature) || !Buffer.isBuffer(req.rawBody)) { + return res.status(401).json({ error: 'Assinatura do webhook inválida' }); + } + + // Authenticate the original bytes captured by the JSON parser, never a reserialized body. + const expected = createHmac('sha256', secret).update(req.rawBody).digest(); + const supplied = Buffer.from(signature.slice('sha256='.length), 'hex'); + if (!timingSafeEqual(expected, supplied)) { + return res.status(401).json({ error: 'Assinatura do webhook inválida' }); + } + + return next(); + }; +} diff --git a/src/api/integrations/channel/meta/meta.router.ts b/src/api/integrations/channel/meta/meta.router.ts index b0fc43ce4d..6a55cb81c6 100644 --- a/src/api/integrations/channel/meta/meta.router.ts +++ b/src/api/integrations/channel/meta/meta.router.ts @@ -1,4 +1,5 @@ import { RouterBroker } from '@api/abstract/abstract.router'; +import { metaWebhookGuard } from '@api/guards/meta-webhook.guard'; import { metaController } from '@api/server.module'; import { ConfigService, WaBusiness } from '@config/env.config'; import { Router } from 'express'; @@ -12,7 +13,7 @@ export class MetaRouter extends RouterBroker { res.send(req.query['hub.challenge']); else res.send('Error, wrong validation token'); }) - .post(this.routerPath('webhook/meta', false), async (req, res) => { + .post(this.routerPath('webhook/meta', false), metaWebhookGuard(configService), async (req, res) => { const { body } = req; const response = await metaController.receiveWebhook(body); diff --git a/src/api/integrations/channel/meta/whatsapp.business.service.ts b/src/api/integrations/channel/meta/whatsapp.business.service.ts index ea3775ecc2..ee8a851a1d 100644 --- a/src/api/integrations/channel/meta/whatsapp.business.service.ts +++ b/src/api/integrations/channel/meta/whatsapp.business.service.ts @@ -171,7 +171,8 @@ export class BusinessStartupService extends ChannelStartupService { const from = this.normalizePhoneNumber(message?.from); const displayPhone = this.normalizePhoneNumber(received?.metadata?.display_phone_number); - return !!from && from === displayPhone; + const phoneNumberId = this.normalizePhoneNumber(received?.metadata?.phone_number_id); + return !!from && (from === displayPhone || from === phoneNumberId); } private isCloudApiStatusFromMe(item: any, received: any) { @@ -449,11 +450,6 @@ export class BusinessStartupService extends ChannelStartupService { let messageRaw: any; let pushName: any; let persistedMedia = false; - const incomingContact = received?.contacts?.[0]; - - if (incomingContact) { - pushName = incomingContact?.profile?.name ?? incomingContact?.name ?? incomingContact?.wa_id ?? undefined; - } if (received.messages) { const message = received.messages[0]; @@ -461,6 +457,14 @@ export class BusinessStartupService extends ChannelStartupService { if (!remoteId) return; const remoteJid = createJid(remoteId); + const incomingContact = Array.isArray(received.contacts) + ? received.contacts.find((item: any) => typeof item.wa_id === 'string' && createJid(item.wa_id) === remoteJid) + : undefined; + + if (incomingContact) { + pushName = incomingContact.profile?.name ?? incomingContact.name ?? incomingContact.wa_id ?? undefined; + } + const contact = await this.prismaRepository.contact.findFirst({ where: { instanceId: this.instanceId, remoteJid }, }); @@ -839,7 +843,7 @@ export class BusinessStartupService extends ChannelStartupService { } if (item.message === null && item.status === undefined) { - this.sendDataWebhook(Events.MESSAGES_DELETE, key); + await this.sendDataWebhook(Events.MESSAGES_DELETE, key); const message: any = { messageId: findMessage.id, @@ -856,7 +860,7 @@ export class BusinessStartupService extends ChannelStartupService { }); if (this.configService.get('CHATWOOT').ENABLED && this.localChatwoot?.enabled) { - this.chatwootService.eventWhatsapp( + await this.chatwootService.eventWhatsapp( Events.MESSAGES_DELETE, { instanceName: this.instance.name, instanceId: this.instanceId }, { key: key }, @@ -876,7 +880,7 @@ export class BusinessStartupService extends ChannelStartupService { instanceId: this.instanceId, }; - this.sendDataWebhook(Events.MESSAGES_UPDATE, message); + await this.sendDataWebhook(Events.MESSAGES_UPDATE, message); await this.prismaRepository.messageUpdate.create({ data: message, diff --git a/src/config/env.config.ts b/src/config/env.config.ts index 1b17a83b53..0f9532605a 100644 --- a/src/config/env.config.ts +++ b/src/config/env.config.ts @@ -192,6 +192,7 @@ export type Websocket = { export type WaBusiness = { TOKEN_WEBHOOK: string; + APP_SECRET: string; URL: string; VERSION: string; LANGUAGE: string; @@ -748,6 +749,7 @@ export class ConfigService { }, WA_BUSINESS: { TOKEN_WEBHOOK: process.env.WA_BUSINESS_TOKEN_WEBHOOK || 'evolution', + APP_SECRET: process.env.WA_BUSINESS_APP_SECRET || '', URL: process.env.WA_BUSINESS_URL || 'https://graph.facebook.com', VERSION: process.env.WA_BUSINESS_VERSION || 'v18.0', LANGUAGE: process.env.WA_BUSINESS_LANGUAGE || 'en', diff --git a/src/main.ts b/src/main.ts index f9884ba34b..3fa7174e56 100644 --- a/src/main.ts +++ b/src/main.ts @@ -77,8 +77,7 @@ async function bootstrap() { json({ limit: '136mb', verify: (req: any, _res, buf) => { - // Captura o RAW body para validação de HMAC (webhook EvoHub X-Hub-Signature-256). - // express.json() re-serializa o body; o HMAC do hub assina os bytes crus. + // Preserve the original bytes for Meta and EvoHub X-Hub-Signature-256 validation. req.rawBody = buf; }, }), diff --git a/test/meta-coexistence.test.ts b/test/meta-coexistence.test.ts index 24bd7cc8e3..589ef53df5 100644 --- a/test/meta-coexistence.test.ts +++ b/test/meta-coexistence.test.ts @@ -1,9 +1,12 @@ import assert from 'node:assert/strict'; +import { createHmac } from 'node:crypto'; +import { once } from 'node:events'; import { createRequire } from 'node:module'; import { resolve } from 'node:path'; -import { test } from 'node:test'; +import { test, TestContext } from 'node:test'; import { build } from 'esbuild'; +import express from 'express'; // Bundle the production controller and service, replacing only infrastructure boundaries. const root = resolve(__dirname, '..'); @@ -11,6 +14,7 @@ const productionClasses = build({ stdin: { contents: `export { BusinessStartupService } from ${JSON.stringify(process.env.COEX_SERVICE_SOURCE ?? './src/api/integrations/channel/meta/whatsapp.business.service')}; export { MetaController } from ${JSON.stringify(process.env.COEX_CONTROLLER_SOURCE ?? './src/api/integrations/channel/meta/meta.controller')}; +export { MetaRouter } from './src/api/integrations/channel/meta/meta.router'; export { EvoHubController } from './src/api/integrations/channel/evohub/evohub.controller'; export { BaseChatbotController } from './src/api/integrations/chatbot/base-chatbot.controller';`, resolveDir: root, @@ -32,7 +36,8 @@ export { BaseChatbotController } from './src/api/integrations/chatbot/base-chatb }`, '@api/server.module': `export const chatbotController = { emit: async (data) => globalThis.__metaTestEmit(data), - };`, + }; + export const metaController = { receiveWebhook: (data) => globalThis.__metaTestReceive(data) };`, '@utils/sendTelemetry': 'export const sendTelemetry = () => {};', '@api/integrations/storage/s3/libs/minio.server': 'export const uploadFile = () => {}; export const getObjectUrl = () => {};', @@ -130,6 +135,117 @@ async function harness(instanceId = 'tenant-a') { return { service, controller, events, emitted, stored, updates, errors }; } +const appSecret = 'synthetic-test-app-secret'; +const signBody = (body: string, secret = appSecret) => + `sha256=${createHmac('sha256', secret).update(body).digest('hex')}`; + +async function routerHarness(t: TestContext, secret = appSecret, captureRaw = true) { + const h = await harness(); + const { MetaRouter } = await productionClasses; + let dispatches = 0; + (globalThis as any).__metaTestReceive = (body: any) => { + dispatches++; + return h.controller.receiveWebhook(body); + }; + const app = express(); + app.use( + express.json({ + verify: (req: any, _res, buffer) => { + if (captureRaw) req.rawBody = buffer; + }, + }), + ); + app.use(new MetaRouter({ get: () => ({ APP_SECRET: secret, TOKEN_WEBHOOK: 'test-verify-token' }) }).router); + const server = app.listen(0, '127.0.0.1'); + await once(server, 'listening'); + t.after( + () => + new Promise((resolve, reject) => { + server.closeAllConnections(); + server.close((error) => (error ? reject(error) : resolve())); + }), + ); + const address = server.address(); + assert.ok(address && typeof address === 'object'); + const url = `http://127.0.0.1:${address.port}/webhook/meta`; + const post = (body: string, signature?: string) => + fetch(url, { + method: 'POST', + body, + headers: { + 'Content-Type': 'application/json', + ...(signature === undefined ? {} : { 'X-Hub-Signature-256': signature }), + }, + }); + return { ...h, url, post, dispatches: () => dispatches }; +} + +test('the Meta POST route rejects missing, invalid and tampered signatures before dispatch', async (t) => { + const h = await routerHarness(t); + const body = JSON.stringify(envelope({ metadata, message_echoes: [message()] }), null, 2); + for (const signature of [ + undefined, + '', + 'sha1=' + 'a'.repeat(40), + 'sha256=invalid', + 'sha256=' + 'f'.repeat(64), + signBody(body, 'wrong-secret'), + ]) { + const response = await h.post(body, signature); + assert.equal(response.status, 401); + await response.text(); + } + for (const tampered of [body.replace(customer, '16505550000'), JSON.stringify(JSON.parse(body))]) { + const response = await h.post(tampered, signBody(body)); + assert.equal(response.status, 401); + await response.text(); + } + assert.equal(h.dispatches(), 0); + assert.equal(h.stored.length, 0); + assert.equal(h.events.length, 0); +}); + +test('a valid signature over the original UTF-8 bytes reaches the Meta controller and service', async (t) => { + const h = await routerHarness(t); + const body = JSON.stringify( + envelope({ metadata, message_echoes: [{ ...message(), text: { body: 'Olá 👋' } }] }), + null, + 2, + ); + const response = await h.post(body, signBody(body)); + assert.equal(response.status, 200); + assert.deepEqual(await response.json(), { status: 'success' }); + assert.equal(h.dispatches(), 1); + assert.equal(h.stored.length, 1); + assert.deepEqual(h.emitted[0].msg.key, { id: 'wamid.phone-1', remoteJid: customerJid, fromMe: true }); +}); + +test('the Meta POST route fails closed when its App Secret is missing', async (t) => { + const h = await routerHarness(t, ''); + const body = JSON.stringify(envelope({ metadata, message_echoes: [message()] })); + const response = await h.post(body, signBody(body)); + assert.equal(response.status, 503); + await response.text(); + assert.equal(h.dispatches(), 0); +}); + +test('the Meta POST route rejects a signature when raw request bytes are unavailable', async (t) => { + const h = await routerHarness(t, appSecret, false); + const body = JSON.stringify(envelope({ metadata, message_echoes: [message()] })); + const response = await h.post(body, signBody(body)); + assert.equal(response.status, 401); + await response.text(); + assert.equal(h.dispatches(), 0); +}); + +test('the Meta GET verification handshake still uses the verification token', async (t) => { + const h = await routerHarness(t, ''); + const response = await fetch(`${h.url}?hub.verify_token=test-verify-token&hub.challenge=12345`); + assert.equal(response.status, 200); + assert.equal(await response.text(), '12345'); + assert.equal(h.dispatches(), 0); +}); + test('a Business app echo reaches the agent and webhook as an outgoing message in the customer conversation', async () => { const h = await harness(); await h.service.connectToWhatsapp(envelope({ metadata, message_echoes: [message()] })); @@ -185,6 +301,78 @@ test('classic inbound messages still reach the customer conversation with fromMe assert.deepEqual(h.errors, []); }); +test('business messages identified by phone_number_id remain outgoing outside echo arrays', async () => { + const h = await harness(); + const outgoing = { ...message('wamid.business-id'), from: metadata.phone_number_id }; + await h.controller.receiveWebhook(envelope({ metadata, messages: [outgoing] }, 'messages')); + assert.deepEqual(h.stored[0].key, { id: outgoing.id, remoteJid: customerJid, fromMe: true }); + assert.equal(h.emitted[0].remoteJid, customerJid); + assert.equal(h.emitted[0].msg.key.fromMe, true); +}); + +for (const isEcho of [false, true]) { + test(`${isEcho ? 'echo' : 'inbound'} batches match each customer to its own contact profile`, async () => { + const h = await harness(); + const secondCustomer = '16505559999'; + const customers = [customer, secondCustomer]; + h.service.prismaRepository.contact.findFirst = async ({ where }: any) => ({ + remoteJid: where.remoteJid, + pushName: 'Previously saved name', + }); + const messages = customers.map((number, index) => ({ + ...message(`wamid.profile-${index}`, isEcho ? number : metadata.display_phone_number), + from: isEcho ? metadata.display_phone_number : number, + })); + await h.controller.receiveWebhook( + envelope({ + metadata, + contacts: [ + { wa_id: secondCustomer, profile: { name: 'Bob' } }, + { wa_id: metadata.display_phone_number, profile: { name: 'Business' } }, + { wa_id: customer, profile: { name: 'Alice' } }, + ], + [isEcho ? 'message_echoes' : 'messages']: messages, + }), + ); + assert.deepEqual( + h.stored.map((item) => [item.key.remoteJid, item.pushName]), + [ + [customerJid, 'Alice'], + [secondCustomer + '@s.whatsapp.net', 'Bob'], + ], + ); + assert.deepEqual( + h.emitted.map((item) => item.pushName), + ['Alice', 'Bob'], + ); + assert.deepEqual( + h.updates.map((item) => [item.where.remoteJid, item.data.pushName]), + [ + [customerJid, 'Alice'], + [secondCustomer + '@s.whatsapp.net', 'Bob'], + ], + ); + assert.deepEqual( + h.events.filter((item) => item.event === 'contacts.update').map((item) => item.data.pushName), + ['Alice', 'Bob'], + ); + }); +} + +test('unmatched contact profiles never replace the saved customer name', async () => { + const h = await harness(); + h.service.prismaRepository.contact.findFirst = async () => ({ remoteJid: customerJid, pushName: 'Saved customer' }); + await h.controller.receiveWebhook( + envelope({ + metadata, + contacts: [{ wa_id: '16505559999', profile: { name: 'Unrelated customer' } }], + message_echoes: [message()], + }), + ); + assert.equal(h.stored[0].pushName, 'Saved customer'); + assert.equal(h.updates[0].data.pushName, 'Saved customer'); +}); + test('empty messages and message_echoes do not hide the alternate smb_message_echoes array', async () => { const h = await harness(); await h.service.connectToWhatsapp( @@ -388,6 +576,69 @@ test('unknown and deleted status items do not discard later updates and saved di assert.equal(h.updates[1].remoteJid, customerJid); }); +for (const dispatch of ['update webhook', 'delete webhook', 'delete Chatwoot']) { + async function statusHarness() { + const h = await harness(); + h.service.prismaRepository.message.findFirst = async () => ({ + id: 'db-status', + key: { fromMe: true, remoteJid: customerJid }, + }); + h.service.configService.get = (key: string) => ({ ENABLED: key === 'CHATWOOT' }); + h.service.localChatwoot.enabled = true; + h.service.chatwootService = { eventWhatsapp: async () => {} }; + const status = { + id: 'wamid.status', + recipient_id: customer, + ...(dispatch.startsWith('update') ? { status: 'read' } : { message: null }), + }; + return { ...h, payload: envelope({ metadata, statuses: [status] }, 'messages') }; + } + + test(`the controller waits for ${dispatch} delivery before acknowledging a status`, async () => { + const h = await statusHarness(); + let release: () => void; + let started: () => void; + const delivery = new Promise((resolve) => { + release = resolve; + }); + const entered = new Promise((resolve) => { + started = resolve; + }); + const handler = () => { + started(); + return delivery; + }; + if (dispatch.endsWith('Chatwoot')) h.service.chatwootService.eventWhatsapp = handler; + else h.service.sendDataWebhook = handler; + let acknowledged = false; + const processing = h.controller.receiveWebhook(h.payload).then(() => { + acknowledged = true; + }); + try { + await entered; + await new Promise((resolve) => setImmediate(resolve)); + assert.equal(acknowledged, false); + } finally { + release(); + await processing; + } + assert.equal(acknowledged, true); + }); + + test(`a rejected ${dispatch} delivery rejects the webhook request`, async () => { + const h = await statusHarness(); + const handler = () => { + const failure = Promise.reject(new Error('status delivery failed')); + // Keep the negative control safe even when the old handler discards this promise. + void failure.catch(() => {}); + return failure; + }; + if (dispatch.endsWith('Chatwoot')) h.service.chatwootService.eventWhatsapp = handler; + else h.service.sendDataWebhook = handler; + await assert.rejects(() => h.controller.receiveWebhook(h.payload), /status delivery failed/); + }); +} + for (const pause of [false, true]) { test(`the native bot ${pause ? 'pauses an active session' : 'ignores the app echo'} without generating a reply`, async () => { const h = await harness();