Compare commits

...
13 Commits
Author SHA1 Message Date
devlikepro faa55ffe2d [core] 2025.5.1
Release / WEBJS - chrome - amd64 - chrome (push) Waiting to run
Release / WEBJS - chromium - amd64 - latest (push) Waiting to run
Release / WEBJS - chromium - linux/arm64 - arm (push) Waiting to run
Release / GOWS - none - amd64 - gows (push) Waiting to run
Release / GOWS - none - linux/arm64 - gows-arm (push) Waiting to run
Release / NOWEB - none - amd64 - noweb (push) Waiting to run
Release / NOWEB - none - linux/arm64 - noweb-arm (push) Waiting to run
2025-05-10 15:55:42 +07:00
devlikepro d45af4191f [core] WEBJS - properly close remote auth storage 2025-05-10 15:55:41 +07:00
devlikepro 9f8dead8b8 [core] NOWEB + SQL - speedup Contacts upsert
fix #952
2025-05-10 15:55:41 +07:00
devlikepro 564037cd53 [core] NOWEB + SQL - speedup Chats upsert
fix #952
2025-05-10 15:55:40 +07:00
devlikepro 434de6292c [core] SQL - Fix inserting duplicate entities in one batch 2025-05-10 15:55:40 +07:00
devlikepro a12e6990e5 [core] WEBJS - fix sorting groups in /chats and /chats/overview
fix #915
2025-05-10 15:55:40 +07:00
devlikepro ce0c4a094b [core] WEBJS - fix groups participant management (add/remove)
fix #944
2025-05-10 15:55:39 +07:00
devlikepro e36728ec08 [core] NOWEB - use ILogger 2025-05-10 15:55:39 +07:00
devlikepro 2c1f3090a0 [core] NOWEB - fix participant for message.ack in groups 2025-05-10 15:55:39 +07:00
devlikepro 0c6d4b3581 [core] NOWEB - update status only if not old 2025-05-10 15:55:39 +07:00
devlikepro 1d5b9eb241 [core] NOWEB - DistinctAck
fix #948
2025-05-10 15:55:38 +07:00
devlikepro b0567dafe3 [core] DistinctAck function 2025-05-10 15:55:37 +07:00
devlikepro 518f78f067 [core] Add StatusToAck and AckToStatus methods 2025-05-10 15:55:37 +07:00
14 changed files with 179 additions and 62 deletions

No files matched your search

+4 -10
View File
@@ -118,11 +118,7 @@ import { sleep, waitUntil } from '@waha/utils/promiseTimeout';
import { onlyEvent } from '@waha/utils/reactive/ops/onlyEvent';
import * as NodeCache from 'node-cache';
import {
debounceTime,
distinct,
filter,
groupBy,
interval,
merge,
mergeMap,
Observable,
@@ -136,6 +132,8 @@ import { promisify } from 'util';
import * as gows from './types';
import { MessageStatus } from './types';
import MessageServiceClient = messages.MessageServiceClient;
import { AckToStatus } from '@waha/core/utils/acks';
import { DistinctAck } from '@waha/core/utils/reactive';
enum WhatsMeowEvent {
CONNECTED = 'gows.ConnectedEventData',
@@ -374,11 +372,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
const receipt$ = all$.pipe(onlyEvent(WhatsMeowEvent.RECEIPT));
const messageAck$ = receipt$.pipe(
mergeMap(this.receiptToMessageAck.bind(this)),
// 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),
),
DistinctAck(),
);
this.events2.get(WAHAEvents.MESSAGE_ACK).switch(messageAck$);
@@ -1342,7 +1336,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
}
const status =
filter['filter.ack'] != null ? filter['filter.ack'] + 1 : null;
filter['filter.ack'] != null ? AckToStatus(filter['filter.ack']) : null;
const request = new messages.GetMessagesRequest({
session: this.session,
filters: new messages.MessageFilters({
+25 -12
View File
@@ -36,11 +36,11 @@ import {
LabelAssociationType,
} from '@adiwajshing/baileys/lib/Types/LabelAssociation';
import { MessageUserReceiptUpdate } from '@adiwajshing/baileys/lib/Types/Message';
import { ILogger } from '@adiwajshing/baileys/lib/Utils/logger';
import {
isJidBroadcast,
isLidUser,
} from '@adiwajshing/baileys/lib/WABinary/jid-utils';
import { Logger as BaileysLogger } from '@adiwajshing/baileys/node_modules/pino';
import { UnprocessableEntityException } from '@nestjs/common';
import {
ensureSuffix,
@@ -65,9 +65,11 @@ import { toVcard } from '@waha/core/helpers';
import { createAgentProxy } from '@waha/core/helpers.proxy';
import { IMediaEngineProcessor } from '@waha/core/media/IMediaEngineProcessor';
import { QR } from '@waha/core/QR';
import { AckToStatus, StatusToAck } from '@waha/core/utils/acks';
import { ExtractMessageKeysForRead } from '@waha/core/utils/convertors';
import { parseMessageIdSerialized } from '@waha/core/utils/ids';
import { isJidNewsletter, toJID } from '@waha/core/utils/jids';
import { DistinctAck } from '@waha/core/utils/reactive';
import { flipObject, splitAt } from '@waha/helpers';
import { PairingCodeResponse } from '@waha/structures/auth.dto';
import { CallData } from '@waha/structures/calls.dto';
@@ -218,7 +220,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
private autoRestartJob: SinglePeriodicJobRunner;
private msgRetryCounterCache: NodeCache;
private placeholderResendCache: NodeCache;
protected engineLogger: BaileysLogger;
protected engineLogger: ILogger;
private authNOWEBStore: any;
@@ -246,7 +248,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
this.engineLogger = this.loggerBuilder.child({
name: 'NOWEBEngine',
}) as unknown as BaileysLogger;
}) as unknown as ILogger;
// Restart job if session failed
this.startDelayedJob = new SingleDelayedJobRunner(
@@ -566,11 +568,11 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
}
private fixMessageUpsertStatus() {
// If no status - set it to WAMessageAck.DEVICE + 1
// If no status - set it to WAMessageAck.DEVICE
this.sock.ev.on('messages.upsert', ({ messages }) => {
for (const message of messages) {
if (message.status == null) {
message.status = WAMessageAck.DEVICE + 1;
message.status = AckToStatus(WAMessageAck.DEVICE);
}
}
});
@@ -920,7 +922,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
// Emit events for our reads
const updates = keys.map((key) => ({
key: key,
update: { status: WAMessageAck.READ + 1 },
update: { status: AckToStatus(WAMessageAck.READ) },
}));
this.sock?.ev.emit('messages.update', updates);
}
@@ -1755,7 +1757,9 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
filter(isMine), // ack comes only for MY messages
map(this.convertMessageReceiptUpdateToMessageAck.bind(this)),
);
const messageAck$ = merge(messageAckDirect$, messageAckGroups$);
const messageAck$ = merge(messageAckDirect$, messageAckGroups$).pipe(
DistinctAck(),
);
this.events2.get(WAHAEvents.MESSAGE_ACK).switch(messageAck$);
//
@@ -2025,7 +2029,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
const id = buildMessageId(message.key);
const body = this.extractBody(message.message);
const replyTo = this.extractReplyTo(message.message);
const ack = message.ack || message.status - 1;
const ack = message.ack || StatusToAck(message.status);
const mediaContent = extractMediaContent(message.message);
const source = this.getMessageSource(message.key.id);
return Promise.resolve({
@@ -2114,7 +2118,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
const message = event;
const fromToParticipant = getFromToParticipant(message.key);
const id = buildMessageId(message.key);
const ack = message.update.status - 1;
const ack = StatusToAck(message.update.status);
const body: WAMessageAckBody = {
id: id,
from: toCusFormat(fromToParticipant.from),
@@ -2129,7 +2133,6 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
protected convertMessageReceiptUpdateToMessageAck(event): WAMessageAckBody {
const fromToParticipant = getFromToParticipant(event.key);
const id = buildMessageId(event.key);
const receipt = event.receipt;
let ack;
@@ -2140,6 +2143,15 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
} else if (receipt.readTimestamp) {
ack = WAMessageAck.READ;
}
const key = { ...event.key };
if (key.fromMe) {
key.participant = this.getSessionMeInfo()?.id;
} else {
key.participant = event.receipt.userJid;
}
const id = buildMessageId(key);
const body: WAMessageAckBody = {
id: id,
from: toCusFormat(fromToParticipant.from),
@@ -2148,6 +2160,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
fromMe: event.key.fromMe,
ack: ack,
ackName: WAMessageAck[ack] || ACK_UNKNOWN,
_data: event,
};
return body;
}
@@ -2323,7 +2336,7 @@ function hasPath(url: string) {
}
export class NOWEBEngineMediaProcessor implements IMediaEngineProcessor<any> {
private readonly logger: BaileysLogger;
private readonly logger: ILogger;
constructor(
public session: WhatsappSessionNoWebCore,
@@ -2331,7 +2344,7 @@ export class NOWEBEngineMediaProcessor implements IMediaEngineProcessor<any> {
) {
this.logger = loggerBuilder.child({
name: NOWEBEngineMediaProcessor.name,
}) as unknown as BaileysLogger;
}) as unknown as ILogger;
}
hasMedia(message: any): boolean {
@@ -18,4 +18,6 @@ export interface IChatRepository {
deleteById(id: string): Promise<void>;
save(chat: Chat): Promise<void>;
upsertMany(chats: Chat[]): Promise<void>;
}
@@ -11,4 +11,8 @@ export interface IContactRepository {
deleteById(id: string): Promise<void>;
save(contact: Contact): Promise<void>;
upsertMany(contacts: Contact[]): Promise<void>;
getEntitiesByIds(ids: string[]): Promise<Map<string, Contact | null>>;
}
@@ -194,9 +194,26 @@ export class NowebPersistentStore implements INowebStore {
const jid = jidNormalizedUser(update.key.remoteJid!);
const message = await this.messagesRepo.getByJidById(jid, update.key.id);
if (!message) {
this.logger.warn(
`got update for non-existent message. update: '${JSON.stringify(
update,
)}'`,
);
continue;
}
const fields = { ...update.update };
// check if fields has only "status" field
const onlyStatusField =
Object.keys(fields).length === 1 &&
'status' in fields &&
fields.status !== null;
if (onlyStatusField) {
// if so, check the message don't have a newer status
if (message.status >= fields.status) {
continue;
}
}
// It can overwrite the key, so we need to delete it
delete fields['key'];
Object.assign(message, fields);
@@ -226,8 +243,8 @@ export class NowebPersistentStore implements INowebStore {
for (const chat of chats) {
delete chat['messages'];
chat.conversationTimestamp = toNumber(chat.conversationTimestamp) || null;
await this.chatRepo.save(chat);
}
await this.chatRepo.upsertMany(chats);
this.logger.info(`store sync - '${chats.length}' synced chats`);
}
@@ -333,15 +350,19 @@ export class NowebPersistentStore implements INowebStore {
}
private async onContactsUpsert(contacts: Contact[]) {
const upserts = [];
const ids = contacts.map((c) => c.id);
const contactById = await this.contactRepo.getEntitiesByIds(ids);
for (const update of contacts) {
const contact = await this.contactRepo.getById(update.id);
const contact = contactById.get(update.id) || {};
// remove undefined from data
Object.keys(update).forEach(
(key) => update[key] === undefined && delete update[key],
);
const result = { ...(contact || {}), ...update };
await this.contactRepo.save(result);
const result = { ...contact, ...update };
upserts.push(result);
}
await this.contactRepo.upsertMany(upserts);
}
private async onContactUpdate(updates: Partial<Contact>[]) {
@@ -1,5 +1,6 @@
import { ALL_JID } from '@waha/core/engines/noweb/session.noweb.core';
import { SqlKVRepository } from '@waha/core/storage/sql/SqlKVRepository';
import { AckToStatus } from '@waha/core/utils/acks';
import { GetChatMessagesFilter } from '@waha/structures/chats.dto';
import { PaginationParams } from '@waha/structures/pagination.dto';
@@ -43,7 +44,7 @@ export class SqlMessagesMethods {
query = query.whereRaw(sql, [value]);
}
if (filter['filter.ack'] != null) {
const status = filter['filter.ack'] + 1;
const status = AckToStatus(filter['filter.ack']);
const [sql, value] = this.repository.filterJson('data', 'status', status);
query = query.whereRaw(sql, [value]);
}
+11 -1
View File
@@ -17,7 +17,7 @@ exports.LoadPaginator = () => {
}
return window.lodash.orderBy(
data,
[this.pagination.sortBy],
[this.NullLast(this.pagination.sortBy)],
[this.pagination.sortOrder || 'asc'],
);
}
@@ -30,6 +30,16 @@ exports.LoadPaginator = () => {
const limit = this.pagination.limit || Infinity;
return data.slice(offset, offset + limit);
}
NullLast(field) {
return (item) => {
const value = item?.[field];
if (value == null) {
return -Infinity;
}
return value;
};
}
}
window.Paginator = Paginator;
};
+33 -26
View File
@@ -30,10 +30,12 @@ import {
} from '@waha/core/exceptions';
import { IMediaEngineProcessor } from '@waha/core/media/IMediaEngineProcessor';
import { QR } from '@waha/core/QR';
import { StatusToAck } from '@waha/core/utils/acks';
import {
parseMessageIdSerialized,
SerializeMessageKey,
} from '@waha/core/utils/ids';
import { DistinctAck } from '@waha/core/utils/reactive';
import { splitAt } from '@waha/helpers';
import { PairingCodeResponse } from '@waha/structures/auth.dto';
import {
@@ -110,19 +112,10 @@ import { sleep, waitUntil } from '@waha/utils/promiseTimeout';
import { SingleDelayedJobRunner } from '@waha/utils/SingleDelayedJobRunner';
import * as lodash from 'lodash';
import { ProtocolError } from 'puppeteer';
import {
debounceTime,
distinct,
filter,
fromEvent,
groupBy,
interval,
merge,
mergeMap,
Observable,
} from 'rxjs';
import { filter, fromEvent, merge, mergeMap, Observable } from 'rxjs';
import { map } from 'rxjs/operators';
import {
AuthStrategy,
Call,
Channel as WEBJSChannel,
Chat,
@@ -351,10 +344,11 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
private async end() {
this.engineStateCheckDelayedJob.cancel();
this.whatsapp?.removeAllListeners();
this.whatsapp?.pupBrowser?.removeAllListeners();
this.whatsapp?.pupPage?.removeAllListeners();
try {
this.whatsapp?.removeAllListeners();
this.whatsapp?.pupBrowser?.removeAllListeners();
this.whatsapp?.pupPage?.removeAllListeners();
// It's possible that browser yet starting
await waitUntil(
async () => {
@@ -365,11 +359,30 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
1_000,
10_000,
);
this.whatsapp?.destroy().catch((error) => {
this.logger.warn(error, 'Failed to destroy the client');
});
this.logger.debug(
'Successfully waited for browser to be ready for closing',
);
} catch (error) {
this.logger.error(error);
this.logger.error(
error,
'Failed while waiting for browser to be ready for closing',
);
}
try {
await this.whatsapp?.destroy();
this.logger.debug('Successfully destroyed whatsapp client');
} catch (error) {
this.logger.error(error, 'Failed to destroy whatsapp client');
}
try {
// @ts-ignore
const strategy: AuthStrategy = this.whatsapp?.authStrategy;
await strategy?.destroy();
this.logger.debug('Successfully destroyed auth strategy');
} catch (error) {
this.logger.error(error, 'Failed to destroy auth strategy');
}
}
@@ -1336,13 +1349,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
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),
),
);
const messageAck$ = messageAckAll$.pipe(DistinctAck());
this.events2.get(WAHAEvents.MESSAGE_ACK).switch(messageAck$);
//
@@ -1487,7 +1494,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
};
const fromToParticipant = getFromToParticipant(messageKey);
const id = SerializeMessageKey(messageKey);
const ack = receipt.status - 1;
const ack = StatusToAck(receipt.status);
acks.push({
id: id,
from: fromToParticipant.from,
+39 -3
View File
@@ -4,6 +4,7 @@ import { ISQLEngine } from '@waha/core/storage/sql/ISQLEngine';
import { PaginationParams } from '@waha/structures/pagination.dto';
import { KnexPaginator } from '@waha/utils/Paginator';
import Knex from 'knex';
import * as lodash from 'lodash';
export type Migration = string;
@@ -79,7 +80,16 @@ export class SqlKVRepository<Entity> {
}
private async upsertBatch(entities: Entity[]): Promise<void> {
const data = entities.map((entity) => this.dump(entity));
const all = entities.map((entity) => this.dump(entity));
// make it unique by .id
const data = lodash.uniqBy(all, (d: any) => d.id);
if (data.length != all.length) {
console.warn(
`WARNING - Duplicated entities for upsert batch: ${JSON.stringify(
entities,
)}`,
);
}
const columns = this.columns.map((c) => `"${c.fieldName}"`);
const values = data.map((d) => Object.values(d)).flat();
const sql = `INSERT INTO "${this.table}" (${columns.join(', ')})
@@ -109,8 +119,34 @@ export class SqlKVRepository<Entity> {
return this.all(query);
}
getAllByIds(ids: string[]) {
return this.all(this.select().whereIn('id', ids));
async getAllByIds(ids: string[]) {
const entitiesMap = await this.getEntitiesByIds(ids);
return Array.from(entitiesMap.values()).filter(
(entity) => entity !== null,
) as Entity[];
}
async getEntitiesByIds(ids: string[]): Promise<Map<string, Entity | null>> {
if (ids.length === 0) {
return new Map();
}
const rows = await this.engine.all(this.select().whereIn('id', ids));
const entitiesMap = new Map<string, Entity | null>();
// Initialize a map with null values for all requested IDs
for (const id of ids) {
entitiesMap.set(id, null);
}
// Fill in the map with found entities
for (const row of rows) {
if (row && row.id) {
entitiesMap.set(row.id, this.parse(row));
}
}
return entitiesMap;
}
getById(id: string): Promise<Entity | null> {
+9
View File
@@ -0,0 +1,9 @@
import { WAMessageAck } from '@waha/structures/enums.dto';
export function StatusToAck(status: number): WAMessageAck {
return status - 1;
}
export function AckToStatus(ack: number): number {
return ack + 1;
}
+10
View File
@@ -0,0 +1,10 @@
import { WAMessageAckBody } from '@waha/structures/webhooks.dto';
import { distinct, interval } from 'rxjs';
export function DistinctAck(flushEvery: number = 60_000) {
// only if we haven’t seen this key since the last flush
return distinct(
(msg: WAMessageAckBody) => `${msg.id}-${msg.ack}-${msg.participant}`,
interval(flushEvery),
);
}
+12 -2
View File
@@ -24,8 +24,8 @@ export class PaginatorInMemory extends Paginator {
}
return lodash.orderBy(
data,
this.pagination.sortBy,
this.pagination.sortOrder || 'asc',
[this.NullLast(this.pagination.sortBy)],
[this.pagination.sortOrder || 'asc'],
);
}
@@ -37,6 +37,16 @@ export class PaginatorInMemory extends Paginator {
const limit = this.pagination.limit || Infinity;
return data.slice(offset, offset + limit);
}
NullLast(field) {
return (item) => {
const value = item?.[field];
if (value == null) {
return -Infinity;
}
return value;
};
}
}
export class KnexPaginator extends Paginator {
+1 -1
View File
@@ -33,7 +33,7 @@ export function getEngineName(): string {
}
export const VERSION: WAHAEnvironment = {
version: '2025.4.2',
version: '2025.5.1',
engine: getEngineName(),
tier: getWAHAVersion(),
browser:
+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=00fbdc333b53c052bf5ed1ecbcc0a9b3616d0e03"
resolution: "whatsapp-web.js@https://github.com/devlikeapro/whatsapp-web.js.git#commit=ae0d467885ac832a4c5804f736f9da226148f3c0"
dependencies:
"@pedroslopez/moduleraid": ^5.0.2
archiver: ^5.3.1
@@ -12897,7 +12897,7 @@ __metadata:
optional: true
unzipper:
optional: true
checksum: 9c96dc9705330372bbbf1e08f565ff0cc175a698b787f4937cf5064714b832958311fee991900e2dcbb9dfba31722940803919ca50ca2dbf51f928a650ba8764
checksum: 60ef6d7d5a19a7575d9eddd9e4a42049715619e6f6e92a6d8558909fbff1de9924b2d82eb8dcbea21598e256335a967be5376421917dd01f1f391bc71bb46387
languageName: node
linkType: hard