From ef939f4357bccb38fb910a5bcfe781af315b4dec Mon Sep 17 00:00:00 2001 From: devlikepro Date: Sun, 4 May 2025 13:13:05 +0700 Subject: [PATCH] [core] Add "message.ack" event for groups and status fix #495 fix #900 --- src/core/engines/noweb/session.noweb.core.ts | 18 +- src/core/engines/webjs/ack.webjs.ts | 163 +++++++++++++++++++ src/core/engines/webjs/session.webjs.core.ts | 83 +++++++++- src/core/utils/ids.ts | 6 + yarn.lock | 4 +- 5 files changed, 258 insertions(+), 16 deletions(-) create mode 100644 src/core/engines/webjs/ack.webjs.ts diff --git a/src/core/engines/noweb/session.noweb.core.ts b/src/core/engines/noweb/session.noweb.core.ts index 3c694328..0e82d991 100644 --- a/src/core/engines/noweb/session.noweb.core.ts +++ b/src/core/engines/noweb/session.noweb.core.ts @@ -1952,7 +1952,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { if (!message.message.reactionMessage) return null; const id = buildMessageId(message.key); - const fromToParticipant = getFromToParticipant(message); + const fromToParticipant = getFromToParticipant(message.key); const reactionMessage = message.message.reactionMessage; const messageId = buildMessageId(reactionMessage.key); const source = this.getMessageSource(message.key.id); @@ -2021,7 +2021,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { } protected toWAMessage(message): Promise { - const fromToParticipant = getFromToParticipant(message); + const fromToParticipant = getFromToParticipant(message.key); const id = buildMessageId(message.key); const body = this.extractBody(message.message); const replyTo = this.extractReplyTo(message.message); @@ -2112,7 +2112,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { protected convertMessageUpdateToMessageAck(event): WAMessageAckBody { const message = event; - const fromToParticipant = getFromToParticipant(message); + const fromToParticipant = getFromToParticipant(message.key); const id = buildMessageId(message.key); const ack = message.update.status - 1; const body: WAMessageAckBody = { @@ -2128,7 +2128,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { } protected convertMessageReceiptUpdateToMessageAck(event): WAMessageAckBody { - const fromToParticipant = getFromToParticipant(event); + const fromToParticipant = getFromToParticipant(event.key); const id = buildMessageId(event.key); const receipt = event.receipt; @@ -2445,15 +2445,15 @@ function isAckUpdateMessageEvent(event) { return event?.update.status != null; } -function getFromToParticipant(message) { - const isGroupMessage = Boolean(message.key.participant); +export function getFromToParticipant(key) { + const isGroupMessage = Boolean(key.participant); let participant: string; let to: string; if (isGroupMessage) { - participant = message.key.participant; - to = message.key.remoteJid; + participant = key.participant; + to = key.remoteJid; } - const from = message.key.remoteJid; + const from = key.remoteJid; return { from: from, to: to, diff --git a/src/core/engines/webjs/ack.webjs.ts b/src/core/engines/webjs/ack.webjs.ts new file mode 100644 index 00000000..7b0cfd53 --- /dev/null +++ b/src/core/engines/webjs/ack.webjs.ts @@ -0,0 +1,163 @@ +import { + areJidsSameUser, + BinaryNode, + getBinaryNodeChildren, + getStatusFromReceiptType, + isJidGroup, + isJidStatusBroadcast, + jidEncode, + jidNormalizedUser, + proto, +} from '@adiwajshing/baileys'; + +export interface ReceiptEvent { + key: proto.IMessageKey; + participant?: string; + messageIds: string[]; + status: number; + _node: any; +} + +interface Me { + id: string; + lid?: string; +} + +function jid(field: any) { + if (!field) { + return field; + } + const data = field['$1']; + if (!data) { + return data; + } + + let server = data.server; + if (!server) { + server = + data.domainType === 0 || data.domainType === 128 + ? 's.whatsapp.net' + : 'lid'; + } + return jidEncode(data.user, server, data.device); +} + +export function TagReceiptNodeToReceiptEvent( + node: BinaryNode, + me: Me, +): ReceiptEvent[] { + const { attrs, content } = node; + const status = getStatusFromReceiptType(attrs.type); + if (status == null) { + return null; + } + + const from = jidNormalizedUser(jid(attrs.from)); + const participant = jidNormalizedUser(jid(attrs.participant)); + const recipient = jidNormalizedUser(jid(attrs.recipient)); + + const isLid = from.includes('lid'); + const isNodeFromMe = areJidsSameUser( + participant || from, + isLid ? me?.lid : me?.id, + ); + const remoteJid = !isNodeFromMe || isJidGroup(from) ? from : recipient; + const fromMe = !recipient || (attrs.type === 'retry' && isNodeFromMe); + + // basically, we only want to know when a message from us has been delivered to/read by the other person + // or another device of ours has read some messages + if (status < proto.WebMessageInfo.Status.SERVER_ACK && isNodeFromMe) { + return []; + } + + const key: proto.IMessageKey = { + remoteJid: remoteJid, + id: '', + fromMe: fromMe, + }; + + const ids = [attrs.id]; + if (Array.isArray(content)) { + const items = getBinaryNodeChildren(content[0], 'item'); + ids.push(...items.map((i) => i.attrs.id)); + } + + if (isJidGroup(remoteJid) || isJidStatusBroadcast(remoteJid)) { + if (participant) { + key.participant = fromMe ? (isLid ? me.lid : me.id) : recipient; + const eventParticipant = fromMe ? participant : isLid ? me.lid : me.id; + return [ + { + key: key, + messageIds: ids, + status: status as any, + participant: eventParticipant, + _node: node, + }, + ]; + } else { + // Handle grouped receipts + return handleGroupedReceipts(node, key, status, fromMe, isLid, me); + } + } + + return [ + { + key: key, + messageIds: ids, + status: status as any, + _node: node, + }, + ]; +} + +function handleGroupedReceipts( + node: BinaryNode, + key: proto.IMessageKey, + status: number, + fromMe: boolean, + isLid: boolean, + me: Me, +): ReceiptEvent[] | null { + const { content } = node; + if (!Array.isArray(content)) { + return []; + } + const participantsTags = content.filter((c) => c.tag === 'participants'); + if (participantsTags.length === 0) { + return null; + } + + const receiptEvents: ReceiptEvent[] = []; + + for (const participants of participantsTags) { + const participantKey = participants.attrs?.key; + if (!participantKey) continue; + + const users = getBinaryNodeChildren(participants, 'user'); + for (const user of users) { + const userAttrs = user.attrs; + if (!userAttrs) continue; + + const userJid = jidNormalizedUser(jid(userAttrs.jid)); + if (!userJid) continue; + + key.participant = fromMe ? (isLid ? me.lid : me.id) : userJid; + const eventParticipant = fromMe ? userJid : isLid ? me.lid : me.id; + const receiptEvent: ReceiptEvent = { + key: { + ...key, + id: participantKey, + }, + messageIds: [participantKey], + status: status as any, + participant: eventParticipant, + _node: node, + }; + + receiptEvents.push(receiptEvent); + } + } + + return receiptEvents; +} diff --git a/src/core/engines/webjs/session.webjs.core.ts b/src/core/engines/webjs/session.webjs.core.ts index 31946e2c..aa85ccc2 100644 --- a/src/core/engines/webjs/session.webjs.core.ts +++ b/src/core/engines/webjs/session.webjs.core.ts @@ -1,8 +1,17 @@ +import { isJidGroup, isJidStatusBroadcast } from '@adiwajshing/baileys'; import { UnprocessableEntityException } from '@nestjs/common'; import { getChannelInviteLink, WhatsappSession, } from '@waha/core/abc/session.abc'; +import { + getFromToParticipant, + toCusFormat, +} from '@waha/core/engines/noweb/session.noweb.core'; +import { + ReceiptEvent, + TagReceiptNodeToReceiptEvent, +} from '@waha/core/engines/webjs/ack.webjs'; import { ToGroupV2JoinEvent, ToGroupV2LeaveEvent, @@ -21,6 +30,10 @@ import { } from '@waha/core/exceptions'; import { IMediaEngineProcessor } from '@waha/core/media/IMediaEngineProcessor'; import { QR } from '@waha/core/QR'; +import { + parseMessageIdSerialized, + SerializeMessageKey, +} from '@waha/core/utils/ids'; import { splitAt } from '@waha/helpers'; import { PairingCodeResponse } from '@waha/structures/auth.dto'; import { @@ -89,6 +102,7 @@ import { MeInfo } from '@waha/structures/sessions.dto'; import { StatusRequest, TextStatus } from '@waha/structures/status.dto'; import { EnginePayload, + WAMessageAckBody, WAMessageRevokedBody, } from '@waha/structures/webhooks.dto'; import { PaginatorInMemory } from '@waha/utils/Paginator'; @@ -96,7 +110,17 @@ import { sleep, waitUntil } from '@waha/utils/promiseTimeout'; import { SingleDelayedJobRunner } from '@waha/utils/SingleDelayedJobRunner'; import * as lodash from 'lodash'; import { ProtocolError } from 'puppeteer'; -import { filter, fromEvent, merge, mergeMap, Observable } from 'rxjs'; +import { + debounceTime, + distinct, + filter, + fromEvent, + groupBy, + interval, + merge, + mergeMap, + Observable, +} from 'rxjs'; import { map } from 'rxjs/operators'; import { Call, @@ -110,6 +134,7 @@ import { Label as WEBJSLabel, Location, Message, + MessageAck, MessageMedia, Reaction, WAState, @@ -1275,18 +1300,37 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { ); this.events2.get(WAHAEvents.MESSAGE_REACTION).switch(messagesReaction$); - const messageAck$ = fromEvent( + const messageAckWEBJS$ = fromEvent( this.whatsapp, Events.MESSAGE_ACK, (message, ack) => { return { message, ack }; }, ); - const messagesAck$ = messageAck$.pipe( + const messagesAckDM$ = messageAckWEBJS$.pipe( map((event) => event.message), - map(this.toWAMessage.bind(this)), + map(this.toWAMessage.bind(this)), + filter((ack) => !isJidGroup(ack.to) && !isJidStatusBroadcast(ack.to)), ); - this.events2.get(WAHAEvents.MESSAGE_ACK).switch(messagesAck$); + const tagReceiptNode$ = fromEvent(this.whatsapp, Events.TAG_RECEIPT); + const messageAckGroups$ = tagReceiptNode$.pipe( + mergeMap((node) => + TagReceiptNodeToReceiptEvent(node as any, this.getSessionMeInfo()), + ), + mergeMap(this.TagReceiptToMessageAck.bind(this)), + filter((ack) => isJidGroup(ack.to) || isJidStatusBroadcast(ack.to)), + ); + + const messageAckAll$ = merge(messagesAckDM$, messageAckGroups$); + + const messageAck$ = messageAckAll$.pipe( + // emit only if we haven’t seen this key since the last flush + distinct( + (msg: WAMessageAckBody) => `${msg.id}-${msg.ack}-${msg.participant}`, + interval(60_000), + ), + ); + this.events2.get(WAHAEvents.MESSAGE_ACK).switch(messageAck$); // // Others @@ -1418,15 +1462,44 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { }; } + protected TagReceiptToMessageAck(receipt: ReceiptEvent): WAMessageAckBody[] { + const ids = receipt.messageIds; + const acks = []; + for (const id_ of ids) { + const messageKey = { + fromMe: receipt.key.fromMe, + remoteJid: toCusFormat(receipt.key.remoteJid), + participant: toCusFormat(receipt.key.participant), + id: id_, + }; + const fromToParticipant = getFromToParticipant(messageKey); + const id = SerializeMessageKey(messageKey); + const ack = receipt.status - 1; + acks.push({ + id: id, + from: fromToParticipant.from, + to: fromToParticipant.to, + participant: toCusFormat(receipt.participant), + fromMe: !receipt.key.fromMe, // reverted, it's right + ack: ack, + ackName: WAMessageAck[ack] || ACK_UNKNOWN, + _data: receipt._node, + }); + } + return acks; + } + protected toWAMessage(message: Message): WAMessage { const replyTo = this.extractReplyTo(message); const source = this.getMessageSource(message.id.id); + const key = parseMessageIdSerialized(message.id._serialized); // @ts-ignore return { id: message.id._serialized, timestamp: message.timestamp, from: message.from, fromMe: message.fromMe, + participant: toCusFormat(key.participant), source: source, to: message.to, body: message.body, diff --git a/src/core/utils/ids.ts b/src/core/utils/ids.ts index a794cce9..93db5c80 100644 --- a/src/core/utils/ids.ts +++ b/src/core/utils/ids.ts @@ -32,3 +32,9 @@ export function parseMessageIdSerialized( participant: participant, }; } + +export function SerializeMessageKey(key: WAMessageKey) { + const { fromMe, id, remoteJid, participant } = key; + const participantStr = participant ? `_${participant}` : ''; + return `${fromMe ? 'true' : 'false'}_${remoteJid}_${id}${participantStr}`; +} diff --git a/yarn.lock b/yarn.lock index ecfe52a0..6c88cdd1 100644 --- a/yarn.lock +++ b/yarn.lock @@ -12879,7 +12879,7 @@ __metadata: "whatsapp-web.js@github:devlikeapro/whatsapp-web.js#fork-main-channels": version: 1.26.0 - resolution: "whatsapp-web.js@https://github.com/devlikeapro/whatsapp-web.js.git#commit=c73f3cfb4e1290103ee946ca7f56d5638f1a66bd" + resolution: "whatsapp-web.js@https://github.com/devlikeapro/whatsapp-web.js.git#commit=00fbdc333b53c052bf5ed1ecbcc0a9b3616d0e03" dependencies: "@pedroslopez/moduleraid": ^5.0.2 archiver: ^5.3.1 @@ -12897,7 +12897,7 @@ __metadata: optional: true unzipper: optional: true - checksum: 831a4a2819041f63c2a51a98e1ad7dd9ff8474b1005b11314b613c6dadac545f3c780e38081b7d260d32afa9ffd333a25f9c6407b630a0660bd08dadbc8ead89 + checksum: 9c96dc9705330372bbbf1e08f565ff0cc175a698b787f4937cf5064714b832958311fee991900e2dcbb9dfba31722940803919ca50ca2dbf51f928a650ba8764 languageName: node linkType: hard