[core] Websockets
This commit is contained in:
1 parent
3ad0a463b5
commit
7b4fa4331f
18 files changed
+705
-372
No files matched your search
+2
-1
@@ -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": {
|
||||
|
||||
+37
-20
@@ -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<any> {
|
||||
return of();
|
||||
}
|
||||
|
||||
public getSessionEvents(
|
||||
session: string,
|
||||
events: WAHAEvents[],
|
||||
): Observable<any> {
|
||||
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<WhatsappSession> {
|
||||
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,
|
||||
};
|
||||
};
|
||||
}
|
||||
@@ -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<WAHAEvents, SwitchObservable<any>>;
|
||||
private status$: Subject<WAHASessionStatus>;
|
||||
|
||||
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<WAHAEvents, SwitchObservable<any>>(
|
||||
(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<WAHASessionStatus, WASessionStatusBody>((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<void>;
|
||||
|
||||
protected stopEvents() {
|
||||
complete(this.events2);
|
||||
}
|
||||
|
||||
/* Unpair the account */
|
||||
async unpair(): Promise<void> {
|
||||
return;
|
||||
|
||||
@@ -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<WebSocket, { session: string; events: string[] }> =
|
||||
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));
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -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<LabelChatAssociation> {
|
||||
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<EnginePayload>((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<WAMessageUpdate> = 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<MessageUserReceiptUpdate> =
|
||||
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<GroupMetadata[]>(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<WACallEvent[]> = 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<NOWEBLabel> = 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<ChatLabelAssociation> =
|
||||
labelsAssociation$.pipe(
|
||||
filter(({ type }: any) => type === 'add'),
|
||||
map((data) => data.association),
|
||||
filter(
|
||||
(association: any) => association.type === LabelAssociationType.Chat,
|
||||
),
|
||||
);
|
||||
|
||||
const labelsAssociationRemove$: Observable<ChatLabelAssociation> =
|
||||
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) {
|
||||
|
||||
@@ -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<EnginePayload>[] = [];
|
||||
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<WAMessage> {
|
||||
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 {
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
@@ -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(),
|
||||
};
|
||||
|
||||
@@ -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<WAHAEvents, SwitchObservable<any>>;
|
||||
|
||||
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<any> {
|
||||
return this.events2.get(event);
|
||||
}
|
||||
|
||||
async stop(name: string, silent: boolean): Promise<void> {
|
||||
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<void> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -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';
|
||||
@@ -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 };
|
||||
@@ -0,0 +1,15 @@
|
||||
export class DefaultMap<K, T> extends Map<K, T> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
+9
-4
@@ -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)];
|
||||
}
|
||||
}
|
||||
+2
-3
@@ -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()}`;
|
||||
}
|
||||
@@ -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<T> extends Observable<T> {
|
||||
private subject$: BehaviorSubject<Observable<T>>;
|
||||
|
||||
/**
|
||||
* fn - Modify inner observable
|
||||
*
|
||||
*/
|
||||
constructor(fn: (source: Observable<T>) => Observable<T> = noop) {
|
||||
const subject$ = new BehaviorSubject<Observable<any>>(EMPTY);
|
||||
let observable$ = subject$.pipe(switchMap((stream$) => stream$));
|
||||
observable$ = fn(observable$);
|
||||
|
||||
super((subscriber) => observable$.subscribe(subscriber));
|
||||
this.subject$ = subject$;
|
||||
}
|
||||
|
||||
switch(newObservable$: Observable<T>) {
|
||||
this.subject$.next(newObservable$);
|
||||
}
|
||||
|
||||
complete() {
|
||||
this.switch(EMPTY);
|
||||
this.subject$.complete();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
import { Subject } from 'rxjs';
|
||||
|
||||
interface WithValues<T> {
|
||||
values(): IterableIterator<T>;
|
||||
}
|
||||
|
||||
interface Completable {
|
||||
complete(): void;
|
||||
}
|
||||
|
||||
/**
|
||||
* Go over collection and complete each item
|
||||
* @param collection
|
||||
*/
|
||||
export function complete(collection: WithValues<Completable>) {
|
||||
for (const item of collection.values()) {
|
||||
item.complete();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
import { filter } from 'rxjs';
|
||||
|
||||
/**
|
||||
* Just as filter, but opposite
|
||||
*/
|
||||
export function exclude<T>(predicate: (...args) => boolean) {
|
||||
return filter((...args) => !predicate(...args));
|
||||
}
|
||||
@@ -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:
|
||||
|
||||
Reference in new issue
Block a user