From e69d992725b4999ad2fe412c13f3c2f12e9db9f5 Mon Sep 17 00:00:00 2001 From: devlikepro Date: Thu, 28 Aug 2025 11:56:12 +0700 Subject: [PATCH] [core] Ignore - Status, Groups, Broadcast fix #1142 fix #1259 fix #1190 --- src/config.service.ts | 18 ++ src/core/abc/manager.abc.ts | 8 + src/core/abc/session.abc.ts | 21 +- src/core/engines/gows/grpc/gows.ts | 153 +++++++++- src/core/engines/gows/grpc/gows_pb.js | 265 +++++++++++++++++- src/core/engines/gows/session.gows.core.ts | 24 +- src/core/engines/noweb/session.noweb.core.ts | 13 +- .../noweb/store/NowebPersistentStore.ts | 40 ++- src/core/engines/webjs/session.webjs.core.ts | 14 + src/core/manager.core.ts | 1 + src/core/utils/jids.ts | 29 +- src/structures/sessions.dto.ts | 36 +++ 12 files changed, 607 insertions(+), 15 deletions(-) diff --git a/src/config.service.ts b/src/config.service.ts index 152e6adc..cf5757af 100644 --- a/src/config.service.ts +++ b/src/config.service.ts @@ -1,6 +1,7 @@ import { Injectable, Logger, OnApplicationBootstrap } from '@nestjs/common'; import { ConfigService } from '@nestjs/config'; import { GlobalWebhookConfigConfig } from '@waha/core/config/GlobalWebhookConfig'; +import { IgnoreJidConfig } from '@waha/core/utils/jids'; import { parseBool } from './helpers'; import { WebhookConfig } from './structures/webhooks.config.dto'; @@ -178,6 +179,23 @@ export class WhatsappConfigService implements OnApplicationBootstrap { return parseBool(value); } + /** + * Global default "ignore settings" for chats. + * If not defined, defaults to false (do not ignore anything). + */ + getIgnoreChatsConfig(): IgnoreJidConfig { + const status = parseBool( + this.configService.get('WAHA_SESSION_CONFIG_IGNORE_STATUS', 'false'), + ); + const groups = parseBool( + this.configService.get('WAHA_SESSION_CONFIG_IGNORE_GROUPS', 'false'), + ); + const channels = parseBool( + this.configService.get('WAHA_SESSION_CONFIG_IGNORE_CHANNELS', 'false'), + ); + return { status, groups, channels }; + } + onApplicationBootstrap() { const error = this.webhookConfig.validateConfig(); if (error) { diff --git a/src/core/abc/manager.abc.ts b/src/core/abc/manager.abc.ts index cd208ee9..6cd188ba 100644 --- a/src/core/abc/manager.abc.ts +++ b/src/core/abc/manager.abc.ts @@ -14,9 +14,11 @@ import { GowsEngineConfigService } from '@waha/core/config/GowsEngineConfigServi import { GowsBootstrap } from '@waha/core/engines/gows/GowsBootstrap'; import { ISessionMeRepository } from '@waha/core/storage/ISessionMeRepository'; import { ISessionWorkerRepository } from '@waha/core/storage/ISessionWorkerRepository'; +import { IgnoreJidConfig } from '@waha/core/utils/jids'; import { WAHAWebhook } from '@waha/structures/webhooks.dto'; import { waitUntil } from '@waha/utils/promiseTimeout'; import { VERSION } from '@waha/version'; +import * as lodash from 'lodash'; import { PinoLogger } from 'nestjs-pino'; import { merge, Observable, of } from 'rxjs'; @@ -235,6 +237,12 @@ export abstract class SessionManager } return new NoopEngineBootstrap(); } + + protected ignoreChatsConfig(config: SessionConfig) { + const ignore: IgnoreJidConfig = this.config.getIgnoreChatsConfig(); + // Given the default, overwrite from the config if any + return lodash.defaults({}, config?.ignore, ignore); + } } export function populateSessionInfo( diff --git a/src/core/abc/session.abc.ts b/src/core/abc/session.abc.ts index df9f8e71..324380d0 100644 --- a/src/core/abc/session.abc.ts +++ b/src/core/abc/session.abc.ts @@ -4,7 +4,11 @@ import { IMediaConverter, } from '@waha/core/media/IConverter'; import { MessagesForRead } from '@waha/core/utils/convertors'; -import { isJidNewsletter } from '@waha/core/utils/jids'; +import { + IgnoreJidConfig, + isJidNewsletter, + JidFilter, +} from '@waha/core/utils/jids'; import { Channel, ChannelListResult, @@ -90,10 +94,7 @@ import { WAHAPresenceStatus, WAHASessionStatus, } from '../../structures/enums.dto'; -import { - EventCancelRequest, - EventMessageRequest, -} from '../../structures/events.dto'; +import { EventMessageRequest } from '../../structures/events.dto'; import { CreateGroupRequest, GroupField, @@ -153,8 +154,11 @@ export interface SessionParams { loggerBuilder: LoggerBuilder; sessionStore: DataStore; proxyConfig?: ProxyConfig; + // Raw unchanged SessionConfig sessionConfig?: SessionConfig; engineConfig?: any; + // Ignore settings + ignore: IgnoreJidConfig; } export abstract class WhatsappSession { @@ -169,6 +173,7 @@ export abstract class WhatsappSession { public sessionConfig?: SessionConfig; protected engineConfig?: any; protected unpairing: boolean = false; + protected jids: JidFilter; private _status: WAHASessionStatus; private shouldPrintQR: boolean; @@ -195,6 +200,7 @@ export abstract class WhatsappSession { mediaManager, sessionConfig, engineConfig, + ignore, }: SessionParams) { this.status$ = new BehaviorSubject(null); @@ -260,6 +266,11 @@ export abstract class WhatsappSession { this.sessionConfig = sessionConfig; this.engineConfig = engineConfig; this.shouldPrintQR = printQR; + this.logger.info( + { ignore: ignore }, + 'The session ignores the following chat ids', + ); + this.jids = new JidFilter(ignore); } public getEventObservable(event: WAHAEvents) { diff --git a/src/core/engines/gows/grpc/gows.ts b/src/core/engines/gows/grpc/gows.ts index 81f36dd0..06b06927 100644 --- a/src/core/engines/gows/grpc/gows.ts +++ b/src/core/engines/gows/grpc/gows.ts @@ -1302,13 +1302,128 @@ export namespace messages { return SessionProxyConfig.deserialize(bytes); } } - export class SessionConfig extends pb_1.Message { + export class SessionIgnoreJidsConfig extends pb_1.Message { #one_of_decls: number[][] = []; constructor(data?: any[] | { + status?: boolean; + groups?: boolean; + newsletters?: boolean; + }) { + super(); + pb_1.Message.initialize(this, Array.isArray(data) ? data : [], 0, -1, [], this.#one_of_decls); + if (!Array.isArray(data) && typeof data == "object") { + if ("status" in data && data.status != undefined) { + this.status = data.status; + } + if ("groups" in data && data.groups != undefined) { + this.groups = data.groups; + } + if ("newsletters" in data && data.newsletters != undefined) { + this.newsletters = data.newsletters; + } + } + } + get status() { + return pb_1.Message.getFieldWithDefault(this, 1, false) as boolean; + } + set status(value: boolean) { + pb_1.Message.setField(this, 1, value); + } + get groups() { + return pb_1.Message.getFieldWithDefault(this, 2, false) as boolean; + } + set groups(value: boolean) { + pb_1.Message.setField(this, 2, value); + } + get newsletters() { + return pb_1.Message.getFieldWithDefault(this, 3, false) as boolean; + } + set newsletters(value: boolean) { + pb_1.Message.setField(this, 3, value); + } + static fromObject(data: { + status?: boolean; + groups?: boolean; + newsletters?: boolean; + }): SessionIgnoreJidsConfig { + const message = new SessionIgnoreJidsConfig({}); + if (data.status != null) { + message.status = data.status; + } + if (data.groups != null) { + message.groups = data.groups; + } + if (data.newsletters != null) { + message.newsletters = data.newsletters; + } + return message; + } + toObject() { + const data: { + status?: boolean; + groups?: boolean; + newsletters?: boolean; + } = {}; + if (this.status != null) { + data.status = this.status; + } + if (this.groups != null) { + data.groups = this.groups; + } + if (this.newsletters != null) { + data.newsletters = this.newsletters; + } + return data; + } + serialize(): Uint8Array; + serialize(w: pb_1.BinaryWriter): void; + serialize(w?: pb_1.BinaryWriter): Uint8Array | void { + const writer = w || new pb_1.BinaryWriter(); + if (this.status != false) + writer.writeBool(1, this.status); + if (this.groups != false) + writer.writeBool(2, this.groups); + if (this.newsletters != false) + writer.writeBool(3, this.newsletters); + if (!w) + return writer.getResultBuffer(); + } + static deserialize(bytes: Uint8Array | pb_1.BinaryReader): SessionIgnoreJidsConfig { + const reader = bytes instanceof pb_1.BinaryReader ? bytes : new pb_1.BinaryReader(bytes), message = new SessionIgnoreJidsConfig(); + while (reader.nextField()) { + if (reader.isEndGroup()) + break; + switch (reader.getFieldNumber()) { + case 1: + message.status = reader.readBool(); + break; + case 2: + message.groups = reader.readBool(); + break; + case 3: + message.newsletters = reader.readBool(); + break; + default: reader.skipField(); + } + } + return message; + } + serializeBinary(): Uint8Array { + return this.serialize(); + } + static deserializeBinary(bytes: Uint8Array): SessionIgnoreJidsConfig { + return SessionIgnoreJidsConfig.deserialize(bytes); + } + } + export class SessionConfig extends pb_1.Message { + #one_of_decls: number[][] = [[4]]; + constructor(data?: any[] | ({ store?: SessionStoreConfig; log?: SessionLogConfig; proxy?: SessionProxyConfig; - }) { + } & (({ + ignore?: SessionIgnoreJidsConfig; + })))) { super(); pb_1.Message.initialize(this, Array.isArray(data) ? data : [], 0, -1, [], this.#one_of_decls); if (!Array.isArray(data) && typeof data == "object") { @@ -1321,6 +1436,9 @@ export namespace messages { if ("proxy" in data && data.proxy != undefined) { this.proxy = data.proxy; } + if ("ignore" in data && data.ignore != undefined) { + this.ignore = data.ignore; + } } } get store() { @@ -1350,10 +1468,29 @@ export namespace messages { get has_proxy() { return pb_1.Message.getField(this, 3) != null; } + get ignore() { + return pb_1.Message.getWrapperField(this, SessionIgnoreJidsConfig, 4) as SessionIgnoreJidsConfig; + } + set ignore(value: SessionIgnoreJidsConfig) { + pb_1.Message.setOneofWrapperField(this, 4, this.#one_of_decls[0], value); + } + get has_ignore() { + return pb_1.Message.getField(this, 4) != null; + } + get _ignore() { + const cases: { + [index: number]: "none" | "ignore"; + } = { + 0: "none", + 4: "ignore" + }; + return cases[pb_1.Message.computeOneofCase(this, [4])]; + } static fromObject(data: { store?: ReturnType; log?: ReturnType; proxy?: ReturnType; + ignore?: ReturnType; }): SessionConfig { const message = new SessionConfig({}); if (data.store != null) { @@ -1365,6 +1502,9 @@ export namespace messages { if (data.proxy != null) { message.proxy = SessionProxyConfig.fromObject(data.proxy); } + if (data.ignore != null) { + message.ignore = SessionIgnoreJidsConfig.fromObject(data.ignore); + } return message; } toObject() { @@ -1372,6 +1512,7 @@ export namespace messages { store?: ReturnType; log?: ReturnType; proxy?: ReturnType; + ignore?: ReturnType; } = {}; if (this.store != null) { data.store = this.store.toObject(); @@ -1382,6 +1523,9 @@ export namespace messages { if (this.proxy != null) { data.proxy = this.proxy.toObject(); } + if (this.ignore != null) { + data.ignore = this.ignore.toObject(); + } return data; } serialize(): Uint8Array; @@ -1394,6 +1538,8 @@ export namespace messages { writer.writeMessage(2, this.log, () => this.log.serialize(writer)); if (this.has_proxy) writer.writeMessage(3, this.proxy, () => this.proxy.serialize(writer)); + if (this.has_ignore) + writer.writeMessage(4, this.ignore, () => this.ignore.serialize(writer)); if (!w) return writer.getResultBuffer(); } @@ -1412,6 +1558,9 @@ export namespace messages { case 3: reader.readMessage(message.proxy, () => message.proxy = SessionProxyConfig.deserialize(reader)); break; + case 4: + reader.readMessage(message.ignore, () => message.ignore = SessionIgnoreJidsConfig.deserialize(reader)); + break; default: reader.skipField(); } } diff --git a/src/core/engines/gows/grpc/gows_pb.js b/src/core/engines/gows/grpc/gows_pb.js index 3d2c205b..dd3d5332 100644 --- a/src/core/engines/gows/grpc/gows_pb.js +++ b/src/core/engines/gows/grpc/gows_pb.js @@ -102,6 +102,7 @@ goog.exportSymbol('proto.messages.SearchPageResult', null, global); goog.exportSymbol('proto.messages.Section', null, global); goog.exportSymbol('proto.messages.Session', null, global); goog.exportSymbol('proto.messages.SessionConfig', null, global); +goog.exportSymbol('proto.messages.SessionIgnoreJidsConfig', null, global); goog.exportSymbol('proto.messages.SessionLogConfig', null, global); goog.exportSymbol('proto.messages.SessionProxyConfig', null, global); goog.exportSymbol('proto.messages.SessionStateResponse', null, global); @@ -452,6 +453,27 @@ if (goog.DEBUG && !COMPILED) { */ proto.messages.SessionProxyConfig.displayName = 'proto.messages.SessionProxyConfig'; } +/** + * Generated by JsPbCodeGenerator. + * @param {Array=} opt_data Optional initial data array, typically from a + * server response, or constructed directly in Javascript. The array is used + * in place and becomes part of the constructed object. It is not cloned. + * If no data is provided, the constructed object will be empty, but still + * valid. + * @extends {jspb.Message} + * @constructor + */ +proto.messages.SessionIgnoreJidsConfig = function(opt_data) { + jspb.Message.initialize(this, opt_data, 0, -1, null, null); +}; +goog.inherits(proto.messages.SessionIgnoreJidsConfig, jspb.Message); +if (goog.DEBUG && !COMPILED) { + /** + * @public + * @override + */ + proto.messages.SessionIgnoreJidsConfig.displayName = 'proto.messages.SessionIgnoreJidsConfig'; +} /** * Generated by JsPbCodeGenerator. * @param {Array=} opt_data Optional initial data array, typically from a @@ -4372,6 +4394,196 @@ proto.messages.SessionProxyConfig.prototype.setUrl = function(value) { +if (jspb.Message.GENERATE_TO_OBJECT) { +/** + * Creates an object representation of this proto. + * Field names that are reserved in JavaScript and will be renamed to pb_name. + * Optional fields that are not set will be set to undefined. + * To access a reserved field use, foo.pb_, eg, foo.pb_default. + * For the list of reserved names please see: + * net/proto2/compiler/js/internal/generator.cc#kKeyword. + * @param {boolean=} opt_includeInstance Deprecated. whether to include the + * JSPB instance for transitional soy proto support: + * http://goto/soy-param-migration + * @return {!Object} + */ +proto.messages.SessionIgnoreJidsConfig.prototype.toObject = function(opt_includeInstance) { + return proto.messages.SessionIgnoreJidsConfig.toObject(opt_includeInstance, this); +}; + + +/** + * Static version of the {@see toObject} method. + * @param {boolean|undefined} includeInstance Deprecated. Whether to include + * the JSPB instance for transitional soy proto support: + * http://goto/soy-param-migration + * @param {!proto.messages.SessionIgnoreJidsConfig} msg The msg instance to transform. + * @return {!Object} + * @suppress {unusedLocalVariables} f is only used for nested messages + */ +proto.messages.SessionIgnoreJidsConfig.toObject = function(includeInstance, msg) { + var f, obj = { + status: jspb.Message.getBooleanFieldWithDefault(msg, 1, false), + groups: jspb.Message.getBooleanFieldWithDefault(msg, 2, false), + newsletters: jspb.Message.getBooleanFieldWithDefault(msg, 3, false) + }; + + if (includeInstance) { + obj.$jspbMessageInstance = msg; + } + return obj; +}; +} + + +/** + * Deserializes binary data (in protobuf wire format). + * @param {jspb.ByteSource} bytes The bytes to deserialize. + * @return {!proto.messages.SessionIgnoreJidsConfig} + */ +proto.messages.SessionIgnoreJidsConfig.deserializeBinary = function(bytes) { + var reader = new jspb.BinaryReader(bytes); + var msg = new proto.messages.SessionIgnoreJidsConfig; + return proto.messages.SessionIgnoreJidsConfig.deserializeBinaryFromReader(msg, reader); +}; + + +/** + * Deserializes binary data (in protobuf wire format) from the + * given reader into the given message object. + * @param {!proto.messages.SessionIgnoreJidsConfig} msg The message object to deserialize into. + * @param {!jspb.BinaryReader} reader The BinaryReader to use. + * @return {!proto.messages.SessionIgnoreJidsConfig} + */ +proto.messages.SessionIgnoreJidsConfig.deserializeBinaryFromReader = function(msg, reader) { + while (reader.nextField()) { + if (reader.isEndGroup()) { + break; + } + var field = reader.getFieldNumber(); + switch (field) { + case 1: + var value = /** @type {boolean} */ (reader.readBool()); + msg.setStatus(value); + break; + case 2: + var value = /** @type {boolean} */ (reader.readBool()); + msg.setGroups(value); + break; + case 3: + var value = /** @type {boolean} */ (reader.readBool()); + msg.setNewsletters(value); + break; + default: + reader.skipField(); + break; + } + } + return msg; +}; + + +/** + * Serializes the message to binary data (in protobuf wire format). + * @return {!Uint8Array} + */ +proto.messages.SessionIgnoreJidsConfig.prototype.serializeBinary = function() { + var writer = new jspb.BinaryWriter(); + proto.messages.SessionIgnoreJidsConfig.serializeBinaryToWriter(this, writer); + return writer.getResultBuffer(); +}; + + +/** + * Serializes the given message to binary data (in protobuf wire + * format), writing to the given BinaryWriter. + * @param {!proto.messages.SessionIgnoreJidsConfig} message + * @param {!jspb.BinaryWriter} writer + * @suppress {unusedLocalVariables} f is only used for nested messages + */ +proto.messages.SessionIgnoreJidsConfig.serializeBinaryToWriter = function(message, writer) { + var f = undefined; + f = message.getStatus(); + if (f) { + writer.writeBool( + 1, + f + ); + } + f = message.getGroups(); + if (f) { + writer.writeBool( + 2, + f + ); + } + f = message.getNewsletters(); + if (f) { + writer.writeBool( + 3, + f + ); + } +}; + + +/** + * optional bool status = 1; + * @return {boolean} + */ +proto.messages.SessionIgnoreJidsConfig.prototype.getStatus = function() { + return /** @type {boolean} */ (jspb.Message.getBooleanFieldWithDefault(this, 1, false)); +}; + + +/** + * @param {boolean} value + * @return {!proto.messages.SessionIgnoreJidsConfig} returns this + */ +proto.messages.SessionIgnoreJidsConfig.prototype.setStatus = function(value) { + return jspb.Message.setProto3BooleanField(this, 1, value); +}; + + +/** + * optional bool groups = 2; + * @return {boolean} + */ +proto.messages.SessionIgnoreJidsConfig.prototype.getGroups = function() { + return /** @type {boolean} */ (jspb.Message.getBooleanFieldWithDefault(this, 2, false)); +}; + + +/** + * @param {boolean} value + * @return {!proto.messages.SessionIgnoreJidsConfig} returns this + */ +proto.messages.SessionIgnoreJidsConfig.prototype.setGroups = function(value) { + return jspb.Message.setProto3BooleanField(this, 2, value); +}; + + +/** + * optional bool newsletters = 3; + * @return {boolean} + */ +proto.messages.SessionIgnoreJidsConfig.prototype.getNewsletters = function() { + return /** @type {boolean} */ (jspb.Message.getBooleanFieldWithDefault(this, 3, false)); +}; + + +/** + * @param {boolean} value + * @return {!proto.messages.SessionIgnoreJidsConfig} returns this + */ +proto.messages.SessionIgnoreJidsConfig.prototype.setNewsletters = function(value) { + return jspb.Message.setProto3BooleanField(this, 3, value); +}; + + + + + if (jspb.Message.GENERATE_TO_OBJECT) { /** * Creates an object representation of this proto. @@ -4403,7 +4615,8 @@ proto.messages.SessionConfig.toObject = function(includeInstance, msg) { var f, obj = { store: (f = msg.getStore()) && proto.messages.SessionStoreConfig.toObject(includeInstance, f), log: (f = msg.getLog()) && proto.messages.SessionLogConfig.toObject(includeInstance, f), - proxy: (f = msg.getProxy()) && proto.messages.SessionProxyConfig.toObject(includeInstance, f) + proxy: (f = msg.getProxy()) && proto.messages.SessionProxyConfig.toObject(includeInstance, f), + ignore: (f = msg.getIgnore()) && proto.messages.SessionIgnoreJidsConfig.toObject(includeInstance, f) }; if (includeInstance) { @@ -4455,6 +4668,11 @@ proto.messages.SessionConfig.deserializeBinaryFromReader = function(msg, reader) reader.readMessage(value,proto.messages.SessionProxyConfig.deserializeBinaryFromReader); msg.setProxy(value); break; + case 4: + var value = new proto.messages.SessionIgnoreJidsConfig; + reader.readMessage(value,proto.messages.SessionIgnoreJidsConfig.deserializeBinaryFromReader); + msg.setIgnore(value); + break; default: reader.skipField(); break; @@ -4508,6 +4726,14 @@ proto.messages.SessionConfig.serializeBinaryToWriter = function(message, writer) proto.messages.SessionProxyConfig.serializeBinaryToWriter ); } + f = message.getIgnore(); + if (f != null) { + writer.writeMessage( + 4, + f, + proto.messages.SessionIgnoreJidsConfig.serializeBinaryToWriter + ); + } }; @@ -4622,6 +4848,43 @@ proto.messages.SessionConfig.prototype.hasProxy = function() { }; +/** + * optional SessionIgnoreJidsConfig ignore = 4; + * @return {?proto.messages.SessionIgnoreJidsConfig} + */ +proto.messages.SessionConfig.prototype.getIgnore = function() { + return /** @type{?proto.messages.SessionIgnoreJidsConfig} */ ( + jspb.Message.getWrapperField(this, proto.messages.SessionIgnoreJidsConfig, 4)); +}; + + +/** + * @param {?proto.messages.SessionIgnoreJidsConfig|undefined} value + * @return {!proto.messages.SessionConfig} returns this +*/ +proto.messages.SessionConfig.prototype.setIgnore = function(value) { + return jspb.Message.setWrapperField(this, 4, value); +}; + + +/** + * Clears the message field making it undefined. + * @return {!proto.messages.SessionConfig} returns this + */ +proto.messages.SessionConfig.prototype.clearIgnore = function() { + return this.setIgnore(undefined); +}; + + +/** + * Returns whether this field is set. + * @return {boolean} + */ +proto.messages.SessionConfig.prototype.hasIgnore = function() { + return jspb.Message.getField(this, 4) != null; +}; + + diff --git a/src/core/engines/gows/session.gows.core.ts b/src/core/engines/gows/session.gows.core.ts index 4ede3749..3fbe7aac 100644 --- a/src/core/engines/gows/session.gows.core.ts +++ b/src/core/engines/gows/session.gows.core.ts @@ -254,6 +254,11 @@ export class WhatsappSessionGoWSCore extends WhatsappSession { proxy: new messages.SessionProxyConfig({ url: this.getProxyUrl(this.proxyConfig), }), + ignore: new messages.SessionIgnoreJidsConfig({ + status: this.jids.ignore.status, + groups: this.jids.ignore.groups, + newsletters: this.jids.ignore.channels, + }), }), }); @@ -391,7 +396,11 @@ export class WhatsappSessionGoWSCore extends WhatsappSession { const all$ = this.all$; this.events2.get(WAHAEvents.ENGINE_EVENT).switch(all$); - const messages$ = all$.pipe(onlyEvent(WhatsMeowEvent.MESSAGE)); + const messages$ = all$.pipe( + onlyEvent(WhatsMeowEvent.MESSAGE), + filter((msg: any) => this.jids.include(msg?.Info?.Chat)), + share(), + ); let [messagesFromMe$, messagesFromOthers$] = partition(messages$, isMine); messagesFromMe$ = messagesFromMe$.pipe( @@ -461,7 +470,10 @@ export class WhatsappSessionGoWSCore extends WhatsappSession { ); this.events2.get(WAHAEvents.MESSAGE_EDITED).switch(messagesEdited$); - const receipt$ = all$.pipe(onlyEvent(WhatsMeowEvent.RECEIPT)); + const receipt$ = all$.pipe( + onlyEvent(WhatsMeowEvent.RECEIPT), + filter((r: any) => this.jids.include(r?.Chat)), + ); const messageAck$ = receipt$.pipe( mergeMap(this.receiptToMessageAck.bind(this)), DistinctAck(), @@ -476,9 +488,13 @@ export class WhatsappSessionGoWSCore extends WhatsappSession { const presence$ = all$.pipe( onlyEvent(WhatsMeowEvent.PRESENCE), + filter((event: any) => this.jids.include(event?.From)), filter((event) => !isJidGroup(event.From)), ); - const chatPresence$ = all$.pipe(onlyEvent(WhatsMeowEvent.CHAT_PRESENCE)); + const chatPresence$ = all$.pipe( + onlyEvent(WhatsMeowEvent.CHAT_PRESENCE), + filter((event: any) => this.jids.include(event?.Chat)), + ); const presenceUpdates$ = merge(presence$, chatPresence$).pipe( map((event) => this.toWahaPresences(event.From || event.Chat, [event])), ); @@ -538,6 +554,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession { // const pollVoteEvent$ = all$.pipe( onlyEvent(WhatsMeowEvent.POLL_VOTE_EVENT), + filter((event: any) => this.jids.include(event?.Info?.Chat)), map(this.toPollVotePayload.bind(this)), filter(Boolean), share(), @@ -557,6 +574,7 @@ export class WhatsappSessionGoWSCore extends WhatsappSession { // const eventMessageResponse$ = all$.pipe( onlyEvent(WhatsMeowEvent.EVENT_MESSAGE_RESPONSE), + filter((event: any) => this.jids.include(event?.Info?.Chat)), map(this.toEventResponsePayload.bind(this)), filter(Boolean), ); diff --git a/src/core/engines/noweb/session.noweb.core.ts b/src/core/engines/noweb/session.noweb.core.ts index ae169e57..08cc515f 100644 --- a/src/core/engines/noweb/session.noweb.core.ts +++ b/src/core/engines/noweb/session.noweb.core.ts @@ -385,6 +385,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { this.store = new NowebPersistentStore( this.loggerBuilder.child({ name: NowebPersistentStore.name }), storage, + this.jids, ); await this.store.init(); } @@ -1846,6 +1847,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { const messagesUpsert$ = fromEvent(this.sock.ev, 'messages.upsert').pipe( map((event: BaileysEventMap['messages.upsert']) => event.messages), mergeAll(), + filter((msg) => this.jids.include(msg.key.remoteJid)), share(), ); let [messagesFromMe$, messagesFromOthers$] = partition( @@ -1933,6 +1935,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { ).pipe( // @ts-ignore mergeAll(), + filter((update) => this.jids.include(update.key.remoteJid)), share(), ); const messageAckDirect$ = messageUpdates$.pipe( @@ -1944,6 +1947,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { fromEvent(this.sock.ev, 'message-receipt.update').pipe( // @ts-ignore mergeAll(), + filter((update) => this.jids.include(update.key.remoteJid)), share(), ); @@ -2011,6 +2015,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { this.events2.get(WAHAEvents.PRESENCE_UPDATE).switch( fromEvent(this.sock.ev, 'presence.update').pipe( + filter((presence: any) => this.jids.include(presence.id)), map((data: any) => this.toWahaPresences(data.id, data.presences)), share(), ), @@ -2041,7 +2046,13 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { // // @ts-ignore const calls$: Observable = fromEvent(this.sock.ev, 'call'); - const call$ = calls$.pipe(mergeMap(identity), share()); + const call$ = calls$.pipe( + mergeMap(identity), + filter((call: WACallEvent) => + this.jids.include(call.groupJid || call.chatId), + ), + share(), + ); this.events2.get(WAHAEvents.CALL_RECEIVED).switch( call$.pipe( filter((call: WACallEvent) => call.status === 'offer'), diff --git a/src/core/engines/noweb/store/NowebPersistentStore.ts b/src/core/engines/noweb/store/NowebPersistentStore.ts index d61e8224..8135a03d 100644 --- a/src/core/engines/noweb/store/NowebPersistentStore.ts +++ b/src/core/engines/noweb/store/NowebPersistentStore.ts @@ -14,6 +14,7 @@ import makeWASocket, { updateMessageWithReceipt, WAMessage, } from '@adiwajshing/baileys'; +import { WACallEvent } from '@adiwajshing/baileys/lib/Types/Call'; import { GroupMetadata } from '@adiwajshing/baileys/lib/Types/GroupMetadata'; import { Label } from '@adiwajshing/baileys/lib/Types/Label'; import { @@ -24,6 +25,7 @@ import { isLidUser } from '@adiwajshing/baileys/lib/WABinary/jid-utils'; import { IGroupRepository } from '@waha/core/engines/noweb/store/IGroupRepository'; import { ILabelAssociationRepository } from '@waha/core/engines/noweb/store/ILabelAssociationsRepository'; import { ILabelsRepository } from '@waha/core/engines/noweb/store/ILabelsRepository'; +import { JidFilter } from '@waha/core/utils/jids'; import { GetChatMessagesFilter, OverviewFilter, @@ -39,6 +41,7 @@ import { waitUntil } from '@waha/utils/promiseTimeout'; import * as lodash from 'lodash'; import { toNumber } from 'lodash'; import { Logger } from 'pino'; +import { filter } from 'rxjs'; import { IChatRepository } from './IChatRepository'; import { IContactRepository } from './IContactRepository'; @@ -82,6 +85,7 @@ export class NowebPersistentStore implements INowebStore { constructor( private logger: Logger, public storage: INowebStorage, + private jids: JidFilter, ) { this.socket = null; this.chatRepo = storage.getChatRepository(); @@ -247,6 +251,7 @@ export class NowebPersistentStore implements INowebStore { private async syncMessagesHistory(messages) { const realMessages = messages.filter(isRealMessage); + messages = messages.filter((msg) => this.jids.include(msg.key.remoteJid)); await this.messagesRepo.upsert(realMessages); this.logger.info( `history sync - '${messages.length}' got messages, '${realMessages.length}' real messages`, @@ -254,11 +259,13 @@ export class NowebPersistentStore implements INowebStore { } private async onMessagesUpsert(update) { - const { messages, type } = update; + const type = update.type; if (type !== 'notify' && type !== 'append') { this.logger.debug(`unexpected type for messages.upsert: '${type}'`); return; } + let messages = update.messages; + messages = messages.filter((msg) => this.jids.include(msg.key.remoteJid)); const realMessages = messages.filter(isRealMessage); await this.messagesRepo.upsert(realMessages); this.logger.debug( @@ -270,6 +277,9 @@ export class NowebPersistentStore implements INowebStore { for (const update of updates) { // eslint-disable-next-line @typescript-eslint/no-non-null-assertion const jid = jidNormalizedUser(update.key.remoteJid!); + if (!this.jids.include(jid)) { + continue; + } if (!update.key.id) { continue; } @@ -333,12 +343,16 @@ export class NowebPersistentStore implements INowebStore { delete chat['messages']; chat.conversationTimestamp = toNumber(chat.conversationTimestamp) || null; } + chats = chats.filter((chat) => this.jids.include(chat.id)); await this.chatRepo.upsertMany(chats); this.logger.info(`store sync - '${chats.length}' synced chats`); } private async onGroupUpsert(groups: GroupMetadata[]) { for (const group of groups) { + if (!this.jids.include(group.id)) { + continue; + } await this.groupRepo.save(group); } this.logger.info(`store sync - '${groups.length}' synced groups`); @@ -346,6 +360,9 @@ export class NowebPersistentStore implements INowebStore { private async onGroupUpdate(groups: Partial[]) { for (const update of groups) { + if (!this.jids.include(update.id)) { + continue; + } let group = await this.groupRepo.getById(update.id); group = Object.assign(group || {}, update) as GroupMetadata; await this.groupRepo.save(group); @@ -356,6 +373,9 @@ export class NowebPersistentStore implements INowebStore { private async onGroupParticipantsUpdate(data) { const id: string = data.id; + if (!this.jids.include(id)) { + return; + } const participants: string[] = data.participants; const action: ParticipantAction = data.action; @@ -419,6 +439,9 @@ export class NowebPersistentStore implements INowebStore { private async onChatUpdate(updates: ChatUpdate[]) { for (const update of updates) { + if (!this.jids.include(update.id)) { + continue; + } const chat = (await this.chatRepo.getById(update.id)) || ({} as Chat); Object.assign(chat, update); chat.conversationTimestamp = toNumber(chat.conversationTimestamp) || null; @@ -447,6 +470,9 @@ export class NowebPersistentStore implements INowebStore { const ids = contacts.map((c) => c.id); const contactById = await this.contactRepo.getEntitiesByIds(ids); for (const update of contacts) { + if (!this.jids.include(update.id)) { + continue; + } const contact = contactById.get(update.id) || {}; // remove undefined from data Object.keys(update).forEach( @@ -460,6 +486,9 @@ export class NowebPersistentStore implements INowebStore { private async onContactUpdate(updates: Partial[]) { for (const update of updates) { + if (!this.jids.include(update.id)) { + continue; + } let contact = await this.contactRepo.getById(update.id); if (!contact) { @@ -493,6 +522,9 @@ export class NowebPersistentStore implements INowebStore { private async onMessageReaction(reactions) { for (const { key, reaction } of reactions) { + if (!this.jids.include(key.remoteJid)) { + continue; + } const msg = await this.messagesRepo.getByJidById(key.remoteJid, key.id); if (!msg) { this.logger.warn( @@ -509,6 +541,9 @@ export class NowebPersistentStore implements INowebStore { private async onMessageReceiptUpdate(updates) { for (const { key, receipt } of updates) { + if (!this.jids.include(key.remoteJid)) { + continue; + } const msg = await this.messagesRepo.getByJidById(key.remoteJid, key.id); if (!msg) { this.logger.warn( @@ -544,6 +579,9 @@ export class NowebPersistentStore implements INowebStore { } private async onPresenceUpdate({ id, presences: update }) { + if (!this.jids.include(id)) { + return; + } this.presences[id] = this.presences[id] || {}; Object.assign(this.presences[id], update); } diff --git a/src/core/engines/webjs/session.webjs.core.ts b/src/core/engines/webjs/session.webjs.core.ts index e375ea02..bdc7a418 100644 --- a/src/core/engines/webjs/session.webjs.core.ts +++ b/src/core/engines/webjs/session.webjs.core.ts @@ -1440,6 +1440,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { // const messageReceived$ = fromEvent(this.whatsapp, Events.MESSAGE_RECEIVED); const messagesFromOthers$ = messageReceived$.pipe( + filter((msg: Message) => this.jids.include(msg?.id?.remote)), mergeMap((msg: any) => this.processIncomingMessage(msg, true)), share(), ); @@ -1447,6 +1448,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { const messageCreate$ = fromEvent(this.whatsapp, Events.MESSAGE_CREATE); const messagesFromAll$ = messageCreate$.pipe( + filter((msg: Message) => this.jids.include(msg?.id?.remote)), mergeMap((msg: any) => this.processIncomingMessage(msg, true)), share(), ); @@ -1457,6 +1459,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { Events.MESSAGE_CIPHERTEXT, ); const messagesWaiting$ = messageCiphertext$.pipe( + filter((msg: Message) => this.jids.include(msg?.id?.remote)), mergeMap((msg: any) => this.processIncomingMessage(msg, false)), share(), ); @@ -1470,6 +1473,9 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { }, ); const messagesRevoked$ = messageRevoked$.pipe( + filter((evt: any) => + this.jids.include(evt?.after?.id?.remote || evt?.before?.id?.remote), + ), map((event): WAMessageRevokedBody => { const afterMessage = event.after ? this.toWAMessage(event.after) : null; const beforeMessage = event.before @@ -1488,6 +1494,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { const messageReaction$ = fromEvent(this.whatsapp, 'message_reaction'); const messagesReaction$ = messageReaction$.pipe( + filter((reaction: Reaction) => this.jids.include(reaction?.id?.remote)), map(this.processMessageReaction.bind(this)), filter(Boolean), ); @@ -1501,6 +1508,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { }, ); const messagesEdit$ = messageEdit$.pipe( + filter((event: any) => this.jids.include(event?.message?.id?.remote)), map((event): WAMessageEditedBody => { const message = this.toWAMessage(event.message); return { @@ -1524,6 +1532,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { map((event) => event.message), map(this.toWAMessage.bind(this)), filter((ack) => !isJidGroup(ack.to) && !isJidStatusBroadcast(ack.to)), + filter((ack) => this.jids.include(ack.to)), ); const tagReceiptNode$ = fromEvent(this.whatsapp, Events.TAG_RECEIPT); const messageAckGroups$ = tagReceiptNode$.pipe( @@ -1533,6 +1542,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { filter(Boolean), mergeMap(this.TagReceiptToMessageAck.bind(this)), filter((ack) => isJidGroup(ack.to) || isJidStatusBroadcast(ack.to)), + filter((ack) => this.jids.include(ack.to)), ); const messageAckAll$ = merge(messagesAckDM$, messageAckGroups$); @@ -1553,11 +1563,13 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { const presences$ = tagPresenceNode$.pipe( map(TagPresenceToPresence), filter(Boolean), + filter((presence: any) => this.jids.include(presence.id)), ); const tagChatstateNode$ = fromEvent(this.whatsapp, 'tag:chatstate'); const chatstatePresences$ = tagChatstateNode$.pipe( map(TagChatstateToPresence), filter(Boolean), + filter((presence: any) => this.jids.include(presence.id)), ); const presenceUpdate$ = merge(presences$, chatstatePresences$); this.events2.get(WAHAEvents.PRESENCE_UPDATE).switch(presenceUpdate$); @@ -1626,6 +1638,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { }, ); const chatsArchived$ = chatArchived$.pipe( + filter((event: any) => this.jids.include(event?.chat?.id?._serialized)), map((event) => { return { id: event.chat.id._serialized, @@ -1641,6 +1654,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { // const call$ = fromEvent(this.whatsapp, 'call'); const calls$ = call$.pipe( + filter((call: Call) => this.jids.include((call as any)?.from)), map((call: Call) => { return { id: call.id, diff --git a/src/core/manager.core.ts b/src/core/manager.core.ts index d8afbe79..6b49a707 100644 --- a/src/core/manager.core.ts +++ b/src/core/manager.core.ts @@ -200,6 +200,7 @@ export class SessionManagerCore extends SessionManager implements OnModuleInit { sessionStore: this.store, proxyConfig: proxyConfig, sessionConfig: this.sessionConfig, + ignore: this.ignoreChatsConfig(this.sessionConfig), }; if (this.EngineClass === WhatsappSessionWebJSCore) { sessionConfig.engineConfig = this.webjsEngineConfigService.getConfig(); diff --git a/src/core/utils/jids.ts b/src/core/utils/jids.ts index 0facf945..78f44920 100644 --- a/src/core/utils/jids.ts +++ b/src/core/utils/jids.ts @@ -1,11 +1,15 @@ -import { isJidGroup, isJidMetaIa } from '@adiwajshing/baileys'; +import { + isJidGroup, + isJidMetaIa, + isJidStatusBroadcast, +} from '@adiwajshing/baileys'; import { isJidBroadcast, isLidUser, } from '@adiwajshing/baileys/lib/WABinary/jid-utils'; export function isJidNewsletter(jid: string) { - return jid.endsWith('@newsletter'); + return jid?.endsWith('@newsletter'); } export function isJidCus(jid: string) { @@ -35,3 +39,24 @@ export function toJID(chatId) { const number = chatId.split('@')[0]; return number + '@s.whatsapp.net'; } + +export interface IgnoreJidConfig { + status: boolean; + groups: boolean; + channels: boolean; +} + +export class JidFilter { + constructor(public ignore: IgnoreJidConfig) {} + + include(jid: string): boolean { + if (this.ignore.status && isJidStatusBroadcast(jid)) { + return false; + } else if (this.ignore.groups && isJidGroup(jid)) { + return false; + } else if (this.ignore.channels && isJidNewsletter(jid)) { + return false; + } + return true; + } +} diff --git a/src/structures/sessions.dto.ts b/src/structures/sessions.dto.ts index ace736c9..e78de33f 100644 --- a/src/structures/sessions.dto.ts +++ b/src/structures/sessions.dto.ts @@ -89,6 +89,29 @@ export class NowebConfig { markOnline: boolean = true; } +export class IgnoreConfig { + @ApiProperty({ + description: 'Ignore a status@broadcast (stories) events', + }) + @IsBoolean() + @IsOptional() + status?: boolean; + + @ApiProperty({ + description: 'Ignore groups events', + }) + @IsBoolean() + @IsOptional() + groups?: boolean; + + @ApiProperty({ + description: 'Ignore channels events', + }) + @IsBoolean() + @IsOptional() + channels?: boolean; +} + export class SessionConfig { @ValidateNested({ each: true }) @Type(() => WebhookConfig) @@ -125,6 +148,19 @@ export class SessionConfig { @IsOptional() debug?: boolean; + @ApiProperty({ + example: { + status: null, + groups: null, + channels: null, + }, + description: 'Ignore some events related to specific chats', + }) + @ValidateNested() + @Type(() => IgnoreConfig) + @IsOptional() + ignore?: IgnoreConfig; + @ApiProperty({ example: { store: {