[core] websockets
This commit is contained in:
1 parent
60e9e53af4
commit
e15ca41d6e
12 files changed
+1170
-1089
No files matched your search
@@ -30,9 +30,11 @@
|
||||
"@nestjs/core": "^9.0.9",
|
||||
"@nestjs/passport": "^9.0.0",
|
||||
"@nestjs/platform-express": "^9.0.9",
|
||||
"@nestjs/platform-ws": "^9.0.9",
|
||||
"@nestjs/serve-static": "^2.1.3",
|
||||
"@nestjs/swagger": "^7.1.11",
|
||||
"@nestjs/terminus": "^10.2.3",
|
||||
"@nestjs/websockets": "^9.0.9",
|
||||
"@types/better-sqlite3": "^7.6.10",
|
||||
"@types/lodash": "^4.14.194",
|
||||
"@types/ws": "^8.5.4",
|
||||
|
||||
@@ -1,6 +1,12 @@
|
||||
import { OnApplicationShutdown } from '@nestjs/common';
|
||||
import {
|
||||
BeforeApplicationShutdown,
|
||||
OnApplicationShutdown,
|
||||
} from '@nestjs/common';
|
||||
import { WAHAWebhook } from '@waha/structures/webhooks.dto';
|
||||
import { VERSION } from '@waha/version';
|
||||
import { EventEmitter } from 'events';
|
||||
|
||||
import { WAHAEngine } from '../../structures/enums.dto';
|
||||
import { WAHAEngine, WAHAEvents } from '../../structures/enums.dto';
|
||||
import {
|
||||
SessionDTO,
|
||||
SessionInfo,
|
||||
@@ -13,10 +19,11 @@ import { ISessionConfigRepository } from '../storage/ISessionConfigRepository';
|
||||
import { WhatsappSession } from './session.abc';
|
||||
import { WebhookConductor } from './webhooks.abc';
|
||||
|
||||
export abstract class SessionManager implements OnApplicationShutdown {
|
||||
export abstract class SessionManager implements BeforeApplicationShutdown {
|
||||
public store: any;
|
||||
public sessionAuthRepository: ISessionAuthRepository;
|
||||
public sessionConfigRepository: ISessionConfigRepository;
|
||||
public events: EventEmitter;
|
||||
|
||||
protected abstract getEngine(engine: WAHAEngine): typeof WhatsappSession;
|
||||
|
||||
@@ -24,7 +31,7 @@ export abstract class SessionManager implements OnApplicationShutdown {
|
||||
|
||||
protected abstract get WebhookConductorClass(): typeof WebhookConductor;
|
||||
|
||||
abstract onApplicationShutdown(signal?: string);
|
||||
abstract beforeApplicationShutdown(signal?: string);
|
||||
|
||||
//
|
||||
// API Methods
|
||||
@@ -40,4 +47,19 @@ export abstract class SessionManager implements OnApplicationShutdown {
|
||||
abstract getSessionInfo(name: string): Promise<SessionInfo | null>;
|
||||
|
||||
abstract getSessions(all: boolean): Promise<SessionInfo[]>;
|
||||
|
||||
handleSessionEvent(event: WAHAEvents, session: WhatsappSession) {
|
||||
return async (payload: any) => {
|
||||
const me = await session.getSessionMeInfo().catch((err) => null);
|
||||
const data: WAHAWebhook = {
|
||||
event: event,
|
||||
session: session.name,
|
||||
me: me,
|
||||
payload: payload,
|
||||
engine: session.engine,
|
||||
environment: VERSION,
|
||||
};
|
||||
this.events.emit(event, data);
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,132 @@
|
||||
import { BeforeApplicationShutdown, ConsoleLogger } from '@nestjs/common';
|
||||
import { sleep } from '@nestjs/terminus/dist/utils';
|
||||
import {
|
||||
OnGatewayConnection,
|
||||
OnGatewayDisconnect,
|
||||
OnGatewayInit,
|
||||
WebSocketGateway,
|
||||
WebSocketServer,
|
||||
} from '@nestjs/websockets';
|
||||
import { SessionManager } from '@waha/core/abc/manager.abc';
|
||||
import { getLogLevels } from '@waha/helpers';
|
||||
import { WAHAEvents } from '@waha/structures/enums.dto';
|
||||
import { HeartbeatJob } from '@waha/utils/HeartbeatJob';
|
||||
import { WebSocket } from '@waha/utils/ws';
|
||||
import { IncomingMessage } from 'http';
|
||||
import * as lodash from 'lodash';
|
||||
import * as url from 'url';
|
||||
import { Server } from 'ws';
|
||||
|
||||
@WebSocketGateway({
|
||||
path: '/ws',
|
||||
cors: true,
|
||||
})
|
||||
export class WebsocketGatewayCore
|
||||
implements
|
||||
OnGatewayInit,
|
||||
OnGatewayConnection,
|
||||
OnGatewayDisconnect,
|
||||
BeforeApplicationShutdown
|
||||
{
|
||||
HEARTBEAT_INTERVAL = 10_000;
|
||||
|
||||
@WebSocketServer()
|
||||
server: Server;
|
||||
|
||||
private listeners: Map<WebSocket, { session: string; events: string[] }> =
|
||||
new Map();
|
||||
|
||||
private log: ConsoleLogger;
|
||||
private heartbeat: HeartbeatJob;
|
||||
|
||||
constructor(private manager: SessionManager) {
|
||||
const levels = getLogLevels(false);
|
||||
this.log = new ConsoleLogger('WebsocketGateway', { logLevels: levels });
|
||||
this.heartbeat = new HeartbeatJob(this.log, this.HEARTBEAT_INTERVAL);
|
||||
}
|
||||
|
||||
handleConnection(socket: WebSocket, request: IncomingMessage, ...args): any {
|
||||
const random = Math.random().toString(36).substring(7);
|
||||
const id = `wsc_${random}`;
|
||||
socket.id = id;
|
||||
this.log.debug(`New client connected: ${request.url}`);
|
||||
const params = this.getParams(id, request, socket);
|
||||
if (!params) {
|
||||
return;
|
||||
}
|
||||
const { session, events } = params;
|
||||
this.log.debug(
|
||||
`Client connected to session: '${session}', events: ${events}, ${id}`,
|
||||
);
|
||||
this.listeners.set(socket, { session, events });
|
||||
}
|
||||
|
||||
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.log.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.log.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;
|
||||
}
|
||||
return { session, events };
|
||||
}
|
||||
|
||||
handleDisconnect(socket: WebSocket): any {
|
||||
this.log.debug(`Client disconnected - ${socket.id}`);
|
||||
this.listeners.delete(socket);
|
||||
}
|
||||
|
||||
async beforeApplicationShutdown(signal?: string) {
|
||||
this.log.log('Shutting down websocket server');
|
||||
// Allow pending messages to be sent, it can be even 1ms, just to release the event loop
|
||||
await sleep(100);
|
||||
this.server.clients.forEach((options, client) => {
|
||||
client.close(1001, 'Server is shutting down');
|
||||
});
|
||||
// Do not turn off heartbeat service here,
|
||||
// it's responsible for terminating the connection that is not alive
|
||||
this.log.log('Websocket server is down');
|
||||
}
|
||||
|
||||
afterInit(server: Server) {
|
||||
this.log.debug('Websocket server initialized');
|
||||
this.manager.events.on(
|
||||
WAHAEvents.SESSION_STATUS,
|
||||
this.sendToAll.bind(this),
|
||||
);
|
||||
this.log.debug('Subscribed to manager events');
|
||||
|
||||
this.log.debug('Starting heartbeat service...');
|
||||
this.heartbeat.start(server);
|
||||
this.log.debug('Heartbeat service started');
|
||||
}
|
||||
|
||||
sendToAll(data: any) {
|
||||
this.listeners.forEach((options, client) => {
|
||||
if (options.events.length === 0) {
|
||||
return;
|
||||
}
|
||||
this.log.debug('Sending data to client', data);
|
||||
client.send(JSON.stringify(data));
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,7 @@ import { PassportModule } from '@nestjs/passport';
|
||||
import { ServeStaticModule } from '@nestjs/serve-static';
|
||||
import { TerminusModule } from '@nestjs/terminus';
|
||||
import { BufferJsonReplacerInterceptor } from '@waha/api/BufferJsonReplacerInterceptor';
|
||||
import { WebsocketGatewayCore } from '@waha/core/api/websocket.gateway.core';
|
||||
import { join } from 'path';
|
||||
|
||||
import { AuthController } from '../api/auth.controller';
|
||||
@@ -93,6 +94,7 @@ const PROVIDERS = [
|
||||
WhatsappConfigService,
|
||||
EngineConfigService,
|
||||
ConsoleLogger,
|
||||
WebsocketGatewayCore,
|
||||
];
|
||||
|
||||
@Module({
|
||||
|
||||
@@ -350,8 +350,9 @@ export class WhatsappSessionNoWebCore extends WhatsappSession {
|
||||
}
|
||||
|
||||
async stop() {
|
||||
await this.end();
|
||||
this.status = WAHASessionStatus.STOPPED;
|
||||
this.events.removeAllListeners();
|
||||
await this.end();
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
@@ -90,8 +90,9 @@ export class WhatsappSessionVenomCore extends WhatsappSession {
|
||||
}
|
||||
|
||||
async stop() {
|
||||
await this.whatsapp.close();
|
||||
this.status = WAHASessionStatus.STOPPED;
|
||||
this.events.removeAllListeners();
|
||||
await this.whatsapp.close();
|
||||
}
|
||||
|
||||
subscribeEngineEvent(event: WAHAEvents | string, handler: (message) => void) {
|
||||
|
||||
@@ -204,6 +204,7 @@ export class WhatsappSessionWebJSCore extends WhatsappSession {
|
||||
async stop() {
|
||||
this.shouldRestart = false;
|
||||
this.status = WAHASessionStatus.STOPPED;
|
||||
this.events.removeAllListeners();
|
||||
await this.end();
|
||||
}
|
||||
|
||||
|
||||
@@ -5,10 +5,15 @@ import {
|
||||
NotFoundException,
|
||||
UnprocessableEntityException,
|
||||
} from '@nestjs/common';
|
||||
import { EventEmitter } from 'events';
|
||||
|
||||
import { WhatsappConfigService } from '../config.service';
|
||||
import { getLogLevels } from '../helpers';
|
||||
import { WAHAEngine, WAHASessionStatus } from '../structures/enums.dto';
|
||||
import {
|
||||
WAHAEngine,
|
||||
WAHAEvents,
|
||||
WAHASessionStatus,
|
||||
} from '../structures/enums.dto';
|
||||
import {
|
||||
ProxyConfig,
|
||||
SessionDTO,
|
||||
@@ -59,7 +64,7 @@ export class SessionManagerCore extends SessionManager {
|
||||
private engineConfigService: EngineConfigService,
|
||||
) {
|
||||
super();
|
||||
|
||||
this.events = new EventEmitter();
|
||||
this.log.setContext('SessionManager');
|
||||
this.session = undefined;
|
||||
const engineName = this.engineConfigService.getDefaultEngineName();
|
||||
@@ -95,7 +100,7 @@ export class SessionManagerCore extends SessionManager {
|
||||
}
|
||||
}
|
||||
|
||||
async onApplicationShutdown(signal?: string) {
|
||||
async beforeApplicationShutdown(signal?: string) {
|
||||
if (!this.session) {
|
||||
return;
|
||||
}
|
||||
@@ -142,6 +147,12 @@ export class SessionManagerCore extends SessionManager {
|
||||
const webhooks = this.getWebhooks(request);
|
||||
webhook.configure(session, webhooks);
|
||||
|
||||
// configure events
|
||||
session.events.on(
|
||||
WAHAEvents.SESSION_STATUS,
|
||||
this.handleSessionEvent(WAHAEvents.SESSION_STATUS, session),
|
||||
);
|
||||
|
||||
// start session
|
||||
await session.start();
|
||||
return {
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import { NestFactory } from '@nestjs/core';
|
||||
import { WsAdapter } from '@nestjs/platform-ws';
|
||||
import { json, urlencoded } from 'express';
|
||||
|
||||
import { AllExceptionsFilter } from './api/exception.filter';
|
||||
@@ -56,6 +57,7 @@ async function bootstrap() {
|
||||
// Allow to send big body - for images and attachments
|
||||
app.use(json({ limit: '50mb' }));
|
||||
app.use(urlencoded({ limit: '50mb', extended: false }));
|
||||
app.useWebSocketAdapter(new WsAdapter(app));
|
||||
|
||||
// Configure swagger
|
||||
const swaggerConfigurator = new SwaggerModule(app);
|
||||
|
||||
@@ -0,0 +1,45 @@
|
||||
import { ConsoleLogger } from '@nestjs/common';
|
||||
import { WebSocket } from '@waha/utils/ws';
|
||||
import { WebSocketServer } from 'ws';
|
||||
|
||||
export class HeartbeatJob {
|
||||
private interval: ReturnType<typeof setInterval>;
|
||||
|
||||
constructor(
|
||||
private log: ConsoleLogger,
|
||||
private intervalTime: number = 10_000,
|
||||
) {}
|
||||
|
||||
start(server: WebSocketServer) {
|
||||
server.on('connection', (ws: WebSocket) => {
|
||||
ws.isAlive = true;
|
||||
ws.on('pong', this.onPong(ws));
|
||||
});
|
||||
|
||||
this.interval = setInterval(() => {
|
||||
server.clients.forEach((client: WebSocket) => {
|
||||
if (client.isAlive === false) {
|
||||
this.log.debug(
|
||||
`Terminating client connection due to heartbeat timeout, ${client.id}`,
|
||||
);
|
||||
return client.terminate();
|
||||
}
|
||||
|
||||
client.isAlive = false;
|
||||
this.log.debug(`Sending heartbeat (ping) to ${client.id}`);
|
||||
client.ping();
|
||||
});
|
||||
}, this.intervalTime);
|
||||
}
|
||||
|
||||
stop() {
|
||||
clearInterval(this.interval);
|
||||
}
|
||||
|
||||
private onPong(ws: WebSocket) {
|
||||
return (event: any) => {
|
||||
ws.isAlive = true;
|
||||
this.log.debug(`Heartbeat (pong) received from ${ws.id}`);
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
import { WebSocket as WSWebSocket } from 'ws';
|
||||
|
||||
export interface WebSocket extends WSWebSocket {
|
||||
id?: string;
|
||||
isAlive?: boolean;
|
||||
}
|
||||
Reference in new issue
Block a user