From 767d1bb3e35e4c6d555cb089ed27d781a759b22d Mon Sep 17 00:00:00 2001 From: devlikepro Date: Sun, 8 Dec 2024 18:20:16 +0700 Subject: [PATCH] [core] Fix 'rate-overlimit' error for groups fix #462 --- src/api/groups.controller.ts | 14 +- src/core/abc/session.abc.ts | 3 +- src/core/engines/noweb/session.noweb.core.ts | 8 +- .../engines/noweb/store/IGroupRepository.ts | 14 ++ src/core/engines/noweb/store/INowebStorage.ts | 3 + src/core/engines/noweb/store/INowebStore.ts | 6 + .../engines/noweb/store/NowebInMemoryStore.ts | 22 ++- .../noweb/store/NowebPersistentStore.ts | 156 +++++++++++++++++- src/core/engines/noweb/store/Schema.ts | 5 + .../store/sqlite3/Sqlite3GroupRepository.ts | 16 ++ .../noweb/store/sqlite3/Sqlite3Storage.ts | 13 ++ src/core/engines/webjs/session.webjs.core.ts | 10 +- src/structures/groups.dto.ts | 39 ++++- 13 files changed, 293 insertions(+), 16 deletions(-) create mode 100644 src/core/engines/noweb/store/IGroupRepository.ts create mode 100644 src/core/engines/noweb/store/sqlite3/Sqlite3GroupRepository.ts diff --git a/src/api/groups.controller.ts b/src/api/groups.controller.ts index dc0fe543..2f87325b 100644 --- a/src/api/groups.controller.ts +++ b/src/api/groups.controller.ts @@ -6,6 +6,9 @@ import { Param, Post, Put, + Query, + UsePipes, + ValidationPipe, } from '@nestjs/common'; import { ApiOperation, ApiSecurity, ApiTags } from '@nestjs/swagger'; import { GroupIdApiParam } from '@waha/nestjs/params/ChatIdApiParam'; @@ -19,6 +22,8 @@ import { WhatsappSession } from '../core/abc/session.abc'; import { CreateGroupRequest, DescriptionRequest, + GetGroupsParams, + GroupsPaginationParams, ParticipantsRequest, SettingsSecurityChangeInfo, SubjectRequest, @@ -43,8 +48,13 @@ export class GroupsController { @Get('') @SessionApiParam @ApiOperation({ summary: 'Get all groups.' }) - getGroups(@WorkingSessionParam session: WhatsappSession) { - return session.getGroups(); + @UsePipes(new ValidationPipe({ transform: true, whitelist: true })) + getGroups( + @WorkingSessionParam session: WhatsappSession, + @Query() pagination: GroupsPaginationParams, + @Query() params: GetGroupsParams, + ) { + return session.getGroups(pagination, params.refresh); } @Get(':id') diff --git a/src/core/abc/session.abc.ts b/src/core/abc/session.abc.ts index 22a60e42..753391bc 100644 --- a/src/core/abc/session.abc.ts +++ b/src/core/abc/session.abc.ts @@ -61,6 +61,7 @@ import { } from '../../structures/enums.dto'; import { CreateGroupRequest, + GroupsPaginationParams, ParticipantsRequest, SettingsSecurityChangeInfo, } from '../../structures/groups.dto'; @@ -522,7 +523,7 @@ export abstract class WhatsappSession { throw new NotImplementedByEngineError(); } - public getGroups() { + public getGroups(pagination: PaginationParams, refresh: boolean) { throw new NotImplementedByEngineError(); } diff --git a/src/core/engines/noweb/session.noweb.core.ts b/src/core/engines/noweb/session.noweb.core.ts index 39d1984b..3ad2feed 100644 --- a/src/core/engines/noweb/session.noweb.core.ts +++ b/src/core/engines/noweb/session.noweb.core.ts @@ -1084,12 +1084,14 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { return this.sock.groupCreate(request.name, participants); } - public async getGroups() { - return await this.sock.groupFetchAllParticipating(); + public async getGroups(pagination: PaginationParams, refresh: boolean) { + const groups = await this.store.getGroups(pagination, refresh); + // return {id: group} mapping for backward compatability + return lodash.keyBy(groups, 'id'); } public async getGroup(id) { - const groups = await this.sock.groupFetchAllParticipating(); + const groups = await this.getGroups({}, false); return groups[id]; } diff --git a/src/core/engines/noweb/store/IGroupRepository.ts b/src/core/engines/noweb/store/IGroupRepository.ts new file mode 100644 index 00000000..b008ba46 --- /dev/null +++ b/src/core/engines/noweb/store/IGroupRepository.ts @@ -0,0 +1,14 @@ +import { GroupMetadata } from '@adiwajshing/baileys'; +import { PaginationParams } from '@waha/structures/pagination.dto'; + +export interface IGroupRepository { + getAll(pagination?: PaginationParams): Promise; + + getById(id: string): Promise; + + deleteAll(): Promise; + + deleteById(id: string): Promise; + + save(group: GroupMetadata): Promise; +} diff --git a/src/core/engines/noweb/store/INowebStorage.ts b/src/core/engines/noweb/store/INowebStorage.ts index cc0b0a39..06f9633a 100644 --- a/src/core/engines/noweb/store/INowebStorage.ts +++ b/src/core/engines/noweb/store/INowebStorage.ts @@ -1,5 +1,6 @@ import { WAMessage } from '@adiwajshing/baileys'; import { LabelAssociation } from '@adiwajshing/baileys/lib/Types/LabelAssociation'; +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'; @@ -16,6 +17,8 @@ export abstract class INowebStorage { abstract getChatRepository(): IChatRepository; + abstract getGroupRepository(): IGroupRepository; + abstract getMessagesRepository(): IMessagesRepository; abstract getLabelsRepository(): ILabelsRepository; diff --git a/src/core/engines/noweb/store/INowebStore.ts b/src/core/engines/noweb/store/INowebStore.ts index 016d07ac..8a8c179e 100644 --- a/src/core/engines/noweb/store/INowebStore.ts +++ b/src/core/engines/noweb/store/INowebStore.ts @@ -4,6 +4,7 @@ import { Contact, proto, } from '@adiwajshing/baileys'; +import { GroupMetadata } from '@adiwajshing/baileys/lib/Types/GroupMetadata'; import { Label } from '@adiwajshing/baileys/lib/Types/Label'; import { GetChatMessagesFilter } from '@waha/structures/chats.dto'; import { PaginationParams } from '@waha/structures/pagination.dto'; @@ -40,4 +41,9 @@ export interface INowebStore { getChatsByLabelId(labelId: string): Promise; getChatLabels(chatId: string): Promise; + + getGroups( + pagination: PaginationParams, + refresh: boolean, + ): Promise; } diff --git a/src/core/engines/noweb/store/NowebInMemoryStore.ts b/src/core/engines/noweb/store/NowebInMemoryStore.ts index f6a82ffc..215ee53a 100644 --- a/src/core/engines/noweb/store/NowebInMemoryStore.ts +++ b/src/core/engines/noweb/store/NowebInMemoryStore.ts @@ -1,8 +1,15 @@ -import { Chat, Contact, makeInMemoryStore, proto } from '@adiwajshing/baileys'; +import makeWASocket, { + Chat, + Contact, + GroupMetadata, + makeInMemoryStore, + proto, +} from '@adiwajshing/baileys'; import { Label } from '@adiwajshing/baileys/lib/Types/Label'; import { BadRequestException } from '@nestjs/common'; import { GetChatMessagesFilter } from '@waha/structures/chats.dto'; import { PaginationParams } from '@waha/structures/pagination.dto'; +import { PaginatorInMemory } from '@waha/utils/Paginator'; import { INowebStore } from './INowebStore'; @@ -10,6 +17,8 @@ import { INowebStore } from './INowebStore'; const logger = require('pino')(); export class NowebInMemoryStore implements INowebStore { + private socket: ReturnType; + private store: ReturnType; errorMessage = 'Enable NOWEB store "config.noweb.store.enabled=True" and "config.noweb.store.full_sync=True" when starting a new session. ' + @@ -33,6 +42,7 @@ export class NowebInMemoryStore implements INowebStore { bind(ev: any, socket: any) { this.store.bind(ev); + this.socket = socket; } loadMessage(jid: string, id: string): Promise { @@ -78,4 +88,14 @@ export class NowebInMemoryStore implements INowebStore { getChatLabels(chatId: string): Promise { throw new BadRequestException(this.errorMessage); } + + async getGroups( + pagination: PaginationParams, + refresh: boolean, + ): Promise { + const response = await this.socket?.groupFetchAllParticipating(); + const groups = Object.values(response); + const paginator = new PaginatorInMemory(pagination); + return paginator.apply(groups); + } } diff --git a/src/core/engines/noweb/store/NowebPersistentStore.ts b/src/core/engines/noweb/store/NowebPersistentStore.ts index 785508f3..1c4dbc02 100644 --- a/src/core/engines/noweb/store/NowebPersistentStore.ts +++ b/src/core/engines/noweb/store/NowebPersistentStore.ts @@ -1,27 +1,34 @@ -import { +import makeWASocket, { + areJidsSameUser, BaileysEventEmitter, Chat, ChatUpdate, Contact, + GroupParticipant, isRealMessage, jidNormalizedUser, + ParticipantAction, proto, updateMessageWithReaction, updateMessageWithReceipt, } from '@adiwajshing/baileys'; +import { GroupMetadata } from '@adiwajshing/baileys/lib/Types/GroupMetadata'; import { Label } from '@adiwajshing/baileys/lib/Types/Label'; import { LabelAssociation, LabelAssociationType, } from '@adiwajshing/baileys/lib/Types/LabelAssociation'; +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 { GetChatMessagesFilter } from '@waha/structures/chats.dto'; import { PaginationParams, SortOrder } from '@waha/structures/pagination.dto'; +import { DefaultMap } from '@waha/utils/DefaultMap'; +import { sleep, waitUntil } from '@waha/utils/promiseTimeout'; +import * as lodash from 'lodash'; import { toNumber } from 'lodash'; import { Logger } from 'pino'; -import { toJID } from '../session.noweb.core'; import { IChatRepository } from './IChatRepository'; import { IContactRepository } from './IContactRepository'; import { IMessagesRepository } from './IMessagesRepository'; @@ -31,15 +38,26 @@ import { INowebStore } from './INowebStore'; // eslint-disable-next-line @typescript-eslint/no-var-requires const AsyncLock = require('async-lock'); +type ms = number; +const HOUR: ms = 60 * 60 * 1000; + export class NowebPersistentStore implements INowebStore { - private socket: any; + private socket: ReturnType; private chatRepo: IChatRepository; + private groupRepo: IGroupRepository; private contactRepo: IContactRepository; private messagesRepo: IMessagesRepository; private labelsRepo: ILabelsRepository; private labelAssociationsRepo: ILabelAssociationRepository; public presences: any; private lock: any; + private groupsFetchLock: any = new AsyncLock({ + maxPending: Infinity, + maxExecutionTime: 60_000, + }); + + private lastTimeGroupUpdate: Date = new Date(0); + private GROUP_METADATA_CACHE_TIME = 24 * HOUR; constructor( private logger: Logger, @@ -47,6 +65,7 @@ export class NowebPersistentStore implements INowebStore { ) { this.socket = null; this.chatRepo = storage.getChatRepository(); + this.groupRepo = storage.getGroupRepository(); this.contactRepo = storage.getContactsRepository(); this.messagesRepo = storage.getMessagesRepository(); this.labelsRepo = storage.getLabelsRepository(); @@ -88,6 +107,19 @@ export class NowebPersistentStore implements INowebStore { ev.on('chats.delete', (data) => this.withLock('chats', () => this.onChatDelete(data)), ); + // Groups + ev.on('groups.upsert', (data) => + this.withLock('groups', () => this.onGroupUpsert(data)), + ); + ev.on('groups.update', (data) => + this.withLock('groups', () => this.onGroupUpdate(data)), + ); + ev.on('group-participants.update', (data) => + this.withLock(`group-${data.id}`, () => + this.onGroupParticipantsUpdate(data), + ), + ); + // Contacts ev.on('contacts.upsert', (data) => this.withLock('contacts', () => this.onContactsUpsert(data)), @@ -195,7 +227,84 @@ export class NowebPersistentStore implements INowebStore { chat.conversationTimestamp = toNumber(chat.conversationTimestamp); await this.chatRepo.save(chat); } - this.logger.info(`history sync - '${chats.length}' synced chats`); + this.logger.info(`store sync - '${chats.length}' synced chats`); + } + + private async onGroupUpsert(groups: GroupMetadata[]) { + for (const group of groups) { + await this.groupRepo.save(group); + } + this.logger.info(`store sync - '${groups.length}' synced groups`); + } + + private async onGroupUpdate(groups: Partial[]) { + for (const update of groups) { + let group = await this.groupRepo.getById(update.id); + group = Object.assign(group || {}, update) as GroupMetadata; + await this.groupRepo.save(group); + } + this.logger.info(`store sync - '${groups.length}' updated groups`); + this.lastTimeGroupUpdate = new Date(); + } + + private async onGroupParticipantsUpdate(data) { + const id: string = data.id; + const participants: string[] = data.participants; + const action: ParticipantAction = data.action; + + if (action == 'remove') { + // Remove the group if the current user is removed + const myJid = this.socket?.authState?.creds?.me?.id; + const participantsIncludesMe = lodash.find(participants, (p) => + areJidsSameUser(p, myJid), + ); + if (participantsIncludesMe) { + await this.groupRepo.deleteById(id); + return; + } + } + + let group = await this.groupRepo.getById(id); + if (!group) { + group = { id: id, participants: [] } as GroupMetadata; + } + + const participantsById = new DefaultMap((key) => { + return { id: key, admin: null } as GroupParticipant; + }); + for (const participant of group.participants) { + participantsById.set(participant.id, participant); + } + + for (const participant of participants) { + this.participantUpdate(participantsById, participant, action); + } + group.participants = Array.from(participantsById.values()); + await this.groupRepo.save(group); + } + + private participantUpdate( + participantsById: DefaultMap, + participant: string, + action: ParticipantAction, + ) { + switch (action) { + case 'add': + // if there's no participant - add it (by id) + participantsById.get(participant); + break; + case 'remove': + // remove the participant (by id) + participantsById.delete(participant); + break; + case 'promote': + // set admin: admin + participantsById.get(participant).admin = 'admin'; + break; + case 'demote': + participantsById.get(participant).admin = null; + break; + } } private async onChatUpdate(updates: ChatUpdate[]) { @@ -344,6 +453,45 @@ export class NowebPersistentStore implements INowebStore { return this.chatRepo.getAllWithMessages(pagination); } + private shouldUpdateGroup(): boolean { + const timePassed = + new Date().getTime() - this.lastTimeGroupUpdate.getTime(); + return timePassed > this.GROUP_METADATA_CACHE_TIME; + } + + private async fetchGroups() { + await this.groupsFetchLock.acquire('groups-fetch', async () => { + if (!this.shouldUpdateGroup()) { + // Update has been done by another request + return; + } + const lastTimeGroupUpdate = this.lastTimeGroupUpdate; + await this.groupRepo.deleteAll(); + await this.socket?.groupFetchAllParticipating(); + // Wait until the groups update is done + await waitUntil( + async () => this.lastTimeGroupUpdate > lastTimeGroupUpdate, + 100, + 5_000, + ); + }); + } + + async getGroups( + pagination: PaginationParams, + refresh: boolean, + ): Promise { + if (refresh) { + // Reset the last update time + this.lastTimeGroupUpdate = new Date(0); + } + + if (this.shouldUpdateGroup()) { + await this.fetchGroups(); + } + return this.groupRepo.getAll(pagination); + } + getContactById(jid) { return this.contactRepo.getById(jid); } diff --git a/src/core/engines/noweb/store/Schema.ts b/src/core/engines/noweb/store/Schema.ts index ef09952d..722c3b15 100644 --- a/src/core/engines/noweb/store/Schema.ts +++ b/src/core/engines/noweb/store/Schema.ts @@ -18,6 +18,11 @@ export const NOWEB_STORE_SCHEMA = [ new Index('chats_conversationTimestamp_index', ['conversationTimestamp']), ], ), + new Schema( + 'groups', + [new Field('id', 'TEXT'), new Field('data', 'TEXT')], + [new Index('groups_id_index', ['id'])], + ), new Schema( 'messages', [ diff --git a/src/core/engines/noweb/store/sqlite3/Sqlite3GroupRepository.ts b/src/core/engines/noweb/store/sqlite3/Sqlite3GroupRepository.ts new file mode 100644 index 00000000..07e11a30 --- /dev/null +++ b/src/core/engines/noweb/store/sqlite3/Sqlite3GroupRepository.ts @@ -0,0 +1,16 @@ +import { GroupMetadata } from '@adiwajshing/baileys/lib/Types/GroupMetadata'; +import { IGroupRepository } from '@waha/core/engines/noweb/store/IGroupRepository'; +import { KnexPaginator } from '@waha/utils/Paginator'; + +import { NOWEBSqlite3KVRepository } from './NOWEBSqlite3KVRepository'; + +class Paginator extends KnexPaginator { + indexes = ['id']; +} + +export class Sqlite3GroupRepository + extends NOWEBSqlite3KVRepository + implements IGroupRepository +{ + protected Paginator = Paginator; +} diff --git a/src/core/engines/noweb/store/sqlite3/Sqlite3Storage.ts b/src/core/engines/noweb/store/sqlite3/Sqlite3Storage.ts index 043c0530..1b7fd8e7 100644 --- a/src/core/engines/noweb/store/sqlite3/Sqlite3Storage.ts +++ b/src/core/engines/noweb/store/sqlite3/Sqlite3Storage.ts @@ -2,6 +2,7 @@ import { WAMessage } from '@adiwajshing/baileys'; import { LabelAssociation } from '@adiwajshing/baileys/lib/Types/LabelAssociation'; import { ILabelAssociationRepository } from '@waha/core/engines/noweb/store/ILabelAssociationsRepository'; import { ILabelsRepository } from '@waha/core/engines/noweb/store/ILabelsRepository'; +import { Sqlite3GroupRepository } from '@waha/core/engines/noweb/store/sqlite3/Sqlite3GroupRepository'; import { Sqlite3LabelAssociationsRepository } from '@waha/core/engines/noweb/store/sqlite3/Sqlite3LabelAssociationsRepository'; import { Sqlite3LabelsRepository } from '@waha/core/engines/noweb/store/sqlite3/Sqlite3LabelsRepository'; import { Field, Index, Schema } from '@waha/core/storage/sqlite3/Schema'; @@ -62,6 +63,14 @@ export class Sqlite3Storage extends INowebStorage { 'CREATE INDEX IF NOT EXISTS chats_conversationTimestamp_index ON chats (conversationTimestamp)', ); + // Groups + this.db.exec( + 'CREATE TABLE IF NOT EXISTS groups (id TEXT PRIMARY KEY, data TEXT)', + ); + this.db.exec( + 'CREATE UNIQUE INDEX IF NOT EXISTS groups_id_index ON groups (id)', + ); + // Messages this.db.exec( 'CREATE TABLE IF NOT EXISTS messages (jid TEXT, id TEXT, messageTimestamp INTEGER, data TEXT)', @@ -119,6 +128,10 @@ export class Sqlite3Storage extends INowebStorage { return new Sqlite3ChatRepository(this.db, this.getSchema('chats')); } + getGroupRepository() { + return new Sqlite3GroupRepository(this.db, this.getSchema('groups')); + } + getLabelsRepository(): ILabelsRepository { return new Sqlite3LabelsRepository(this.db, this.getSchema('labels')); } diff --git a/src/core/engines/webjs/session.webjs.core.ts b/src/core/engines/webjs/session.webjs.core.ts index ecdaa0e6..539f4182 100644 --- a/src/core/engines/webjs/session.webjs.core.ts +++ b/src/core/engines/webjs/session.webjs.core.ts @@ -75,6 +75,7 @@ import { import { PaginatorInMemory } from '@waha/utils/Paginator'; import { sleep, waitUntil } from '@waha/utils/promiseTimeout'; import { SingleDelayedJobRunner } from '@waha/utils/SingleDelayedJobRunner'; +import * as lodash from 'lodash'; import { fromEvent, merge, mergeMap, Observable } from 'rxjs'; import { map } from 'rxjs/operators'; import { @@ -789,10 +790,11 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { return groupChat.setMessagesAdminsOnly(value); } - public getGroups() { - return this.whatsapp - .getChats() - .then((chats) => chats.filter((chat) => chat.isGroup)); + public async getGroups(pagination: PaginationParams, refresh: boolean) { + const chats = await this.whatsapp.getChats(); + const groups = lodash.filter(chats, (chat) => chat.isGroup); + const paginator = new PaginatorInMemory(pagination); + return paginator.apply(groups); } public getGroup(id) { diff --git a/src/structures/groups.dto.ts b/src/structures/groups.dto.ts index 402e1494..12e89873 100644 --- a/src/structures/groups.dto.ts +++ b/src/structures/groups.dto.ts @@ -1,5 +1,15 @@ import { ApiProperty } from '@nestjs/swagger'; -import { IsArray, IsString } from 'class-validator'; +import { BooleanString } from '@waha/nestjs/validation/BooleanString'; +import { Transform } from 'class-transformer'; +import { + IsArray, + IsBoolean, + IsEnum, + IsOptional, + IsString, +} from 'class-validator'; + +import { PaginationParams } from './pagination.dto'; /** * Structures @@ -46,3 +56,30 @@ export class CreateGroupRequest { @IsArray() participants: Array; } + +enum GroupSortField { + ID = 'id', + SUBJECT = 'subject', +} + +export class GroupsPaginationParams extends PaginationParams { + @ApiProperty({ + description: 'Sort by field', + enum: GroupSortField, + }) + @IsOptional() + @IsEnum(GroupSortField) + sortBy?: string; +} + +export class GetGroupsParams { + @ApiProperty({ + description: 'Refresh the groups list and participants from the server', + example: false, + required: false, + }) + @Transform(BooleanString) + @IsBoolean() + @IsOptional() + refresh: boolean = false; +}