1 parent
05132ac0f5
commit
767d1bb3e3
13 files changed
+293
-16
No files matched your search
@@ -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')
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
|
||||
|
||||
@@ -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];
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
import { GroupMetadata } from '@adiwajshing/baileys';
|
||||
import { PaginationParams } from '@waha/structures/pagination.dto';
|
||||
|
||||
export interface IGroupRepository {
|
||||
getAll(pagination?: PaginationParams): Promise<GroupMetadata[]>;
|
||||
|
||||
getById(id: string): Promise<GroupMetadata | null>;
|
||||
|
||||
deleteAll(): Promise<void>;
|
||||
|
||||
deleteById(id: string): Promise<void>;
|
||||
|
||||
save(group: GroupMetadata): Promise<void>;
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
@@ -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<Chat[]>;
|
||||
|
||||
getChatLabels(chatId: string): Promise<Label[]>;
|
||||
|
||||
getGroups(
|
||||
pagination: PaginationParams,
|
||||
refresh: boolean,
|
||||
): Promise<GroupMetadata[]>;
|
||||
}
|
||||
@@ -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<typeof makeWASocket>;
|
||||
|
||||
private store: ReturnType<typeof makeInMemoryStore>;
|
||||
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<proto.IWebMessageInfo> {
|
||||
@@ -78,4 +88,14 @@ export class NowebInMemoryStore implements INowebStore {
|
||||
getChatLabels(chatId: string): Promise<Label[]> {
|
||||
throw new BadRequestException(this.errorMessage);
|
||||
}
|
||||
|
||||
async getGroups(
|
||||
pagination: PaginationParams,
|
||||
refresh: boolean,
|
||||
): Promise<GroupMetadata[]> {
|
||||
const response = await this.socket?.groupFetchAllParticipating();
|
||||
const groups = Object.values(response);
|
||||
const paginator = new PaginatorInMemory(pagination);
|
||||
return paginator.apply(groups);
|
||||
}
|
||||
}
|
||||
@@ -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<typeof makeWASocket>;
|
||||
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<GroupMetadata>[]) {
|
||||
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<string, GroupParticipant>((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<string, GroupParticipant>,
|
||||
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<GroupMetadata[]> {
|
||||
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);
|
||||
}
|
||||
|
||||
@@ -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',
|
||||
[
|
||||
|
||||
@@ -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<GroupMetadata>
|
||||
implements IGroupRepository
|
||||
{
|
||||
protected Paginator = Paginator;
|
||||
}
|
||||
@@ -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'));
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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<Participant>;
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
Reference in new issue
Block a user