diff --git a/package.json b/package.json index 928bcb83..215626f0 100644 --- a/package.json +++ b/package.json @@ -75,9 +75,10 @@ "qrcode-terminal": "^0.12.0", "reflect-metadata": "^0.1.13", "rimraf": "^3.0.2", - "rxjs": "^7.1.0", + "rxjs": "^7.8.1", "sharp": "^0.33.4", "swagger-ui-express": "^4.1.4", + "ulid": "^2.3.0", "whatsapp-web.js": "github:devlikeapro/whatsapp-web.js#fork-main-channels" }, "optionalDependencies": { diff --git a/src/core/abc/manager.abc.ts b/src/core/abc/manager.abc.ts index 1e60e275..4e9ff5dd 100644 --- a/src/core/abc/manager.abc.ts +++ b/src/core/abc/manager.abc.ts @@ -8,8 +8,8 @@ import { ISessionWorkerRepository } from '@waha/core/storage/ISessionWorkerRepos import { WAHAWebhook } from '@waha/structures/webhooks.dto'; import { waitUntil } from '@waha/utils/promiseTimeout'; import { VERSION } from '@waha/version'; -import { EventEmitter } from 'events'; import { PinoLogger } from 'nestjs-pino'; +import { merge, Observable, of } from 'rxjs'; import { WAHAEngine, @@ -35,7 +35,6 @@ export abstract class SessionManager implements BeforeApplicationShutdown { public sessionConfigRepository: ISessionConfigRepository; protected sessionMeRepository: ISessionMeRepository; protected sessionWorkerRepository: ISessionWorkerRepository; - public events: EventEmitter; private lock: any; WAIT_STATUS_INTERVAL = 500; @@ -71,6 +70,19 @@ export abstract class SessionManager implements BeforeApplicationShutdown { protected abstract get EngineClass(): typeof WhatsappSession; + public getSessionEvent(session: string, event: WAHAEvents): Observable { + return of(); + } + + public getSessionEvents( + session: string, + events: WAHAEvents[], + ): Observable { + return merge( + ...events.map((event) => this.getSessionEvent(session, event)), + ); + } + // // API Methods // @@ -111,22 +123,6 @@ export abstract class SessionManager implements BeforeApplicationShutdown { await this.sessionWorkerRepository?.unassign(name, this.workerId); } - protected handleSessionEvent(event: WAHAEvents, session: WhatsappSession) { - return (payload: any) => { - const me = session.getSessionMeInfo(); - const data: WAHAWebhook = { - event: event, - session: session.name, - metadata: session.sessionConfig?.metadata, - me: me, - payload: payload, - engine: session.engine, - environment: VERSION, - }; - this.events.emit(event, data); - }; - } - async getWorkingSession(sessionName: string): Promise { return this.waitUntilStatus(sessionName, [WAHASessionStatus.WORKING]); } @@ -157,7 +153,28 @@ export abstract class SessionManager implements BeforeApplicationShutdown { return session; } - async beforeApplicationShutdown(signal?: string) { - this.events.removeAllListeners(); + beforeApplicationShutdown(signal?: string) { + return; } } +export function populateSessionInfo( + event: WAHAEvents, + session: WhatsappSession, +) { + return (payload: any): WAHAWebhook => { + const id = payload._eventId; + const data = { ...payload }; + delete data._eventId; + const me = session.getSessionMeInfo(); + return { + id: id, + event: event, + session: session.name, + metadata: session.sessionConfig?.metadata, + me: me, + payload: data, + engine: session.engine, + environment: VERSION, + }; + }; +} diff --git a/src/core/abc/session.abc.ts b/src/core/abc/session.abc.ts index 4705b878..a511e9c4 100644 --- a/src/core/abc/session.abc.ts +++ b/src/core/abc/session.abc.ts @@ -12,11 +12,26 @@ import { SendButtonsRequest } from '@waha/structures/chatting.buttons.dto'; import { Label, LabelDTO, LabelID } from '@waha/structures/labels.dto'; import { PaginationParams } from '@waha/structures/pagination.dto'; import { WAMessage } from '@waha/structures/responses.dto'; +import { DefaultMap } from '@waha/utils/DefaultMap'; +import { generatePrefixedId } from '@waha/utils/ids'; import { LoggerBuilder } from '@waha/utils/logging'; -import { EventEmitter } from 'events'; +import { complete } from '@waha/utils/reactive/complete'; +import { SwitchObservable } from '@waha/utils/reactive/SwitchObservable'; import * as fs from 'fs'; import * as lodash from 'lodash'; import { Logger } from 'pino'; +import { + BehaviorSubject, + catchError, + delay, + filter, + of, + retry, + share, + Subject, + switchMap, +} from 'rxjs'; +import { distinctUntilChanged, map } from 'rxjs/operators'; import { MessageId } from 'whatsapp-web.js'; import { @@ -100,7 +115,6 @@ export interface SessionParams { } export abstract class WhatsappSession { - public events: EventEmitter; public engine: WAHAEngine; public name: string; @@ -115,6 +129,8 @@ export abstract class WhatsappSession { private _status: WAHASessionStatus; private shouldPrintQR: boolean; + protected events2: DefaultMap>; + private status$: Subject; public constructor({ name, @@ -126,11 +142,64 @@ export abstract class WhatsappSession { sessionConfig, engineConfig, }: SessionParams) { - this.events = new EventEmitter(); + this.status$ = new BehaviorSubject(null); + this.name = name; this.proxyConfig = proxyConfig; this.loggerBuilder = loggerBuilder; this.logger = loggerBuilder.child({ name: 'WhatsappSession' }); + this.events2 = new DefaultMap>( + (key) => + new SwitchObservable((obs$) => { + return obs$.pipe( + catchError((err) => { + this.logger.error( + `Caught error, dropping value from, event: '${key}'`, + ); + this.logger.error(err, err.stack); + throw err; + }), + filter(Boolean), + map((data) => { + data._eventId = generatePrefixedId('evt'); + return data; + }), + retry(), + share(), + ); + }), + ); + + this.events2.get(WAHAEvents.SESSION_STATUS).switch( + this.status$ + // initial value is null + .pipe(filter(Boolean)) + // Wait for WORKING status to get all the info + // https://github.com/devlikeapro/waha/issues/409 + .pipe( + switchMap((status: WAHASessionStatus) => { + const me = this.getSessionMeInfo(); + const hasMe = !!me?.pushName && !!me?.id; + // Delay WORKING by 1 second if condition is met + // Usually we get WORKING with all the info after + if (status === WAHASessionStatus.WORKING && !hasMe) { + return of(status).pipe(delay(2000)); + } + return of(status); + }), + // Remove consecutive duplicate WORKING statuses + distinctUntilChanged( + (prev, curr) => prev === curr && curr === WAHASessionStatus.WORKING, + ), + ) + // Populate the session info + .pipe( + map((status) => { + return { name: this.name, status: status }; + }), + ), + ); + this.sessionStore = sessionStore; this.mediaManager = mediaManager; this.sessionConfig = sessionConfig; @@ -138,6 +207,10 @@ export abstract class WhatsappSession { this.shouldPrintQR = printQR; } + public getEventObservable(event: WAHAEvents) { + return this.events2.get(event); + } + protected set status(value: WAHASessionStatus) { if (this.unpairing && value !== WAHASessionStatus.STOPPED) { // In case of unpairing @@ -145,8 +218,7 @@ export abstract class WhatsappSession { return; } this._status = value; - const body: WASessionStatusBody = { name: this.name, status: value }; - this.events.emit(WAHAEvents.SESSION_STATUS, body); + this.status$.next(value); } public get status() { @@ -219,6 +291,10 @@ export abstract class WhatsappSession { /** Stop the session */ abstract stop(): Promise; + protected stopEvents() { + complete(this.events2); + } + /* Unpair the account */ async unpair(): Promise { return; diff --git a/src/core/api/websocket.gateway.core.ts b/src/core/api/websocket.gateway.core.ts index f94531cd..8d1ef1c9 100644 --- a/src/core/api/websocket.gateway.core.ts +++ b/src/core/api/websocket.gateway.core.ts @@ -14,10 +14,10 @@ import { import { SessionManager } from '@waha/core/abc/manager.abc'; import { WebsocketHeartbeatJob } from '@waha/nestjs/ws/WebsocketHeartbeatJob'; import { WebSocket } from '@waha/nestjs/ws/ws'; -import { WAHAEvents } from '@waha/structures/enums.dto'; +import { WAHAEvents, WAHAEventsWild } from '@waha/structures/enums.dto'; +import { EventWildUnmask } from '@waha/utils/events'; import { generatePrefixedId } from '@waha/utils/ids'; import { IncomingMessage } from 'http'; -import * as lodash from 'lodash'; import * as url from 'url'; import { Server } from 'ws'; @@ -37,11 +37,9 @@ export class WebsocketGatewayCore @WebSocketServer() server: Server; - private listeners: Map = - new Map(); - private readonly logger: LoggerService; private heartbeat: WebsocketHeartbeatJob; + private eventUnmask = new EventWildUnmask(WAHAEvents, WAHAEventsWild); constructor(private manager: SessionManager) { this.logger = new Logger('WebsocketGateway'); @@ -58,47 +56,42 @@ export class WebsocketGatewayCore this.logger.debug(`New client connected: ${request.url}`); const params = this.getParams(id, request, socket); if (!params) { + const error = "Invalid parameters. Provide 'session' and 'events' params"; + socket.close(4001, JSON.stringify({ error })); return; } - const { session, events } = params; + const session: string = params.session; + const events: WAHAEvents[] = params.events; this.logger.debug( `Client connected to session: '${session}', events: ${events}, ${id}`, ); - this.listeners.set(socket, { session, events }); + + const sub = this.manager + .getSessionEvents(session, events) + .subscribe((data) => { + this.logger.debug(`Sending data to client, event.id: ${data.id}`, data); + socket.send(JSON.stringify(data)); + }); + socket.on('close', () => { + this.logger.debug(`Client disconnected - ${socket.id}`); + sub.unsubscribe(); + }); } private getParams(id: string, request: IncomingMessage, socket: WebSocket) { const query = url.parse(request.url, true).query; const session = (query.session as string) || '*'; - if (session !== '*') { - this.logger.warn( - `Only connecting to all sessions is allowed for now, use session=*, ${id}`, - ); - const error = - 'Only connecting to all sessions is allowed for now, use session=*'; - socket.close(4001, JSON.stringify({ error })); - return null; - } - - const events = ((query.events as string) || '*').split(','); - if ( - !lodash.isEqual(events, ['*']) && - !lodash.isEqual(events, [WAHAEvents.SESSION_STATUS]) - ) { - this.logger.warn( - `Only \'session.status\' event is allowed for now, use events=session.status or events=*, ${id}`, - ); - const error = - "Only 'session.status' event is allowed for now, use events=session.status or events=*"; - socket.close(4001, JSON.stringify({ error })); - return null; + let paramsEvents = (query.events as string[]) || '*'; + // if params events string - split by "," + if (typeof paramsEvents === 'string') { + paramsEvents = paramsEvents.split(','); } + const events = this.eventUnmask.unmask(paramsEvents); return { session, events }; } handleDisconnect(socket: WebSocket): any { this.logger.debug(`Client disconnected - ${socket.id}`); - this.listeners.delete(socket); } async beforeApplicationShutdown(signal?: string) { @@ -118,24 +111,9 @@ export class WebsocketGatewayCore afterInit(server: Server) { this.logger.debug('Websocket server initialized'); - this.manager.events.on( - WAHAEvents.SESSION_STATUS, - this.sendToAll.bind(this), - ); - this.logger.debug('Subscribed to manager events'); this.logger.debug('Starting heartbeat service...'); this.heartbeat.start(server); this.logger.debug('Heartbeat service started'); } - - sendToAll(data: any) { - this.listeners.forEach((options, client) => { - if (options.events.length === 0) { - return; - } - this.logger.debug('Sending data to client', data); - client.send(JSON.stringify(data)); - }); - } } diff --git a/src/core/engines/noweb/session.noweb.core.ts b/src/core/engines/noweb/session.noweb.core.ts index c24aefb9..2d7be885 100644 --- a/src/core/engines/noweb/session.noweb.core.ts +++ b/src/core/engines/noweb/session.noweb.core.ts @@ -19,13 +19,20 @@ import makeWASocket, { proto, WAMessageContent, WAMessageKey, + WAMessageUpdate, } from '@adiwajshing/baileys'; import { WACallEvent } from '@adiwajshing/baileys/lib/Types/Call'; +import { BaileysEventMap } from '@adiwajshing/baileys/lib/Types/Events'; +import { GroupMetadata } from '@adiwajshing/baileys/lib/Types/GroupMetadata'; import { Label as NOWEBLabel, LabelActionBody, } from '@adiwajshing/baileys/lib/Types/Label'; -import { LabelAssociationType } from '@adiwajshing/baileys/lib/Types/LabelAssociation'; +import { + ChatLabelAssociation, + LabelAssociationType, +} from '@adiwajshing/baileys/lib/Types/LabelAssociation'; +import { MessageUserReceiptUpdate } from '@adiwajshing/baileys/lib/Types/Message'; import { isLidUser } from '@adiwajshing/baileys/lib/WABinary/jid-utils'; import { Logger as BaileysLogger } from '@adiwajshing/baileys/node_modules/pino'; import { UnprocessableEntityException } from '@nestjs/common'; @@ -58,18 +65,33 @@ import { import { ReplyToMessage } from '@waha/structures/message.dto'; import { PaginationParams } from '@waha/structures/pagination.dto'; import { + EnginePayload, PollVote, PollVotePayload, WAMessageAckBody, } from '@waha/structures/webhooks.dto'; import { LoggerBuilder } from '@waha/utils/logging'; import { sleep, waitUntil } from '@waha/utils/promiseTimeout'; +import { exclude } from '@waha/utils/reactive/ops/exclude'; import { SingleDelayedJobRunner } from '@waha/utils/SingleDelayedJobRunner'; import { SinglePeriodicJobRunner } from '@waha/utils/SinglePeriodicJobRunner'; import * as Buffer from 'buffer'; import { Agent } from 'https'; import * as lodash from 'lodash'; import * as NodeCache from 'node-cache'; +import { + filter, + fromEvent, + identity, + merge, + mergeAll, + mergeMap, + Observable, + partition, + share, + Subject, +} from 'rxjs'; +import { map } from 'rxjs/operators'; import { ChatRequest, @@ -333,7 +355,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { } this.connectStore(); this.listenConnectionEvents(); - this.subscribeEngineEvents(); + this.subscribeEngineEvents2(); this.enableAutoRestart(); } @@ -441,7 +463,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { this.authNOWEBStore = null; } this.status = WAHASessionStatus.STOPPED; - this.events.removeAllListeners(); + this.stopEvents(); await this.end(); await this.store?.close(); @@ -966,6 +988,18 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { }; } + private async toLabelChatAssociation( + association: ChatLabelAssociation, + ): Promise { + const labelData = await this.store.getLabelById(association.labelId); + const label = labelData ? this.toLabel(labelData) : null; + return { + labelId: association.labelId, + chatId: toCusFormat(association.chatId), + label: label, + }; + } + /** * Contacts methods */ @@ -1229,204 +1263,232 @@ export class WhatsappSessionNoWebCore extends WhatsappSession { return await this.sock.newsletterAction(id, 'unmute'); } - /** - * END - Methods for API - */ - subscribeEngineEvents() { + subscribeEngineEvents2() { + // + // All + // + const all$ = new Observable((subscriber) => { + return this.sock.ev.process((events) => { + // iterate over keys + for (const event in events) { + const data = events[event]; + subscriber.next({ event: event, data: data }); + } + }); + }); + this.events2.get(WAHAEvents.ENGINE_EVENT).switch(all$); + // // Messages // - this.sock.ev.on('messages.upsert', async ({ messages }) => { - // Check there's some listeners, because we download media during that process - if ( - this.events.listenerCount(WAHAEvents.MESSAGE) == 0 && - this.events.listenerCount(WAHAEvents.MESSAGE_ANY) == 0 - ) { - this.logger.trace('No listeners for messages, skipping...'); - return; - } + const messagesUpsert$ = fromEvent(this.sock.ev, 'messages.upsert').pipe( + map((event: BaileysEventMap['messages.upsert']) => event.messages), + mergeAll(), + ); + let [messagesFromMe$, messagesFromOthers$] = partition( + messagesUpsert$, + isMine, + ); + messagesFromMe$ = messagesFromMe$.pipe( + mergeMap((msg) => this.processIncomingMessage(msg, true)), + share(), // share it so we don't process twice in message.any + ); + messagesFromOthers$ = messagesFromOthers$.pipe( + mergeMap((msg) => this.processIncomingMessage(msg, true)), + share(), // share it so we don't process twice in message.any + ); + const messagesFromAll$ = merge(messagesFromMe$, messagesFromOthers$); + this.events2.get(WAHAEvents.MESSAGE).switch(messagesFromOthers$); + this.events2.get(WAHAEvents.MESSAGE_ANY).switch(messagesFromAll$); - for (const message of messages) { - const payload = await this.processIncomingMessage(message); - if (!message.key.fromMe) { - this.events.emit(WAHAEvents.MESSAGE, payload); - } - this.events.emit(WAHAEvents.MESSAGE_ANY, payload); - } - }); + // // Message Reactions - this.sock.ev.on('messages.upsert', ({ messages }) => { - const reactions = this.processMessageReaction(messages); - for (const reaction of reactions) { - this.events.emit(WAHAEvents.MESSAGE_REACTION, reaction); - } - }); - // Message Ack - direct messages - this.sock.ev.on('messages.update', (events) => { - events - .filter(isMine) - .filter(isAckUpdateMessageEvent) - .map(this.convertMessageUpdateToMessageAck) - .forEach((payload) => - this.events.emit(WAHAEvents.MESSAGE_ACK, payload), - ); - }); - // Message Ack - groups - this.sock.ev.on('message-receipt.update', (events) => { - events - .filter(isMine) - .map(this.convertMessageReceiptUpdateToMessageAck) - .forEach((payload) => - this.events.emit(WAHAEvents.MESSAGE_ACK, payload), - ); - }); + // + const messageReactions$ = messagesUpsert$.pipe( + map(this.processMessageReaction.bind(this)), + filter(Boolean), + ); + this.events2.get(WAHAEvents.MESSAGE_REACTION).switch(messageReactions$); + + // + // Message Ack + // + const messageUpdates$: Observable = fromEvent( + this.sock.ev, + 'messages.update', + ).pipe( + // @ts-ignore + mergeAll(), + ); + const messageAckDirect$ = messageUpdates$.pipe( + filter(isMine), // ack comes only for MY messages + filter(isAckUpdateMessageEvent), + map(this.convertMessageUpdateToMessageAck.bind(this)), + ); + const messageReceiptUpdate$: Observable = + fromEvent(this.sock.ev, 'message-receipt.update').pipe( + // @ts-ignore + mergeAll(), + ); + + const messageAckGroups$ = messageReceiptUpdate$.pipe( + filter(isMine), // ack comes only for MY messages + map(this.convertMessageReceiptUpdateToMessageAck.bind(this)), + ); + const messageAck$ = merge(messageAckDirect$, messageAckGroups$); + this.events2.get(WAHAEvents.MESSAGE_ACK).switch(messageAck$); // // Other // - this.sock.ev.on('connection.update', (event) => - this.events.emit(WAHAEvents.STATE_CHANGE, event), + this.events2 + .get(WAHAEvents.STATE_CHANGE) + .switch(fromEvent(this.sock.ev, 'connection.update')); + this.events2.get(WAHAEvents.GROUP_JOIN).switch( + // @ts-ignore + fromEvent(this.sock.ev, 'groups.upsert').pipe( + mergeAll(), + ), ); - this.sock.ev.on('groups.upsert', (event) => - this.events.emit(WAHAEvents.GROUP_JOIN, event), - ); - this.sock.ev.on('presence.update', (data) => { - this.events.emit( - WAHAEvents.PRESENCE_UPDATE, - this.toWahaPresences(data.id, data.presences), + this.events2 + .get(WAHAEvents.PRESENCE_UPDATE) + .switch( + fromEvent(this.sock.ev, 'presence.update').pipe( + map((data: any) => this.toWahaPresences(data.id, data.presences)), + ), ); - }); // // Poll votes // - this.sock.ev.on('messages.update', async (messages) => { - for (const message of messages) { - const payload = await this.handleMessagesUpdatePollVote(message); - this.events.emit(WAHAEvents.POLL_VOTE, payload); - } - }); - this.sock.ev.on('messages.upsert', ({ messages }) => { - for (const message of messages) { - const payload = this.handleMessageUpsertPollVoteFailed(message); - this.events.emit(WAHAEvents.POLL_VOTE_FAILED, payload); - } - }); + this.events2 + .get(WAHAEvents.POLL_VOTE) + .switch( + messageUpdates$.pipe( + mergeMap(this.handleMessagesUpdatePollVote.bind(this)), + filter(Boolean), + ), + ); + this.events2 + .get(WAHAEvents.POLL_VOTE_FAILED) + .switch( + messagesUpsert$.pipe( + mergeMap(this.handleMessageUpsertPollVoteFailed.bind(this)), + filter(Boolean), + ), + ); // // Calls // - this.sock.ev.on('call', (calls: WACallEvent[]) => { - calls = lodash.filter(calls, { status: 'offer' }); - for (const call of calls) { - const body = this.toCallData(call); - this.events.emit(WAHAEvents.CALL_RECEIVED, body); - } - }); - this.sock.ev.on('call', (calls: WACallEvent[]) => { - calls = lodash.filter(calls, { status: 'accept' }); - for (const call of calls) { - const body = this.toCallData(call); - this.events.emit(WAHAEvents.CALL_ACCEPTED, body); - } - }); - this.sock.ev.on('call', (calls: WACallEvent[]) => { - const acceptCalls = lodash.filter(calls, { status: 'accept' }); - if (acceptCalls.length > 0) { - // We got two events when accepting calls - reject and accept - // Like for each device - // So if we see accepted call - ignore rejected - return; - } - - calls = lodash.filter(calls, { status: 'reject' }); - for (const call of calls) { - const body = this.toCallData(call); - if (body.isGroup == null) { - // We get two "reject" events, one with null property, ignore it - return; - } - this.events.emit(WAHAEvents.CALL_REJECTED, body); - } - }); + // @ts-ignore + const calls$: Observable = fromEvent(this.sock.ev, 'call'); + const call$ = calls$.pipe(mergeMap(identity)); + this.events2.get(WAHAEvents.CALL_RECEIVED).switch( + call$.pipe( + filter((call: WACallEvent) => call.status === 'offer'), + map(this.toCallData.bind(this)), + ), + ); + this.events2.get(WAHAEvents.CALL_ACCEPTED).switch( + call$.pipe( + filter((call: WACallEvent) => call.status === 'accept'), + map(this.toCallData.bind(this)), + ), + ); + this.events2.get(WAHAEvents.CALL_REJECTED).switch( + calls$.pipe( + // Filter out if there's any "accept" events. + // Meaning it's been accepted on one device, but rejected on another + exclude((calls) => calls.some((call) => call.status === 'accept')), + mergeAll(), + filter((call: WACallEvent) => call.status === 'reject'), + // We get two "reject" events, one with null isGroup property, ignore it + exclude((call: WACallEvent) => call.isGroup == null), + map(this.toCallData.bind(this)), + ), + ); // // Labels // - this.sock.ev.on('labels.edit', (data: NOWEBLabel) => { - if (data.deleted) { - return; - } - const body = this.toLabel(data); - this.events.emit(WAHAEvents.LABEL_UPSERT, body); - }); - this.sock.ev.on('labels.edit', (data: NOWEBLabel) => { - if (!data.deleted) { - return; - } - const body = this.toLabel(data); - this.events.emit(WAHAEvents.LABEL_DELETED, body); - }); - this.sock.ev.on('labels.association', async ({ association, type }) => { - if (type !== 'add') { - return; - } - if (association.type !== LabelAssociationType.Chat) { - return; - } - const labelData = await this.store.getLabelById(association.labelId); - const label = labelData ? this.toLabel(labelData) : null; - const body: LabelChatAssociation = { - labelId: association.labelId, - chatId: toCusFormat(association.chatId), - label: label, - }; - this.events.emit(WAHAEvents.LABEL_CHAT_ADDED, body); - }); - this.sock.ev.on('labels.association', async ({ association, type }) => { - if (type !== 'remove') { - return; - } - if (association.type !== LabelAssociationType.Chat) { - return; - } - const labelData = await this.store.getLabelById(association.labelId); - const label = labelData ? this.toLabel(labelData) : null; - const body: LabelChatAssociation = { - labelId: association.labelId, - chatId: toCusFormat(association.chatId), - label: label, - }; - this.events.emit(WAHAEvents.LABEL_CHAT_DELETED, body); - }); + // @ts-ignore + const labelsEdit$: Observable = fromEvent( + this.sock.ev, + 'labels.edit', + ); + this.events2.get(WAHAEvents.LABEL_UPSERT).switch( + labelsEdit$.pipe( + exclude((data: NOWEBLabel) => data.deleted), + map(this.toLabel.bind(this)), + ), + ); + this.events2.get(WAHAEvents.LABEL_DELETED).switch( + labelsEdit$.pipe( + filter((data: NOWEBLabel) => data.deleted), + map(this.toLabel.bind(this)), + ), + ); + const labelsAssociation$ = fromEvent(this.sock.ev, 'labels.association'); + const labelsAssociationAdd$: Observable = + labelsAssociation$.pipe( + filter(({ type }: any) => type === 'add'), + map((data) => data.association), + filter( + (association: any) => association.type === LabelAssociationType.Chat, + ), + ); + + const labelsAssociationRemove$: Observable = + labelsAssociation$.pipe( + filter(({ type }: any) => type === 'remove'), + map((data) => data.association), + filter( + (association: any) => association.type === LabelAssociationType.Chat, + ), + ); + this.events2 + .get(WAHAEvents.LABEL_CHAT_ADDED) + .switch( + labelsAssociationAdd$.pipe( + mergeMap(this.toLabelChatAssociation.bind(this)), + ), + ); + this.events2 + .get(WAHAEvents.LABEL_CHAT_DELETED) + .switch( + labelsAssociationRemove$.pipe( + mergeMap(this.toLabelChatAssociation.bind(this)), + ), + ); } - private processMessageReaction(messages: any[]): WAMessageReaction[] { - const reactions = []; - for (const message of messages) { - if (!message) return []; - if (!message.message) return []; - if (!message.message.reactionMessage) return []; + /** + * END - Methods for API + */ - const id = buildMessageId(message.key); - const fromToParticipant = getFromToParticipant(message); - const reactionMessage = message.message.reactionMessage; - const messageId = buildMessageId(reactionMessage.key); - const reaction: WAMessageReaction = { - id: id, - timestamp: message.messageTimestamp, - from: toCusFormat(fromToParticipant.from), - fromMe: message.key.fromMe, - to: toCusFormat(fromToParticipant.to), - participant: toCusFormat(fromToParticipant.participant), - reaction: { - text: reactionMessage.text, - messageId: messageId, - }, - }; - reactions.push(reaction); - } - return reactions; + private processMessageReaction(message): WAMessageReaction | null { + if (!message) return null; + if (!message.message) return null; + if (!message.message.reactionMessage) return null; + + const id = buildMessageId(message.key); + const fromToParticipant = getFromToParticipant(message); + const reactionMessage = message.message.reactionMessage; + const messageId = buildMessageId(reactionMessage.key); + const reaction: WAMessageReaction = { + id: id, + timestamp: message.messageTimestamp, + from: toCusFormat(fromToParticipant.from), + fromMe: message.key.fromMe, + to: toCusFormat(fromToParticipant.to), + participant: toCusFormat(fromToParticipant.participant), + reaction: { + text: reactionMessage.text, + messageId: messageId, + }, + }; + return reaction; } private async processIncomingMessage(message, downloadMedia = true) { diff --git a/src/core/engines/webjs/session.webjs.core.ts b/src/core/engines/webjs/session.webjs.core.ts index 0f94b4a9..163f68ec 100644 --- a/src/core/engines/webjs/session.webjs.core.ts +++ b/src/core/engines/webjs/session.webjs.core.ts @@ -68,10 +68,15 @@ import { StatusRequest, TextStatus, } from '@waha/structures/status.dto'; -import { WAMessageRevokedBody } from '@waha/structures/webhooks.dto'; +import { + EnginePayload, + WAMessageRevokedBody, +} from '@waha/structures/webhooks.dto'; import { PaginatorInMemory } from '@waha/utils/Paginator'; import { sleep, waitUntil } from '@waha/utils/promiseTimeout'; import { SingleDelayedJobRunner } from '@waha/utils/SingleDelayedJobRunner'; +import { fromEvent, merge, mergeMap, Observable } from 'rxjs'; +import { map } from 'rxjs/operators'; import { Call, Channel as WEBJSChannel, @@ -242,7 +247,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { this.listenEngineEventsInDebugMode(); } this.listenConnectionEvents(); - this.subscribeEngineEvents(); + this.subscribeEngineEvents2(); } async start() { @@ -258,7 +263,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { async stop() { this.shouldRestart = false; this.status = WAHASessionStatus.STOPPED; - this.events.removeAllListeners(); + this.stopEvents(); this.startDelayedJob.cancel(); await this.end(); } @@ -992,90 +997,146 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { /** * END - Methods for API */ + subscribeEngineEvents2() { + // + // All + // + const events: Observable[] = []; + for (const key in Events) { + const event = Events[key]; + const event$ = fromEvent(this.whatsapp, event); + events.push( + event$.pipe( + map((data) => { + return { + event: event, + data: data, + }; + }), + ), + ); + } + const all$ = merge(...events); + this.events2.get(WAHAEvents.ENGINE_EVENT).switch(all$); - subscribeEngineEvents() { // // Messages // - this.whatsapp.on(Events.MESSAGE_RECEIVED, async (message) => { - // Check there's some listeners, because we download media during that process - if (this.events.listenerCount(WAHAEvents.MESSAGE) == 0) { - this.logger.trace('No listeners for messages, skipping...'); - return; - } - const payload = await this.processIncomingMessage(message); - this.events.emit(WAHAEvents.MESSAGE, payload); - }); - this.whatsapp.on(Events.MESSAGE_CIPHERTEXT, async (message) => { - const payload = await this.processIncomingMessage(message); - this.events.emit(WAHAEvents.MESSAGE_WAITING, payload); - }); - this.whatsapp.on(Events.MESSAGE_REVOKED_EVERYONE, async (after, before) => { - const afterMessage = after ? await this.toWAMessage(after) : null; - const beforeMessage = before ? await this.toWAMessage(before) : null; - const body: WAMessageRevokedBody = { - after: afterMessage, - before: beforeMessage, - }; - this.events.emit(WAHAEvents.MESSAGE_REVOKED, body); - }); - this.whatsapp.on('message_reaction', (message) => { - const payload = this.processMessageReaction(message); - this.events.emit(WAHAEvents.MESSAGE_REACTION, payload); - }); - this.whatsapp.on(Events.MESSAGE_CREATE, async (message) => { - // Check there's some listeners, because we download media during that process - if (this.events.listenerCount(WAHAEvents.MESSAGE_ANY) == 0) { - this.logger.trace('No listeners for messages, skipping...'); - return; - } - const payload = await this.processIncomingMessage(message); - this.events.emit(WAHAEvents.MESSAGE_ANY, payload); - }); - this.whatsapp.on(Events.STATE_CHANGED, (event) => { - this.events.emit(WAHAEvents.STATE_CHANGE, event); - }); - this.whatsapp.on(Events.MESSAGE_ACK, (message) => { - // We do not download media here - const payload = this.toWAMessage(message); - this.events.emit(WAHAEvents.MESSAGE_ACK, payload); - }); + const messageReceived$ = fromEvent(this.whatsapp, Events.MESSAGE_RECEIVED); + const messagesFromOthers$ = messageReceived$.pipe( + mergeMap((msg: any) => this.processIncomingMessage(msg, true)), + ); + this.events2.get(WAHAEvents.MESSAGE).switch(messagesFromOthers$); + + const messageCreate$ = fromEvent(this.whatsapp, Events.MESSAGE_CREATE); + const messagesFromAll$ = messageCreate$.pipe( + mergeMap((msg: any) => this.processIncomingMessage(msg, true)), + ); + this.events2.get(WAHAEvents.MESSAGE_ANY).switch(messagesFromAll$); + + const messageCiphertext$ = fromEvent( + this.whatsapp, + Events.MESSAGE_CIPHERTEXT, + ); + const messagesWaiting$ = messageCiphertext$.pipe( + mergeMap((msg: any) => this.processIncomingMessage(msg, true)), + ); + this.events2.get(WAHAEvents.MESSAGE_WAITING).switch(messagesWaiting$); + + const messageRevoked$ = fromEvent( + this.whatsapp, + Events.MESSAGE_REVOKED_EVERYONE, + (after, before) => { + return { after, before }; + }, + ); + const messagesRevoked$ = messageRevoked$.pipe( + map((event): WAMessageRevokedBody => { + const afterMessage = event.after ? this.toWAMessage(event.after) : null; + const beforeMessage = event.before + ? this.toWAMessage(event.before) + : null; + return { + after: afterMessage, + before: beforeMessage, + }; + }), + ); + this.events2.get(WAHAEvents.MESSAGE_REVOKED).switch(messagesRevoked$); + + const messageReaction$ = fromEvent(this.whatsapp, 'message_reaction'); + const messagesReaction$ = messageReaction$.pipe( + map(this.processMessageReaction.bind(this)), + ); + this.events2.get(WAHAEvents.MESSAGE_REACTION).switch(messagesReaction$); + + const messageAck$ = fromEvent( + this.whatsapp, + Events.MESSAGE_ACK, + (message, ack) => { + return { message, ack }; + }, + ); + const messagesAck$ = messageAck$.pipe( + map((event) => event.message), + map(this.toWAMessage.bind(this)), + ); + this.events2.get(WAHAEvents.MESSAGE_ACK).switch(messagesAck$); + + // + // Others + // + const stateChanged$ = fromEvent(this.whatsapp, Events.STATE_CHANGED); + this.events2.get(WAHAEvents.STATE_CHANGE).switch(stateChanged$); // // Groups // - this.whatsapp.on(Events.GROUP_JOIN, (event) => { - this.events.emit(WAHAEvents.GROUP_JOIN, event); - }); - this.whatsapp.on(Events.GROUP_LEAVE, (event) => { - this.events.emit(WAHAEvents.GROUP_LEAVE, event); - }); + const groupJoin$ = fromEvent(this.whatsapp, Events.GROUP_JOIN); + this.events2.get(WAHAEvents.GROUP_JOIN).switch(groupJoin$); + const groupLeave$ = fromEvent(this.whatsapp, Events.GROUP_LEAVE); + this.events2.get(WAHAEvents.GROUP_LEAVE).switch(groupLeave$); // // Chats // - this.whatsapp.on('chat_archived', (chat, archived, _) => { - const body: ChatArchiveEvent = { - id: chat.id._serialized, - archived: archived, - timestamp: chat.timestamp, - }; - this.events.emit(WAHAEvents.CHAT_ARCHIVE, body); - }); + const chatArchived$ = fromEvent( + this.whatsapp, + 'chat_archived', + (chat, archived, _) => { + return { + chat: chat, + archived: archived, + }; + }, + ); + const chatsArchived$ = chatArchived$.pipe( + map((event) => { + return { + id: event.chat.id._serialized, + archived: event.archived, + timestamp: event.chat.timestamp, + }; + }), + ); + this.events2.get(WAHAEvents.CHAT_ARCHIVE).switch(chatsArchived$); // // Calls // - this.whatsapp.on('call', (call: Call) => { - const body: CallData = { - id: call.id, - from: call.from, - timestamp: call.timestamp, - isVideo: call.isVideo, - isGroup: call.isGroup, - }; - this.events.emit(WAHAEvents.CALL_RECEIVED, body); - }); + const call$ = fromEvent(this.whatsapp, 'call'); + const calls$ = call$.pipe( + map((call: Call) => { + return { + id: call.id, + from: call.from, + timestamp: call.timestamp, + isVideo: call.isVideo, + isGroup: call.isGroup, + }; + }), + ); + this.events2.get(WAHAEvents.CALL_RECEIVED).switch(calls$); } private async processIncomingMessage(message: Message, downloadMedia = true) { @@ -1087,7 +1148,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { this.logger.error(e, e.stack); } } - return await this.toWAMessage(message); + return this.toWAMessage(message); } private processMessageReaction(reaction: Reaction): WAMessageReaction { @@ -1105,10 +1166,10 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { }; } - protected toWAMessage(message: Message): Promise { + protected toWAMessage(message: Message): WAMessage { const replyTo = this.extractReplyTo(message); // @ts-ignore - return Promise.resolve({ + return { id: message.id._serialized, timestamp: message.timestamp, from: message.from, @@ -1129,7 +1190,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession { vCards: message.vCards, replyTo: replyTo, _data: message.rawData, - }); + }; } protected extractReplyTo(message: Message): ReplyToMessage | null { diff --git a/src/core/integrations/webhooks/WebhookConductor.ts b/src/core/integrations/webhooks/WebhookConductor.ts index fd04a98e..3fd826fa 100644 --- a/src/core/integrations/webhooks/WebhookConductor.ts +++ b/src/core/integrations/webhooks/WebhookConductor.ts @@ -1,20 +1,18 @@ +import { populateSessionInfo } from '@waha/core/abc/manager.abc'; import { WhatsappSession } from '@waha/core/abc/session.abc'; import { WebhookSender } from '@waha/core/integrations/webhooks/WebhookSender'; -import { WAHAEvents } from '@waha/structures/enums.dto'; +import { WAHAEvents, WAHAEventsWild } from '@waha/structures/enums.dto'; import { WebhookConfig } from '@waha/structures/webhooks.config.dto'; -import { WAHAWebhook } from '@waha/structures/webhooks.dto'; import { EventWildUnmask } from '@waha/utils/events'; import { LoggerBuilder } from '@waha/utils/logging'; -import { VERSION } from '@waha/version'; import { Logger } from 'pino'; export class WebhookConductor { private logger: Logger; - private eventUnmask: EventWildUnmask; + private eventUnmask = new EventWildUnmask(WAHAEvents, WAHAEventsWild); constructor(protected loggerBuilder: LoggerBuilder) { this.logger = loggerBuilder.child({ name: WebhookConductor.name }); - this.eventUnmask = new EventWildUnmask(WAHAEvents); } protected buildSender(webhookConfig: WebhookConfig): WebhookSender { @@ -44,30 +42,13 @@ export class WebhookConductor { const events = this.getSuitableEvents(webhook.events); const sender = this.buildSender(webhook); for (const event of events) { - session.events.on(event, (data: any) => - this.callWebhook(event, session, data, sender), - ); + const obs$ = session.getEventObservable(event); + obs$.subscribe((payload) => { + const data = populateSessionInfo(event, session)(payload); + sender.send(data); + }); this.logger.debug(`Event '${event}' is enabled for url: ${url}`); } this.logger.info(`Webhooks were configured for ${url}.`); } - - public async callWebhook( - event, - session: WhatsappSession, - data: any, - sender: WebhookSender, - ) { - const me = session.getSessionMeInfo(); - const json: WAHAWebhook = { - event: event, - session: session.name, - metadata: session.sessionConfig?.metadata, - me: me, - payload: data, - engine: session.engine, - environment: VERSION, - }; - sender.send(json); - } } diff --git a/src/core/integrations/webhooks/WebhookSender.ts b/src/core/integrations/webhooks/WebhookSender.ts index e908c014..4d4f8875 100644 --- a/src/core/integrations/webhooks/WebhookSender.ts +++ b/src/core/integrations/webhooks/WebhookSender.ts @@ -9,7 +9,7 @@ import axios, { AxiosInstance } from 'axios'; import axiosRetry, { retryAfter } from 'axios-retry'; import * as crypto from 'crypto'; import { Logger } from 'pino'; -import { v4 as uuid4 } from 'uuid'; +import { ulid } from 'ulid'; // eslint-disable-next-line @typescript-eslint/no-var-requires const HttpAgent = require('agentkeepalive'); @@ -69,11 +69,19 @@ export class WebhookSender { Object.assign(headers, this.getWebhookHeader()); Object.assign(headers, this.getHMACHeaders(body)); this.logger.info( - { id: headers['X-Webhook-Request-Id'], url: this.url }, + { + id: headers['X-Webhook-Request-Id'], + ['event.id']: json.id, + url: this.url, + }, `Sending POST...`, ); this.logger.debug( - { id: headers['X-Webhook-Request-Id'], data: json }, + { + id: headers['X-Webhook-Request-Id'], + data: json, + ['event.id']: json.id, + }, `POST DATA`, ); @@ -81,11 +89,18 @@ export class WebhookSender { .post(this.url, body, { headers: headers }) .then((response) => { this.logger.info( - { id: headers['X-Webhook-Request-Id'] }, + { + id: headers['X-Webhook-Request-Id'], + ['event.id']: json.id, + }, `POST request was sent with status code: ${response.status}`, ); this.logger.debug( - { id: headers['X-Webhook-Request-Id'], body: response.data }, + { + id: headers['X-Webhook-Request-Id'], + ['event.id']: json.id, + body: response.data, + }, `Response`, ); }) @@ -93,6 +108,7 @@ export class WebhookSender { this.logger.error( { id: headers['X-Webhook-Request-Id'], + ['event.id']: json.id, error: error.message, data: error.response?.data, }, @@ -154,7 +170,7 @@ export class WebhookSender { protected getWebhookHeader() { return { // UUID, no '-' in it - 'X-Webhook-Request-Id': uuid4().replace(/-/g, ''), + 'X-Webhook-Request-Id': ulid(), // unix timestamp with ms 'X-Webhook-Timestamp': Date.now().toString(), }; diff --git a/src/core/manager.core.ts b/src/core/manager.core.ts index 4f250276..2d65f757 100644 --- a/src/core/manager.core.ts +++ b/src/core/manager.core.ts @@ -5,10 +5,14 @@ import { } from '@nestjs/common'; import { WebhookConductor } from '@waha/core/integrations/webhooks/WebhookConductor'; import { MediaStorageFactory } from '@waha/core/media/MediaStorageFactory'; +import { DefaultMap } from '@waha/utils/DefaultMap'; import { getPinoLogLevel, LoggerBuilder } from '@waha/utils/logging'; import { promiseTimeout, sleep } from '@waha/utils/promiseTimeout'; -import { EventEmitter } from 'events'; +import { complete } from '@waha/utils/reactive/complete'; +import { SwitchObservable } from '@waha/utils/reactive/SwitchObservable'; import { PinoLogger } from 'nestjs-pino'; +import { Observable } from 'rxjs'; +import { map } from 'rxjs/operators'; import { WhatsappConfigService } from '../config.service'; import { @@ -24,7 +28,7 @@ import { SessionInfo, } from '../structures/sessions.dto'; import { WebhookConfig } from '../structures/webhooks.config.dto'; -import { SessionManager } from './abc/manager.abc'; +import { populateSessionInfo, SessionManager } from './abc/manager.abc'; import { SessionParams, WhatsappSession } from './abc/session.abc'; import { EngineConfigService } from './config/EngineConfigService'; import { WhatsappSessionNoWebCore } from './engines/noweb/session.noweb.core'; @@ -60,6 +64,7 @@ export class SessionManagerCore extends SessionManager { DEFAULT = 'default'; protected readonly EngineClass: typeof WhatsappSession; + protected events2: DefaultMap>; constructor( config: WhatsappConfigService, @@ -68,7 +73,6 @@ export class SessionManagerCore extends SessionManager { private mediaStorageFactory: MediaStorageFactory, ) { super(config, log); - this.events = new EventEmitter(); this.session = DefaultSessionStatus.STOPPED; this.sessionConfig = null; const engineName = this.engineConfigService.getDefaultEngineName(); @@ -98,6 +102,7 @@ export class SessionManagerCore extends SessionManager { } async beforeApplicationShutdown(signal?: string) { + this.stopEvents(); if (!this.session) { return; } @@ -166,17 +171,12 @@ export class SessionManagerCore extends SessionManager { // @ts-ignore const session = new this.EngineClass(sessionConfig); this.session = session; + this.updateSession(); // configure webhooks const webhooks = this.getWebhooks(); webhook.configure(session, webhooks); - // configure events - session.events.on( - WAHAEvents.SESSION_STATUS, - this.handleSessionEvent(WAHAEvents.SESSION_STATUS, session), - ); - // start session await session.start(); logger.info('Session has been started.'); @@ -187,6 +187,24 @@ export class SessionManagerCore extends SessionManager { }; } + private updateSession() { + if (!this.session) { + return; + } + const session: WhatsappSession = this.session as WhatsappSession; + for (const eventName in WAHAEvents) { + const event = WAHAEvents[eventName]; + const stream$ = session + .getEventObservable(event) + .pipe(map(populateSessionInfo(event, session))); + this.events2.get(event).switch(stream$); + } + } + + getSessionEvent(session: string, event: WAHAEvents): Observable { + return this.events2.get(event); + } + async stop(name: string, silent: boolean): Promise { this.onlyDefault(name); if (!this.isRunning(name)) { @@ -206,6 +224,7 @@ export class SessionManagerCore extends SessionManager { } this.log.info(`Session has been stopped.`, { session: name }); this.session = DefaultSessionStatus.STOPPED; + this.updateSession(); await sleep(this.SESSION_STOP_TIMEOUT); } @@ -230,6 +249,7 @@ export class SessionManagerCore extends SessionManager { async delete(name: string): Promise { this.onlyDefault(name); this.session = DefaultSessionStatus.REMOVED; + this.updateSession(); this.sessionConfig = undefined; } @@ -337,4 +357,8 @@ export class SessionManagerCore extends SessionManager { const engine = await this.fetchEngineInfo(); return { ...session, engine: engine }; } + + protected stopEvents() { + complete(this.events2); + } } diff --git a/src/structures/enums.dto.ts b/src/structures/enums.dto.ts index 6e5391ab..6dd5e89c 100644 --- a/src/structures/enums.dto.ts +++ b/src/structures/enums.dto.ts @@ -22,8 +22,14 @@ export enum WAHAEvents { LABEL_DELETED = 'label.deleted', LABEL_CHAT_ADDED = 'label.chat.added', LABEL_CHAT_DELETED = 'label.chat.deleted', + ENGINE_EVENT = 'engine.event', } +// All but no state.change, it's internal one +export const WAHAEventsWild = Object.values(WAHAEvents).filter( + (e) => e !== WAHAEvents.STATE_CHANGE && e !== WAHAEvents.ENGINE_EVENT, +); + export enum WAHASessionStatus { STOPPED = 'STOPPED', STARTING = 'STARTING', @@ -53,4 +59,5 @@ export enum WAMessageAck { READ = 3, PLAYED = 4, } + export const ACK_UNKNOWN = 'UNKNOWN'; diff --git a/src/structures/webhooks.dto.ts b/src/structures/webhooks.dto.ts index 78be36cb..59f938e6 100644 --- a/src/structures/webhooks.dto.ts +++ b/src/structures/webhooks.dto.ts @@ -101,6 +101,11 @@ export class WASessionStatusBody { } export class WAHAWebhook { + @ApiProperty({ + example: 'evt_01jcn4pjwwg47bwy2gsey6q5sx', + }) + id: string; + @ApiProperty({ example: 'default', }) @@ -320,6 +325,20 @@ class WAHAWebhookLabelChatDeleted extends WAHAWebhook { payload: LabelChatAssociation; } +export class EnginePayload { + event: string; + data: any; +} + +class WAHAWebhookEngineEvent extends WAHAWebhook { + @ApiProperty({ + description: 'Internal engine event.', + }) + event = WAHAEvents.ENGINE_EVENT; + + payload: EnginePayload; +} + const WAHA_WEBHOOKS = [ WAHAWebhookSessionStatus, WAHAWebhookMessage, @@ -341,5 +360,6 @@ const WAHA_WEBHOOKS = [ WAHAWebhookLabelDeleted, WAHAWebhookLabelChatAdded, WAHAWebhookLabelChatDeleted, + WAHAWebhookEngineEvent, ]; export { WAHA_WEBHOOKS }; diff --git a/src/utils/DefaultMap.ts b/src/utils/DefaultMap.ts new file mode 100644 index 00000000..7e9dad77 --- /dev/null +++ b/src/utils/DefaultMap.ts @@ -0,0 +1,15 @@ +export class DefaultMap extends Map { + private readonly factory: (key: K) => T; + + constructor(factory: (key: K) => T) { + super(); + this.factory = factory; + } + + get(key: K): T { + if (!this.has(key)) { + this.set(key, this.factory(key)); + } + return super.get(key); + } +} diff --git a/src/utils/events.ts b/src/utils/events.ts index 8f62784d..8bfefbde 100644 --- a/src/utils/events.ts +++ b/src/utils/events.ts @@ -3,24 +3,29 @@ * Remove duplicates if any */ export class EventWildUnmask { - constructor(private readonly events: string[] | any) { + constructor( + private readonly events: string[] | any, + private readonly all: string[] | any | null = null, + ) { // in case of enum - convert to array this.events = Object.values(events); + this.all = all ? Object.values(all) : this.events; } unmask(events: string[]) { + const rightEvents = []; if (events.includes('*')) { - return this.events; + rightEvents.push(...this.all); } // Get only known events, log and ignore others - const rightEvents = []; for (const event of events) { if (!this.events.includes(event)) { continue; } rightEvents.push(event); } - return rightEvents; + // return unique values + return [...new Set(rightEvents)]; } } diff --git a/src/utils/ids.ts b/src/utils/ids.ts index c5271142..a78399cd 100644 --- a/src/utils/ids.ts +++ b/src/utils/ids.ts @@ -1,10 +1,9 @@ -import { v4 as uuid4 } from 'uuid'; +import { ulid } from 'ulid'; /** * Generate prefix uuid (but remove -) * @param prefix */ export function generatePrefixedId(prefix: string) { - const id = uuid4().replace(/-/g, ''); - return `${prefix}_${id}`; + return `${prefix}_${ulid().toLowerCase()}`; } diff --git a/src/utils/reactive/SwitchObservable.ts b/src/utils/reactive/SwitchObservable.ts new file mode 100644 index 00000000..f22a7f5f --- /dev/null +++ b/src/utils/reactive/SwitchObservable.ts @@ -0,0 +1,34 @@ +import { BehaviorSubject, EMPTY, Observable, share, switchMap } from 'rxjs'; + +const noop = (x: any) => x; + +/** + * Observable that can easily be switched to new observable + * So you can change "source" of the observable, + * but clients will use the same observable (but consume NEW data from NEW source) + */ +export class SwitchObservable extends Observable { + private subject$: BehaviorSubject>; + + /** + * fn - Modify inner observable + * + */ + constructor(fn: (source: Observable) => Observable = noop) { + const subject$ = new BehaviorSubject>(EMPTY); + let observable$ = subject$.pipe(switchMap((stream$) => stream$)); + observable$ = fn(observable$); + + super((subscriber) => observable$.subscribe(subscriber)); + this.subject$ = subject$; + } + + switch(newObservable$: Observable) { + this.subject$.next(newObservable$); + } + + complete() { + this.switch(EMPTY); + this.subject$.complete(); + } +} diff --git a/src/utils/reactive/complete.ts b/src/utils/reactive/complete.ts new file mode 100644 index 00000000..cfcb414f --- /dev/null +++ b/src/utils/reactive/complete.ts @@ -0,0 +1,19 @@ +import { Subject } from 'rxjs'; + +interface WithValues { + values(): IterableIterator; +} + +interface Completable { + complete(): void; +} + +/** + * Go over collection and complete each item + * @param collection + */ +export function complete(collection: WithValues) { + for (const item of collection.values()) { + item.complete(); + } +} diff --git a/src/utils/reactive/ops/exclude.ts b/src/utils/reactive/ops/exclude.ts new file mode 100644 index 00000000..277cf355 --- /dev/null +++ b/src/utils/reactive/ops/exclude.ts @@ -0,0 +1,8 @@ +import { filter } from 'rxjs'; + +/** + * Just as filter, but opposite + */ +export function exclude(predicate: (...args) => boolean) { + return filter((...args) => !predicate(...args)); +} diff --git a/yarn.lock b/yarn.lock index 8b615020..3dd4ff54 100644 --- a/yarn.lock +++ b/yarn.lock @@ -10619,7 +10619,7 @@ __metadata: languageName: node linkType: hard -"rxjs@npm:7.8.1, rxjs@npm:^7.1.0, rxjs@npm:^7.5.5": +"rxjs@npm:7.8.1, rxjs@npm:^7.5.5, rxjs@npm:^7.8.1": version: 7.8.1 resolution: "rxjs@npm:7.8.1" dependencies: @@ -12018,6 +12018,15 @@ __metadata: languageName: node linkType: hard +"ulid@npm:^2.3.0": + version: 2.3.0 + resolution: "ulid@npm:2.3.0" + bin: + ulid: ./bin/cli.js + checksum: d6dbf253fdc189f60fe2829d934ee5447b3dab62d05449a2e0fe89670d77087dd6eba4f844a69f9ffdb01384ec6fd97bdd9be638fc67d593569a45e8969f1e69 + languageName: node + linkType: hard + "unbox-primitive@npm:^1.0.2": version: 1.0.2 resolution: "unbox-primitive@npm:1.0.2" @@ -12280,7 +12289,7 @@ __metadata: qrcode-terminal: ^0.12.0 reflect-metadata: ^0.1.13 rimraf: ^3.0.2 - rxjs: ^7.1.0 + rxjs: ^7.8.1 sharp: ^0.33.4 supertest: ^4.0.2 swagger-ui-express: ^4.1.4 @@ -12289,6 +12298,7 @@ __metadata: ts-node: ^10.9.1 tsconfig-paths: ^4.1.0 typescript: ^4.8.4 + ulid: ^2.3.0 whatsapp-web.js: "github:devlikeapro/whatsapp-web.js#fork-main-channels" dependenciesMeta: bufferutil: