[core] Add "message.ack" event for groups and status

fix #495
fix #900
This commit is contained in:
devlikepro committed 2025-05-05 13:08:41 +07:00
1 parent bfb9dd7009
commit ef939f4357
5 files changed
+258 -16

No files matched your search

+9 -9
View File
@@ -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<WAMessage> {
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,
+163
View File
@@ -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;
}
+78 -5
View File
@@ -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<any, WAMessage>(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,
+6
View File
@@ -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}`;
}
+2 -2
View File
@@ -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