Compare commits

...
18 Commits
Author SHA1 Message Date
devlikepro 913b8606d9 [core] 2024.11.7
Release / amd64 - chrome - chrome (push) Waiting to run
Release / amd64 - chromium - latest (push) Waiting to run
Release / linux/arm64 - chromium - arm (push) Waiting to run
Release / amd64 - none - noweb (push) Waiting to run
Release / linux/arm64 - none - noweb-arm (push) Waiting to run
2024-11-24 14:06:33 +07:00
devlikepro 7cadbad36e [core] Add WebJSEngineConfigService.ts
fix #654 fix #653
2024-11-24 14:06:28 +07:00
devlikepro abb5bdffcc [core] 2024.11.6
Release / amd64 - chrome - chrome (push) Waiting to run
Release / amd64 - chromium - latest (push) Waiting to run
Release / linux/arm64 - chromium - arm (push) Waiting to run
Release / amd64 - none - noweb (push) Waiting to run
Release / linux/arm64 - none - noweb-arm (push) Waiting to run
2024-11-20 16:38:59 +07:00
devlikepro 6a44c15dd3 [core] fix undefined.get
fix #645
2024-11-20 16:38:43 +07:00
devlikepro e669dfe533 [core] Update dashboard
Release / amd64 - chrome - chrome (push) Waiting to run
Release / amd64 - chromium - latest (push) Waiting to run
Release / linux/arm64 - chromium - arm (push) Waiting to run
Release / amd64 - none - noweb (push) Waiting to run
Release / linux/arm64 - none - noweb-arm (push) Waiting to run
2024-11-19 13:17:40 +07:00
devlikepro be888f195d [core] NOWEB - fix timestamp 2024-11-19 13:17:39 +07:00
devlikepro d87cdef311 [core] Add participant for group 2024-11-19 13:17:39 +07:00
devlikepro ecf07cb520 [core] Pin Unpin message
fix #613
2024-11-19 13:17:39 +07:00
devlikepro 10829d53b9 [core] Check quoted message better 2024-11-19 13:17:38 +07:00
devlikepro 71093f2f31 [core] Rm unused code 2024-11-19 13:17:38 +07:00
devlikepro e88192002a [core] Fix log ctx 2024-11-19 13:17:38 +07:00
devlikepro 3e14282603 [core] 2024.11.5 2024-11-19 13:17:37 +07:00
devlikepro 5f16f4ca77 [core] Update dashboard 2024-11-19 13:17:37 +07:00
devlikepro 7b4fa4331f [core] Websockets 2024-11-19 13:17:36 +07:00
devlikepro 3ad0a463b5 [core] Lazy message and message.any 2024-11-19 13:17:36 +07:00
devlikepro d28d640f77 [core] Use EventWildUmask 2024-11-19 13:17:35 +07:00
devlikepro 594c0a6122 [core] Emit events in engine 2024-11-19 13:17:35 +07:00
devlikepro 843a826437 [core] Merge plus webhooks to core 2024-11-19 13:17:35 +07:00
27 changed files with 1228 additions and 682 deletions

No files matched your search

+1 -1
View File
@@ -26,7 +26,7 @@ RUN yarn build && find ./dist -name "*.d.ts" -delete
FROM node:${NODE_VERSION} AS dashboard
# Download WAHA Dashboard
ENV WAHA_DASHBOARD_SHA e6a72c5de2effc97d83e997cb1f361851604bae8
ENV WAHA_DASHBOARD_SHA 45db5b46a944f4320136f2409d8b970483c8bd9d
RUN \
wget https://github.com/devlikeapro/dashboard/archive/${WAHA_DASHBOARD_SHA}.zip \
&& unzip ${WAHA_DASHBOARD_SHA}.zip -d /tmp/dashboard \
+2 -1
View File
@@ -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": {
+28
View File
@@ -27,6 +27,7 @@ import {
GetChatMessageQuery,
GetChatMessagesFilter,
GetChatMessagesQuery,
PinMessageRequest,
} from '../structures/chats.dto';
import { EditMessageRequest } from '../structures/chatting.dto';
@@ -89,6 +90,33 @@ class ChatsController {
return message;
}
@Post(':chatId/messages/:messageId/pin')
@SessionApiParam
@ApiOperation({ summary: 'Pins a message in the chat' })
@ChatIdApiParam
async pinMessage(
@WorkingSessionParam session: WhatsappSession,
@Param('chatId') chatId: string,
@Param('messageId') messageId: string,
@Body() body: PinMessageRequest,
) {
const result = await session.pinMessage(chatId, messageId, body.duration);
return { success: result };
}
@Post(':chatId/messages/:messageId/unpin')
@SessionApiParam
@ApiOperation({ summary: 'Unpins a message in the chat' })
@ChatIdApiParam
async unpinMessage(
@WorkingSessionParam session: WhatsappSession,
@Param('chatId') chatId: string,
@Param('messageId') messageId: string,
) {
const result = await session.unpinMessage(chatId, messageId);
return { success: result };
}
@Delete(':chatId/messages')
@SessionApiParam
@ApiOperation({ summary: 'Clears all messages from the chat' })
+36 -22
View File
@@ -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,
@@ -25,7 +25,6 @@ import {
import { ISessionAuthRepository } from '../storage/ISessionAuthRepository';
import { ISessionConfigRepository } from '../storage/ISessionConfigRepository';
import { WhatsappSession } from './session.abc';
import { WebhookConductor } from './webhooks.abc';
// eslint-disable-next-line @typescript-eslint/no-var-requires
const AsyncLock = require('async-lock');
@@ -36,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;
@@ -72,7 +70,18 @@ export abstract class SessionManager implements BeforeApplicationShutdown {
protected abstract get EngineClass(): typeof WhatsappSession;
protected abstract get WebhookConductorClass(): typeof WebhookConductor;
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
@@ -114,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]);
}
@@ -160,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,
};
};
}
+93 -28
View File
@@ -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 {
@@ -88,10 +103,6 @@ export function ensureSuffix(phone) {
return phone + suffix;
}
export enum WAHAInternalEvent {
ENGINE_START = 'engine.start',
}
export interface SessionParams {
name: string;
printQR: boolean;
@@ -104,7 +115,6 @@ export interface SessionParams {
}
export abstract class WhatsappSession {
public events: EventEmitter;
public engine: WAHAEngine;
public name: string;
@@ -119,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,
@@ -130,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;
@@ -142,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
@@ -149,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() {
@@ -223,30 +291,15 @@ export abstract class WhatsappSession {
/** Stop the session */
abstract stop(): Promise<void>;
protected stopEvents() {
complete(this.events2);
}
/* Unpair the account */
async unpair(): Promise<void> {
return;
}
/** Subscribe the handler to specific hook */
subscribeSessionEvent(
hook: WAHAEvents | string,
handler: (message) => void,
): boolean {
switch (hook) {
case WAHAEvents.SESSION_STATUS:
this.events.on(WAHAEvents.SESSION_STATUS, handler);
return true;
default:
return false;
}
}
abstract subscribeEngineEvent(
hook: WAHAEvents | string,
handler: (message) => void,
): boolean;
/**
* START - Methods for API
*/
@@ -359,6 +412,18 @@ export abstract class WhatsappSession {
throw new NotImplementedByEngineError();
}
public pinMessage(
chatId: string,
messageId: string,
duration: number,
): Promise<boolean> {
throw new NotImplementedByEngineError();
}
public unpinMessage(chatId: string, messageId: string): Promise<boolean> {
throw new NotImplementedByEngineError();
}
public deleteMessage(chatId: string, messageId: string) {
throw new NotImplementedByEngineError();
}
-26
View File
@@ -1,26 +0,0 @@
import { LoggerBuilder } from '@waha/utils/logging';
import { Logger } from 'pino';
import { WebhookConfig } from '../../structures/webhooks.config.dto';
import { WhatsappSession } from './session.abc';
export abstract class WebhookSender {
protected url: string;
protected logger: Logger;
protected readonly config: WebhookConfig;
constructor(
loggerBuilder: LoggerBuilder,
protected webhookConfig: WebhookConfig,
) {
this.url = webhookConfig.url;
this.logger = loggerBuilder.child({ name: WebhookSender.name });
this.config = webhookConfig;
}
abstract send(json: any): void;
}
export abstract class WebhookConductor {
abstract configure(session: WhatsappSession, webhooks: WebhookConfig[]);
}
+23 -50
View File
@@ -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');
@@ -56,49 +54,39 @@ export class WebsocketGatewayCore
const id = generatePrefixedId('wsc');
socket.id = id;
this.logger.debug(`New client connected: ${request.url}`);
const params = this.getParams(id, request, socket);
if (!params) {
return;
}
const { session, events } = params;
const params = this.getParams(request);
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) {
private getParams(request: IncomingMessage) {
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 +106,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));
});
}
}
+2
View File
@@ -10,6 +10,7 @@ import {
ServerDebugController,
} from '@waha/api/server.controller';
import { WebsocketGatewayCore } from '@waha/core/api/websocket.gateway.core';
import { WebJSEngineConfigService } from '@waha/core/config/WebJSEngineConfigService';
import { MediaLocalStorageModule } from '@waha/core/media/local/media.local.storage.module';
import { MediaLocalStorageConfig } from '@waha/core/media/local/MediaLocalStorageConfig';
import { BufferJsonReplacerInterceptor } from '@waha/nestjs/BufferJsonReplacerInterceptor';
@@ -146,6 +147,7 @@ const PROVIDERS = [
},
DashboardConfigServiceCore,
SwaggerConfigServiceCore,
WebJSEngineConfigService,
WhatsappConfigService,
EngineConfigService,
WebsocketGatewayCore,
@@ -0,0 +1,37 @@
import { Injectable } from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import { WebJSConfig } from '../../core/engines/webjs/session.webjs.core';
@Injectable()
export class WebJSEngineConfigService {
constructor(protected configService: ConfigService) {}
getConfig(): WebJSConfig {
let webVersion = this.configService.get<string>(
'WAHA_WEBJS_WEB_VERSION',
undefined,
);
if (webVersion === '2.2412.54-videofix') {
// Deprecated version
webVersion = undefined;
}
return {
webVersion: webVersion,
cacheType: this.getCacheType(),
};
}
getCacheType(): 'local' | 'none' {
const cacheType = this.configService
.get<string>('WAHA_WEBJS_CACHE_TYPE', 'local')
.toLowerCase();
if (cacheType != 'local' && cacheType != 'none') {
throw new Error(
'Invalid cache type, only "local" and "none" are allowed',
);
}
return cacheType;
}
}
+348 -243
View File
@@ -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';
@@ -45,6 +52,7 @@ import {
GetChatMessageQuery,
GetChatMessagesFilter,
GetChatMessagesQuery,
PinDuration,
} from '@waha/structures/chats.dto';
import { SendButtonsRequest } from '@waha/structures/chatting.buttons.dto';
import { ContactQuery, ContactRequest } from '@waha/structures/contacts.dto';
@@ -58,18 +66,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,
@@ -122,7 +145,6 @@ import {
ensureSuffix,
getChannelInviteLink,
isNewsletter,
WAHAInternalEvent,
WhatsappSession,
} from '../../abc/session.abc';
import {
@@ -136,7 +158,7 @@ import { NowebAuthFactoryCore } from './NowebAuthFactoryCore';
import { INowebStore } from './store/INowebStore';
import { NowebPersistentStore } from './store/NowebPersistentStore';
import { NowebStorageFactoryCore } from './store/NowebStorageFactoryCore';
import { extractMediaContent } from './utils';
import { ensureNumber, extractMediaContent } from './utils';
// eslint-disable-next-line @typescript-eslint/no-var-requires
const QRCode = require('qrcode');
@@ -176,10 +198,6 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
private authNOWEBStore: any;
get listenConnectionEventsFromTheStart() {
return true;
}
sock: ReturnType<typeof makeWASocket>;
store: INowebStore;
private qr: QR;
@@ -286,27 +304,33 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return createAgentProxy(this.proxyConfig);
}
async connectStore() {
this.logger.debug(`Connecting store...`);
if (!this.store) {
this.logger.debug(`Making a new store...`);
const storeEnabled = this.sessionConfig?.noweb?.store?.enabled || false;
if (storeEnabled) {
this.logger.debug('Using NowebPersistentStore');
const storage = this.storageFactory.createStorage(
this.sessionStore,
this.name,
);
this.store = new NowebPersistentStore(
this.loggerBuilder.child({ name: NowebPersistentStore.name }),
storage,
);
await this.store.init();
} else {
this.logger.debug('Using NowebInMemoryStore');
this.store = new NowebInMemoryStore();
}
private async ensureStore() {
if (this.store) {
return;
}
this.logger.debug(`Making a new store...`);
const storeEnabled = this.sessionConfig?.noweb?.store?.enabled || false;
if (!storeEnabled) {
this.logger.debug('Using NowebInMemoryStore');
this.store = new NowebInMemoryStore();
return;
}
this.logger.debug('Using NowebPersistentStore');
const storage = this.storageFactory.createStorage(
this.sessionStore,
this.name,
);
this.store = new NowebPersistentStore(
this.loggerBuilder.child({ name: NowebPersistentStore.name }),
storage,
);
await this.store.init();
}
connectStore() {
this.logger.debug(`Connecting store...`);
this.logger.debug(`Binding store to socket...`);
this.store.bind(this.sock.ev, this.sock);
}
@@ -321,17 +345,18 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
this.shouldRestart = true;
// @ts-ignore
this.sock?.ev?.removeAllListeners();
await this.ensureStore();
this.sock = await this.makeSocket();
this.issueMessageUpdateOnEdits();
this.issuePresenceUpdateOnMessageUpsert();
if (this.isDebugEnabled()) {
this.listenEngineEventsInDebugMode();
}
await this.connectStore();
if (this.listenConnectionEventsFromTheStart) {
this.listenConnectionEvents();
this.events.emit(WAHAInternalEvent.ENGINE_START);
}
this.connectStore();
this.listenConnectionEvents();
this.subscribeEngineEvents2();
this.enableAutoRestart();
}
@@ -439,7 +464,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();
@@ -807,6 +832,34 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return await this.processIncomingMessage(message, query.downloadMedia);
}
public async pinMessage(
chatId: string,
messageId: string,
duration: PinDuration,
): Promise<boolean> {
const jid = toJID(chatId);
const key = parseMessageIdSerialized(messageId);
await this.sock.sendMessage(jid, {
pin: key,
type: proto.PinInChat.Type.PIN_FOR_ALL,
time: duration,
});
return true;
}
public async unpinMessage(
chatId: string,
messageId: string,
): Promise<boolean> {
const jid = toJID(chatId);
const key = parseMessageIdSerialized(messageId);
await this.sock.sendMessage(jid, {
pin: key,
type: proto.PinInChat.Type.UNPIN_FOR_ALL,
});
return true;
}
async setReaction(request: MessageReactionRequest) {
const key = parseMessageIdSerialized(request.messageId);
const reactionMessage = {
@@ -964,6 +1017,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
*/
@@ -1227,208 +1292,232 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return await this.sock.newsletterAction(id, 'unmute');
}
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
//
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$);
//
// Message Reactions
//
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.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.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.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
//
// @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
//
// @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)),
),
);
}
/**
* END - Methods for API
*/
subscribeEngineEvent(event, handler): boolean {
switch (event) {
case WAHAEvents.MESSAGE:
this.sock.ev.on('messages.upsert', ({ messages }) => {
this.handleIncomingMessages(messages, handler, false);
});
return true;
case WAHAEvents.MESSAGE_REACTION:
this.sock.ev.on('messages.upsert', ({ messages }) => {
const reactions = this.processMessageReaction(messages);
reactions.map(handler);
});
return true;
case WAHAEvents.MESSAGE_ANY:
this.sock.ev.on('messages.upsert', ({ messages }) =>
this.handleIncomingMessages(messages, handler, true),
);
return true;
case WAHAEvents.MESSAGE_ACK: // Direct message ack
this.sock.ev.on('messages.update', (events) => {
events
.filter(isMine)
.filter(isAckUpdateMessageEvent)
.map(this.convertMessageUpdateToMessageAck)
.forEach(handler);
});
// Group message ack
this.sock.ev.on('message-receipt.update', (events) => {
events
.filter(isMine)
.map(this.convertMessageReceiptUpdateToMessageAck)
.forEach(handler);
});
return true;
case WAHAEvents.STATE_CHANGE:
this.sock.ev.on('connection.update', handler);
return true;
case WAHAEvents.GROUP_JOIN:
this.sock.ev.on('groups.upsert', handler);
return true;
case WAHAEvents.PRESENCE_UPDATE:
this.sock.ev.on('presence.update', (data) =>
handler(this.toWahaPresences(data.id, data.presences)),
);
return true;
case WAHAEvents.POLL_VOTE:
this.sock.ev.on('messages.update', (events) => {
events.forEach((event) =>
this.handleMessagesUpdatePollVote(event, handler),
);
});
return true;
case WAHAEvents.POLL_VOTE_FAILED:
this.sock.ev.on('messages.upsert', ({ messages }) => {
messages.forEach((message) =>
this.handleMessageUpsertPollVoteFailed(message, handler),
);
});
return true;
case WAHAEvents.CALL_RECEIVED:
this.sock.ev.on('call', (calls: WACallEvent[]) => {
calls = lodash.filter(calls, { status: 'offer' });
for (const call of calls) {
const body = this.toCallData(call);
handler(body);
}
});
return true;
case WAHAEvents.CALL_ACCEPTED:
this.sock.ev.on('call', (calls: WACallEvent[]) => {
calls = lodash.filter(calls, { status: 'accept' });
for (const call of calls) {
const body = this.toCallData(call);
handler(body);
}
});
return true;
case WAHAEvents.CALL_REJECTED:
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;
}
private processMessageReaction(message): WAMessageReaction | null {
if (!message) return null;
if (!message.message) return null;
if (!message.message.reactionMessage) return null;
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;
}
handler(body);
}
});
return true;
case WAHAEvents.LABEL_UPSERT:
this.sock.ev.on('labels.edit', (data: NOWEBLabel) => {
if (data.deleted) {
return;
}
const body = this.toLabel(data);
handler(body);
});
return true;
case WAHAEvents.LABEL_DELETED:
this.sock.ev.on('labels.edit', (data: NOWEBLabel) => {
if (!data.deleted) {
return;
}
const body = this.toLabel(data);
handler(body);
});
return true;
case WAHAEvents.LABEL_CHAT_ADDED:
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,
};
handler(body);
});
return true;
case WAHAEvents.LABEL_CHAT_DELETED:
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,
};
handler(body);
});
return true;
default:
return false;
}
}
private handleIncomingMessages(messages, handler, includeFromMe) {
for (const message of messages) {
// Do not include my messages
if (!includeFromMe && message.key.fromMe) {
continue;
}
this.processIncomingMessage(message).then((msg) => {
if (!msg) {
return;
}
handler(msg);
});
}
}
private processMessageReaction(messages: any[]): WAMessageReaction[] {
const reactions = [];
for (const message of messages) {
if (!message) return [];
if (!message.message) return [];
if (!message.message.reactionMessage) return [];
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;
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: ensureNumber(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) {
@@ -1467,7 +1556,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
const ack = message.ack || message.status - 1;
return Promise.resolve({
id: id,
timestamp: message.messageTimestamp,
timestamp: ensureNumber(message.messageTimestamp),
from: toCusFormat(fromToParticipant.from),
fromMe: message.key.fromMe,
body: body,
@@ -1518,6 +1607,9 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return null;
}
const quotedMessage = contextInfo.quotedMessage;
if (!quotedMessage) {
return null;
}
const body = this.extractBody(quotedMessage);
return {
id: contextInfo.stanzaId,
@@ -1577,7 +1669,7 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
return body;
}
protected async handleMessagesUpdatePollVote(event, handler) {
protected async handleMessagesUpdatePollVote(event) {
const { key, update } = event;
const pollUpdates = update?.pollUpdates;
if (!pollUpdates) {
@@ -1608,17 +1700,17 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
const pollVote: PollVote = {
...voteDestination,
selectedOptions: selectedOptions,
timestamp: pollUpdate.senderTimestampMs,
timestamp: ensureNumber(pollUpdate.senderTimestampMs),
};
const payload: PollVotePayload = {
vote: pollVote,
poll: getDestination(pollCreationMessageKey),
};
handler(payload);
return payload;
}
}
protected async handleMessageUpsertPollVoteFailed(message, handler) {
protected async handleMessageUpsertPollVoteFailed(message) {
const pollUpdateMessage = message.message?.pollUpdateMessage;
if (!pollUpdateMessage) {
return;
@@ -1640,13 +1732,13 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
// https://github.com/WhiskeySockets/Baileys/pull/348
// Or without toNumber() - it depends on the PR above
// timestamp: pollUpdateMessage.senderTimestampMs.toNumber()
timestamp: message.messageTimestamp,
timestamp: ensureNumber(message.messageTimestamp),
};
const payload: PollVotePayload = {
vote: pollVote,
poll: getDestination(pollCreationMessageKey),
};
handler(payload);
return payload;
}
private toCallData(call: WACallEvent): CallData {
@@ -1843,9 +1935,13 @@ export function toJID(chatId) {
* {id: "AAA", remoteJid: "11111111111@s.whatsapp.net", "fromMe": false}
* false_11111111111@c.us_AA
*/
function buildMessageId({ id, remoteJid, fromMe }) {
function buildMessageId({ id, remoteJid, fromMe, participant }: WAMessageKey) {
const chatId = toCusFormat(remoteJid);
return `${fromMe || false}_${chatId}_${id}`;
const parts = [fromMe || false, chatId, id];
if (participant) {
parts.push(toCusFormat(participant));
}
return parts.join('_');
}
/**
@@ -1853,22 +1949,31 @@ function buildMessageId({ id, remoteJid, fromMe }) {
* false_11111111111@c.us_AAA
* {id: "AAA", remoteJid: "11111111111@s.whatsapp.net", "fromMe": false}
*/
function parseMessageIdSerialized(messageId: string, soft: boolean = false) {
function parseMessageIdSerialized(
messageId: string,
soft: boolean = false,
): WAMessageKey {
if (!messageId.includes('_') && soft) {
return { id: messageId };
}
const parts = messageId.split('_');
if (parts.length != 3) {
if (parts.length != 3 && parts.length != 4) {
throw new Error(
'Message id be in format false_11111111111@c.us_AAAAAAAAAAAAAAAAAAAA',
'Message id be in format false_11111111111@c.us_AAAAAAAAAAAAAAAAAAAA[_participant]',
);
}
const fromMe = parts[0] == 'true';
const chatId = parts[1];
const remoteJid = toJID(chatId);
const id = parts[2];
return { fromMe: fromMe, id: id, remoteJid: remoteJid };
const participant = parts[3] ? toJID(parts[3]) : undefined;
return {
fromMe: fromMe,
id: id,
remoteJid: remoteJid,
participant: participant,
};
}
function getId(object) {
+8
View File
@@ -78,6 +78,14 @@ const isObjectALong = (value: any): value is Long => {
);
};
export function ensureNumber(value: number | Long): number {
if (!value) {
// @ts-ignore
return value;
}
return typeof value === 'number' ? value : toNumber(value);
}
const toNumber = (longValue: Long): number => {
const { low, high, unsigned } = longValue;
const result = unsigned ? low >>> 0 : low + high * 0x100000000;
+174 -90
View File
@@ -1,10 +1,8 @@
import { UnprocessableEntityException } from '@nestjs/common';
import {
getChannelInviteLink,
WAHAInternalEvent,
WhatsappSession,
} from '@waha/core/abc/session.abc';
import { toJID } from '@waha/core/engines/noweb/session.noweb.core';
import { LocalAuth } from '@waha/core/engines/webjs/LocalAuth';
import { WebjsClient } from '@waha/core/engines/webjs/WebjsClient';
import {
@@ -70,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,
@@ -96,6 +99,7 @@ const QRCode = require('qrcode');
export interface WebJSConfig {
webVersion?: string;
cacheType: 'local' | 'none';
}
export class WhatsappSessionWebJSCore extends WhatsappSession {
@@ -140,7 +144,11 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
const path = this.getClassDirName();
const webVersion =
this.engineConfig?.webVersion || '2.3000.1018072227-alpha';
this.logger.info(`Using web version: '${webVersion}'`);
const cacheType = this.engineConfig?.cacheType || 'local';
this.logger.info(`Using cache type: '${cacheType}'`);
if (cacheType === 'local') {
this.logger.info(`Using web version: '${webVersion}'`);
}
return {
puppeteer: {
headless: true,
@@ -150,8 +158,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
},
webVersion: webVersion,
webVersionCache: {
// type: 'none',
type: 'local',
type: cacheType,
path: path,
strict: true,
},
@@ -244,7 +251,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
this.listenEngineEventsInDebugMode();
}
this.listenConnectionEvents();
this.events.emit(WAHAInternalEvent.ENGINE_START);
this.subscribeEngineEvents2();
}
async start() {
@@ -260,7 +267,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();
}
@@ -575,7 +582,6 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
);
}
const limit = query.limit;
const downloadMedia = query.downloadMedia;
// Test there's chat with id
await this.whatsapp.getChatById(this.ensureSuffix(chatId));
@@ -604,6 +610,23 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
return await this.processIncomingMessage(message, query.downloadMedia);
}
public async pinMessage(
chatId: string,
messageId: string,
duration: number,
): Promise<boolean> {
const message = await this.whatsapp.getMessageById(messageId);
return message.pin(duration);
}
public async unpinMessage(
chatId: string,
messageId: string,
): Promise<boolean> {
const message = await this.whatsapp.getMessageById(messageId);
return message.unpin();
}
async deleteChat(chatId) {
const chat = await this.whatsapp.getChatById(chatId);
return chat.delete();
@@ -994,85 +1017,146 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
/**
* END - Methods for API
*/
subscribeEngineEvent(event, handler): boolean {
switch (event) {
case WAHAEvents.MESSAGE:
this.whatsapp.on(Events.MESSAGE_RECEIVED, (message) =>
this.processIncomingMessage(message).then(handler),
);
return true;
case WAHAEvents.MESSAGE_WAITING:
this.whatsapp.on(Events.MESSAGE_CIPHERTEXT, (message) =>
this.processIncomingMessage(message).then(handler),
);
return true;
case WAHAEvents.MESSAGE_REVOKED:
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,
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,
};
handler(body);
},
);
return true;
case WAHAEvents.MESSAGE_REACTION:
this.whatsapp.on('message_reaction', (message) =>
handler(this.processMessageReaction(message)),
);
return true;
case WAHAEvents.MESSAGE_ANY:
this.whatsapp.on(Events.MESSAGE_CREATE, (message) =>
this.processIncomingMessage(message).then(handler),
);
return true;
case WAHAEvents.STATE_CHANGE:
this.whatsapp.on(Events.STATE_CHANGED, handler);
return true;
case WAHAEvents.MESSAGE_ACK:
// We do not download media here
this.whatsapp.on(Events.MESSAGE_ACK, (message) =>
this.toWAMessage(message).then(handler),
);
return true;
case WAHAEvents.GROUP_JOIN:
this.whatsapp.on(Events.GROUP_JOIN, handler);
return true;
case WAHAEvents.GROUP_LEAVE:
this.whatsapp.on(Events.GROUP_LEAVE, handler);
return true;
case WAHAEvents.CHAT_ARCHIVE:
this.whatsapp.on('chat_archived', (chat, archived, _) => {
const body: ChatArchiveEvent = {
id: chat.id._serialized,
archived: archived,
timestamp: chat.timestamp,
};
handler(body);
});
return true;
case WAHAEvents.CALL_RECEIVED:
this.whatsapp.on('call', (call: Call) => {
const body: CallData = {
id: call.id,
from: call.from,
timestamp: call.timestamp,
isVideo: call.isVideo,
isGroup: call.isGroup,
};
handler(body);
});
return true;
default:
return false;
}),
),
);
}
const all$ = merge(...events);
this.events2.get(WAHAEvents.ENGINE_EVENT).switch(all$);
//
// Messages
//
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
//
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
//
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
//
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) {
@@ -1084,7 +1168,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 {
@@ -1102,10 +1186,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,
@@ -1126,7 +1210,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
vCards: message.vCards,
replyTo: replyTo,
_data: message.rawData,
});
};
}
protected extractReplyTo(message: Message): ReplyToMessage | null {
@@ -0,0 +1,54 @@
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, WAHAEventsWild } from '@waha/structures/enums.dto';
import { WebhookConfig } from '@waha/structures/webhooks.config.dto';
import { EventWildUnmask } from '@waha/utils/events';
import { LoggerBuilder } from '@waha/utils/logging';
import { Logger } from 'pino';
export class WebhookConductor {
private logger: Logger;
private eventUnmask = new EventWildUnmask(WAHAEvents, WAHAEventsWild);
constructor(protected loggerBuilder: LoggerBuilder) {
this.logger = loggerBuilder.child({ name: WebhookConductor.name });
}
protected buildSender(webhookConfig: WebhookConfig): WebhookSender {
return new WebhookSender(this.loggerBuilder, webhookConfig);
}
private getSuitableEvents(events: WAHAEvents[] | string[]): WAHAEvents[] {
return this.eventUnmask.unmask(events);
}
public configure(session: WhatsappSession, webhooks: WebhookConfig[]) {
for (const webhookConfig of webhooks) {
this.configureSingleWebhook(session, webhookConfig);
}
}
private configureSingleWebhook(
session: WhatsappSession,
webhook: WebhookConfig,
) {
if (!webhook || !webhook.url || webhook.events.length === 0) {
return;
}
const url = webhook.url;
this.logger.info(`Configuring webhooks for ${url}...`);
const events = this.getSuitableEvents(webhook.events);
const sender = this.buildSender(webhook);
for (const event of events) {
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}.`);
}
}
@@ -0,0 +1,203 @@
import { SECOND } from '@waha/structures/enums.dto';
import {
RetryPolicy,
WebhookConfig,
} from '@waha/structures/webhooks.config.dto';
import { LoggerBuilder } from '@waha/utils/logging';
import { VERSION } from '@waha/version';
import axios, { AxiosInstance } from 'axios';
import axiosRetry, { retryAfter } from 'axios-retry';
import * as crypto from 'crypto';
import { Logger } from 'pino';
import { ulid } from 'ulid';
// eslint-disable-next-line @typescript-eslint/no-var-requires
const HttpAgent = require('agentkeepalive');
// eslint-disable-next-line @typescript-eslint/no-var-requires
const HttpsAgent = require('agentkeepalive').HttpsAgent;
const DEFAULT_RETRY_DELAY_SECONDS = 2;
const DEFAULT_RETRY_ATTEMPTS = 15;
const DEFAULT_HMAC_ALGORITHM = 'sha512';
function noDelay(_retryNumber = 0, error: any) {
return Math.max(0, retryAfter(error));
}
function constantDelay(delayFactor: number) {
return (_retryNumber = 0, error = undefined) => {
return Math.max(delayFactor, retryAfter(error));
};
}
export function exponentialDelay(delayFactor: number) {
return (retryNumber = 0, error = undefined) => {
const calculatedDelay = 2 ** retryNumber * delayFactor;
const delay = Math.max(calculatedDelay, retryAfter(error));
const randomSum = delay * 0.2 * Math.random(); // 0-20% of the delay
return delay + randomSum;
};
}
export class WebhookSender {
protected static AGENTS = {
http: new HttpAgent({}),
https: new HttpsAgent({}),
};
protected url: string;
protected logger: Logger;
protected readonly config: WebhookConfig;
protected axios: AxiosInstance;
constructor(
loggerBuilder: LoggerBuilder,
protected webhookConfig: WebhookConfig,
) {
this.url = webhookConfig.url;
this.logger = loggerBuilder.child({ name: WebhookSender.name });
this.config = webhookConfig;
this.axios = this.buildAxiosInstance();
}
send(json: any) {
const body = JSON.stringify(json);
const headers = {
'content-type': 'application/json',
};
Object.assign(headers, this.getWebhookHeader());
Object.assign(headers, this.getHMACHeaders(body));
const ctx = {
id: headers['X-Webhook-Request-Id'],
['event.id']: json.id,
url: this.url,
};
this.logger.info(ctx, `Sending POST...`);
this.logger.debug(ctx, `POST DATA`);
this.axios
.post(this.url, body, { headers: headers })
.then((response) => {
this.logger.info(
ctx,
`POST request was sent with status code: ${response.status}`,
);
this.logger.debug(
{
...ctx,
body: response.data,
},
`Response`,
);
})
.catch((error) => {
this.logger.error(
{
...ctx,
error: error.message,
data: error.response?.data,
},
`POST request failed: ${error.message}`,
);
});
}
protected buildAxiosInstance(): AxiosInstance {
// configure headers
const customHeaders = this.config.customHeaders || [];
const headers = {
'content-type': 'application/json',
'User-Agent': `WAHA/${VERSION.version}`,
};
customHeaders.forEach((header) => {
headers[header.name] = header.value;
});
// configure retry
const attempts = this.config.retries?.attempts ?? DEFAULT_RETRY_ATTEMPTS;
const delaySeconds =
this.config.retries?.delaySeconds ?? DEFAULT_RETRY_DELAY_SECONDS;
const delayMs = delaySeconds * SECOND;
const policy = this.config.retries?.policy;
const retryDelay = this.buildRetryDelay(policy, delayMs);
const instance = axios.create({
headers: headers,
httpAgent: WebhookSender.AGENTS.http,
httpsAgent: WebhookSender.AGENTS.https,
});
axiosRetry(instance, {
retries: attempts,
retryDelay: retryDelay,
retryCondition: (error) => true,
onRetry: (retryCount, error, requestConfig) => {
this.logger.warn(
{ id: requestConfig.headers['X-Webhook-Request-Id'] },
`Error sending POST request: '${error.message}'. Retrying ${retryCount}/${attempts}...`,
);
},
});
return instance;
}
protected getHMACHeaders(body: string) {
// HMAC
const hmac = this.calculateHmac(body, DEFAULT_HMAC_ALGORITHM);
if (!hmac) {
return {};
}
return {
'X-Webhook-Hmac': hmac,
'X-Webhook-Hmac-Algorithm': DEFAULT_HMAC_ALGORITHM,
};
}
protected getWebhookHeader() {
return {
// UUID, no '-' in it
'X-Webhook-Request-Id': ulid(),
// unix timestamp with ms
'X-Webhook-Timestamp': Date.now().toString(),
};
}
private calculateHmac(body, algorithm) {
if (!this.config.hmac || !this.config.hmac.key) {
return undefined;
}
return crypto
.createHmac(algorithm, this.config.hmac.key)
.update(body)
.digest('hex');
}
private buildRetryDelay(
policy: RetryPolicy | null,
ms: number,
): (retryNumber: number, error: any) => number {
if (!ms) {
this.logger.debug(`Using no delay, because delaySeconds set to 0`);
return noDelay;
}
switch (policy) {
case RetryPolicy.CONSTANT:
this.logger.debug(`Using constant delay with '${ms}' ms factor`);
return constantDelay(ms);
case RetryPolicy.LINEAR:
this.logger.debug(`Using linear delay with '${ms}' ms factor`);
return axiosRetry.linearDelay(ms);
case RetryPolicy.EXPONENTIAL:
this.logger.debug(`Using exponential delay with '${ms}' ms factor`);
return exponentialDelay(ms);
default:
this.logger.debug('No delay policy specified, using constant delay');
return constantDelay(ms);
}
}
}
+53 -18
View File
@@ -3,11 +3,17 @@ import {
NotFoundException,
UnprocessableEntityException,
} from '@nestjs/common';
import { WebJSEngineConfigService } from '@waha/core/config/WebJSEngineConfigService';
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, retry, share } from 'rxjs';
import { map } from 'rxjs/operators';
import { WhatsappConfigService } from '../config.service';
import {
@@ -23,7 +29,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';
@@ -33,7 +39,6 @@ import { getProxyConfig } from './helpers.proxy';
import { MediaManager } from './media/MediaManager';
import { LocalSessionAuthRepository } from './storage/LocalSessionAuthRepository';
import { LocalStoreCore } from './storage/LocalStoreCore';
import { WebhookConductorCore } from './webhooks.core';
export class OnlyDefaultSessionIsAllowed extends UnprocessableEntityException {
constructor() {
@@ -59,22 +64,29 @@ export class SessionManagerCore extends SessionManager {
private sessionConfig?: SessionConfig;
DEFAULT = 'default';
// @ts-ignore
protected WebhookConductorClass = WebhookConductorCore;
protected readonly EngineClass: typeof WhatsappSession;
protected events2: DefaultMap<WAHAEvents, SwitchObservable<any>>;
constructor(
config: WhatsappConfigService,
private engineConfigService: EngineConfigService,
private webjsEngineConfigService: WebJSEngineConfigService,
log: PinoLogger,
private mediaStorageFactory: MediaStorageFactory,
) {
super(config, log);
this.events = new EventEmitter();
this.session = DefaultSessionStatus.STOPPED;
this.sessionConfig = null;
const engineName = this.engineConfigService.getDefaultEngineName();
this.EngineClass = this.getEngine(engineName);
this.events2 = new DefaultMap<WAHAEvents, SwitchObservable<any>>(
(key) =>
new SwitchObservable((obs$) => {
return obs$.pipe(retry(), share());
}),
);
this.store = new LocalStoreCore(engineName.toLowerCase());
this.sessionAuthRepository = new LocalSessionAuthRepository(this.store);
this.startPredefinedSessions();
@@ -100,6 +112,7 @@ export class SessionManagerCore extends SessionManager {
}
async beforeApplicationShutdown(signal?: string) {
this.stopEvents();
if (!this.session) {
return;
}
@@ -138,7 +151,7 @@ export class SessionManagerCore extends SessionManager {
`Session '${this.DEFAULT}' is already started.`,
);
}
this.log.info(`'${name}' - starting session...`);
this.log.info({ session: name }, `Starting session...`);
const logger = this.log.logger.child({ session: name });
logger.level = getPinoLogLevel(this.sessionConfig?.debug);
const loggerBuilder: LoggerBuilder = logger;
@@ -153,7 +166,7 @@ export class SessionManagerCore extends SessionManager {
loggerBuilder.child({ name: 'MediaManager' }),
);
const webhook = new this.WebhookConductorClass(loggerBuilder);
const webhook = new WebhookConductor(loggerBuilder);
const proxyConfig = this.getProxyConfig();
const sessionConfig: SessionParams = {
name,
@@ -164,21 +177,19 @@ export class SessionManagerCore extends SessionManager {
proxyConfig: proxyConfig,
sessionConfig: this.sessionConfig,
};
if (this.EngineClass === WhatsappSessionWebJSCore) {
sessionConfig.engineConfig = this.webjsEngineConfigService.getConfig();
}
await this.sessionAuthRepository.init(name);
// @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.');
@@ -189,14 +200,32 @@ 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)) {
this.log.debug(`Session is not running.`, { session: name });
this.log.debug({ session: name }, `Session is not running.`);
return;
}
this.log.info(`Stopping session...`, { session: name });
this.log.info({ session: name }, `Stopping session...`);
try {
const session = this.getSession(name);
await session.stop();
@@ -206,8 +235,9 @@ export class SessionManagerCore extends SessionManager {
throw err;
}
}
this.log.info(`Session has been stopped.`, { session: name });
this.log.info({ session: name }, `Session has been stopped.`);
this.session = DefaultSessionStatus.STOPPED;
this.updateSession();
await sleep(this.SESSION_STOP_TIMEOUT);
}
@@ -217,7 +247,7 @@ export class SessionManagerCore extends SessionManager {
}
const session = this.session as WhatsappSession;
this.log.info('Unpairing the device from account...', { session: name });
this.log.info({ session: name }, 'Unpairing the device from account...');
await session.unpair().catch((err) => {
this.log.warn(`Error while unpairing from device: ${err}`);
});
@@ -232,6 +262,7 @@ export class SessionManagerCore extends SessionManager {
async delete(name: string): Promise<void> {
this.onlyDefault(name);
this.session = DefaultSessionStatus.REMOVED;
this.updateSession();
this.sessionConfig = undefined;
}
@@ -339,4 +370,8 @@ export class SessionManagerCore extends SessionManager {
const engine = await this.fetchEngineInfo();
return { ...session, engine: engine };
}
protected stopEvents() {
complete(this.events2);
}
}
-197
View File
@@ -1,197 +0,0 @@
import { LoggerBuilder } from '@waha/utils/logging';
import Agent from 'agentkeepalive';
import axios from 'axios';
import { AxiosInstance } from 'axios';
import { Logger } from 'pino';
import { WAHAEvents } from '../structures/enums.dto';
import { WebhookConfig } from '../structures/webhooks.config.dto';
import { WAHAWebhook } from '../structures/webhooks.dto';
import { VERSION } from '../version';
import { WAHAInternalEvent, WhatsappSession } from './abc/session.abc';
import { WebhookConductor, WebhookSender } from './abc/webhooks.abc';
// eslint-disable-next-line @typescript-eslint/no-var-requires
const HttpAgent = require('agentkeepalive');
// eslint-disable-next-line @typescript-eslint/no-var-requires
const HttpsAgent = require('agentkeepalive').HttpsAgent;
export class WebhookSenderCore extends WebhookSender {
protected static AGENTS = {
http: new HttpAgent({}),
https: new HttpsAgent({}),
};
protected axios: AxiosInstance;
constructor(loggerBuilder: LoggerBuilder, webhookConfig: WebhookConfig) {
super(loggerBuilder, webhookConfig);
this.axios = this.buildAxiosInstance();
}
protected buildAxiosInstance() {
const headers = {
'content-type': 'application/json',
'User-Agent': `WAHA/${VERSION.version}`,
};
return axios.create({
headers: headers,
httpAgent: WebhookSenderCore.AGENTS.http,
httpsAgent: WebhookSenderCore.AGENTS.https,
});
}
send(json: any, headers?: Record<string, string>) {
headers = headers || {};
this.logger.info(
{ id: headers['X-Webhook-Request-Id'], url: this.url },
`Sending POST...`,
);
this.logger.debug(
{ id: headers['X-Webhook-Request-Id'], data: json },
`POST DATA`,
);
this.axios
.post(this.url, json, { headers: headers })
.then((response) => {
this.logger.info(
{ id: headers['X-Webhook-Request-Id'] },
`POST request was sent with status code: ${response.status}`,
);
this.logger.debug(
{ id: headers['X-Webhook-Request-Id'], body: response.data },
`Response`,
);
})
.catch((error) => {
this.logger.error(
{
id: headers['X-Webhook-Request-Id'],
error: error.message,
data: error.response?.data,
},
`POST request failed: ${error.message}`,
);
});
}
}
export class WebhookConductorCore implements WebhookConductor {
private logger: Logger;
constructor(protected loggerBuilder: LoggerBuilder) {
this.logger = loggerBuilder.child({ name: WebhookConductor.name });
}
protected buildSender(webhookConfig: WebhookConfig): WebhookSender {
return new WebhookSenderCore(this.loggerBuilder, webhookConfig);
}
private getSuitableEvents(events: WAHAEvents[] | string[]): WAHAEvents[] {
const allEvents = Object.values(WAHAEvents);
// Enable all events if * in the events
// @ts-ignore
if (events.includes('*')) {
return allEvents;
}
// Get only known events, log and ignore others
const rightEvents = [];
for (const event of events) {
// @ts-ignore
if (!allEvents.includes(event)) {
this.logger.error(`Unknown event for webhook: '${event}'`);
continue;
}
rightEvents.push(event);
}
return rightEvents;
}
public configure(session: WhatsappSession, webhooks: WebhookConfig[]) {
for (const webhookConfig of webhooks) {
this.configureSingleWebhook(session, webhookConfig);
}
}
private configureSingleWebhook(
session: WhatsappSession,
webhook: WebhookConfig,
) {
if (!webhook || !webhook.url || webhook.events.length === 0) {
return;
}
const url = webhook.url;
this.logger.info(`Configuring webhooks for ${url}...`);
const events = this.getSuitableEvents(webhook.events);
const sender = this.buildSender(webhook);
for (const event of events) {
const found = this.configureSingleEvent(
session.subscribeSessionEvent,
session,
event,
sender,
url,
);
if (!found) {
// Postpone for ENGINE_START event and configure engine events
session.events.on(WAHAInternalEvent.ENGINE_START, () => {
const found = this.configureSingleEvent(
session.subscribeEngineEvent,
session,
event,
sender,
url,
);
if (!found) {
this.logger.error(
`Engine does not support webhook event: '${event}' for url '${url}'`,
);
}
});
}
}
this.logger.info(`Webhooks were configured for ${url}.`);
}
private configureSingleEvent(
subscribeMethod,
session: WhatsappSession,
event: WAHAEvents,
sender: WebhookSender,
url: string,
) {
const found = subscribeMethod.apply(session, [
event,
(data: any) => this.callWebhook(event, session, data, sender),
]);
if (!found) {
return false;
}
this.logger.debug(`Event '${event}' is enabled for url: ${url}`);
return true;
}
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);
}
}
+17
View File
@@ -88,6 +88,23 @@ export class ChatsPaginationParams extends PaginationParams {
sortBy?: string;
}
export enum PinDuration {
DAY = 86400,
WEEK = 604800,
MONTH = 2592000,
}
export class PinMessageRequest {
@IsNumber()
@IsEnum(PinDuration)
@ApiProperty({
description:
'Duration in seconds. 24 hours (86400), 7 days (604800), 30 days (2592000)',
example: 86400,
})
duration: number;
}
/**
* Events
*/
+7
View File
@@ -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';
+20
View File
@@ -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 };
+15
View File
@@ -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);
}
}
+31
View File
@@ -0,0 +1,31 @@
/**
* Unmask "*" in events list to exact values
* Remove duplicates if any
*/
export class EventWildUnmask {
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('*')) {
rightEvents.push(...this.all);
}
// Get only known events, log and ignore others
for (const event of events) {
if (!this.events.includes(event)) {
continue;
}
rightEvents.push(event);
}
// return unique values
return [...new Set(rightEvents)];
}
}
+2 -3
View File
@@ -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()}`;
}
+34
View File
@@ -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();
}
}
+19
View File
@@ -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();
}
}
+8
View File
@@ -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));
}
+1 -1
View File
@@ -33,7 +33,7 @@ export function getEngineName(): string {
}
export const VERSION: WAHAEnvironment = {
version: '2024.11.4',
version: '2024.11.7',
engine: getEngineName(),
tier: getWAHAVersion(),
browser: getBrowserExecutablePath(),
+12 -2
View File
@@ -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: