diff --git a/src/api/server.controller.ts b/src/api/server.controller.ts index c7d5c1b8..362b2d90 100644 --- a/src/api/server.controller.ts +++ b/src/api/server.controller.ts @@ -32,7 +32,7 @@ import * as lodash from 'lodash'; export class ServerController { private logger: Logger; - constructor() { + constructor(private config: WhatsappConfigService) { this.logger = new Logger('ServerController'); } @@ -77,6 +77,9 @@ export class ServerController { return { startTimestamp: startTimestamp, uptime: uptime, + worker: { + id: this.config.workerId, + }, }; } diff --git a/src/api/sessions.controller.ts b/src/api/sessions.controller.ts index 2fb0c3bb..40a979d4 100644 --- a/src/api/sessions.controller.ts +++ b/src/api/sessions.controller.ts @@ -91,6 +91,7 @@ class SessionsController { const start = request.start || false; await this.manager.upsert(name, config); if (start) { + await this.manager.assign(name); await this.manager.start(name); } }); @@ -133,6 +134,7 @@ class SessionsController { @UsePipes(new WAHAValidationPipe()) async delete(@Param('session') name: string): Promise { await this.withLock(name, async () => { + await this.manager.unassign(name); await this.manager.stop(name, true); await this.manager.logout(name); await this.manager.delete(name); @@ -153,6 +155,7 @@ class SessionsController { if (!exists) { throw new NotFoundException('Session not found'); } + await this.manager.assign(name); await this.manager.start(name); }); return await this.manager.getSessionInfo(name); @@ -167,6 +170,7 @@ class SessionsController { @UsePipes(new WAHAValidationPipe()) async stop(@Param('session') name: string): Promise { await this.withLock(name, async () => { + await this.manager.unassign(name); await this.manager.stop(name, false); }); return await this.manager.getSessionInfo(name); @@ -208,6 +212,7 @@ class SessionsController { if (!exists) { throw new NotFoundException('Session not found'); } + await this.manager.assign(name); await this.manager.stop(name, true); await this.manager.start(name); }); @@ -234,6 +239,7 @@ class SessionsController { return await this.withLock(name, async () => { const config = request.config; await this.manager.upsert(name, config); + await this.manager.assign(name); return await this.manager.start(name); }); } @@ -251,12 +257,14 @@ class SessionsController { if (request.logout) { // Old API did remove the session complete await this.withLock(name, async () => { + await this.manager.unassign(name); await this.manager.stop(name, true); await this.manager.logout(name); await this.manager.delete(name); }); } else { await this.withLock(name, async () => { + await this.manager.unassign(name); await this.manager.stop(name, false); }); } @@ -274,6 +282,7 @@ class SessionsController { ): Promise { const name = request.name; await this.withLock(name, async () => { + await this.manager.unassign(name); await this.manager.stop(name, true); await this.manager.logout(name); await this.manager.delete(name); diff --git a/src/config.service.ts b/src/config.service.ts index 28df30c0..6b028499 100644 --- a/src/config.service.ts +++ b/src/config.service.ts @@ -34,6 +34,18 @@ export class WhatsappConfigService { return baseUrl.replace(/\/$/, ''); } + get workerId(): string { + return this.configService.get('WAHA_WORKER_ID', ''); + } + + get shouldRestartWorkerSessions(): boolean { + const value = this.configService.get( + 'WAHA_WORKER_RESTART_SESSIONS', + 'true', + ); + return parseBool(value); + } + get mimetypes(): string[] { if (!this.shouldDownloadMedia) { return ['mimetype/ignore-all-media']; diff --git a/src/core/abc/manager.abc.ts b/src/core/abc/manager.abc.ts index ca4241cb..fd4431dc 100644 --- a/src/core/abc/manager.abc.ts +++ b/src/core/abc/manager.abc.ts @@ -4,6 +4,7 @@ import { } from '@nestjs/common'; import { WhatsappConfigService } from '@waha/config.service'; import { ISessionMeRepository } from '@waha/core/storage/ISessionMeRepository'; +import { ISessionWorkerRepository } from '@waha/core/storage/ISessionWorkerRepository'; import { WAHAWebhook } from '@waha/structures/webhooks.dto'; import { waitUntil } from '@waha/utils/promiseTimeout'; import { VERSION } from '@waha/version'; @@ -34,6 +35,7 @@ export abstract class SessionManager implements BeforeApplicationShutdown { public sessionAuthRepository: ISessionAuthRepository; public sessionConfigRepository: ISessionConfigRepository; protected sessionMeRepository: ISessionMeRepository; + protected sessionWorkerRepository: ISessionWorkerRepository; public events: EventEmitter; private lock: any; @@ -52,16 +54,11 @@ export abstract class SessionManager implements BeforeApplicationShutdown { const startSessions = this.config.startSessions; startSessions.forEach((sessionName) => { this.withLock(sessionName, async () => { - this.log.info( - { session: sessionName }, - `Restarting PREDEFINED session...`, - ); + const log = this.log.logger.child({ session: sessionName }); + log.info(`Restarting PREDEFINED session...`); this.start(sessionName).catch((error) => { - this.log.error( - { session: sessionName }, - `Failed to start predefined session: ${error}`, - ); - this.log.error(error.stack); + log.error(`Failed to start PREDEFINED session: ${error}`); + log.error(error.stack); }); }); }); @@ -103,6 +100,18 @@ export abstract class SessionManager implements BeforeApplicationShutdown { abstract getSessions(all: boolean): Promise; + get workerId() { + return this.config.workerId; + } + + async assign(name: string) { + await this.sessionWorkerRepository?.assign(name, this.workerId); + } + + async unassign(name: string) { + await this.sessionWorkerRepository?.unassign(name, this.workerId); + } + protected handleSessionEvent(event: WAHAEvents, session: WhatsappSession) { return (payload: any) => { const me = session.getSessionMeInfo(); diff --git a/src/core/storage/ISessionWorkerRepository.ts b/src/core/storage/ISessionWorkerRepository.ts new file mode 100644 index 00000000..c5d83b80 --- /dev/null +++ b/src/core/storage/ISessionWorkerRepository.ts @@ -0,0 +1,18 @@ +export class SessionWorkerInfo { + id: string; // aka session name here + worker: string; +} + +export abstract class ISessionWorkerRepository { + abstract assign(session: string, worker: string): Promise; + + abstract unassign(session: string, worker: string): Promise; + + abstract remove(session: string): Promise; + + abstract getAll(): Promise; + + abstract getSessionsByWorker(worker: string): Promise; + + abstract init(): Promise; +} diff --git a/src/core/storage/LocalSessionWorkerRepository.ts b/src/core/storage/LocalSessionWorkerRepository.ts new file mode 100644 index 00000000..d965ff02 --- /dev/null +++ b/src/core/storage/LocalSessionWorkerRepository.ts @@ -0,0 +1,74 @@ +import { Sqlite3SchemaValidation } from '@waha/core/engines/noweb/store/sqlite3/Sqlite3SchemaValidation'; +import { + ISessionWorkerRepository, + SessionWorkerInfo, +} from '@waha/core/storage/ISessionWorkerRepository'; +import { LocalStore } from '@waha/core/storage/LocalStore'; +import { Field, Index, Schema } from '@waha/core/storage/sqlite3/Schema'; +import { Sqlite3KVRepository } from '@waha/core/storage/sqlite3/Sqlite3KVRepository'; + +// eslint-disable-next-line @typescript-eslint/no-var-requires +const Database = require('better-sqlite3'); + +const SCHEMA = new Schema( + 'session_worker', + [ + new Field('id', 'TEXT'), + new Field('worker', 'TEXT'), + new Field('data', 'TEXT'), + ], + [ + new Index('session_worker_id_idx', ['id']), + new Index('session_worker_worker_idx', ['worker']), + ], +); + +export class LocalSessionWorkerRepository + extends Sqlite3KVRepository + implements ISessionWorkerRepository +{ + constructor(store: LocalStore) { + const db = store.getWAHADatabase(); + super(db, SCHEMA); + } + + assign(session: string, worker: string): Promise { + return this.upsertOne({ id: session, worker: worker }); + } + + unassign(session: string, worker: string): Promise { + return this.deleteBy({ id: session, worker: worker }); + } + + remove(session: string) { + return this.deleteById(session); + } + + async getSessionsByWorker(worker: string): Promise { + const data = await this.getAllBy({ worker: worker }); + return data.map((d) => d.id); + } + + async init(): Promise { + this.migrations(); + this.validateSchema(); + } + + private migrations() { + this.db.exec( + 'CREATE TABLE IF NOT EXISTS session_worker (id TEXT, worker TEXT, data TEXT)', + ); + // Session can have only one record + this.db.exec( + 'CREATE UNIQUE INDEX IF NOT EXISTS session_worker_id_idx ON session_worker (id)', + ); + // Worker can have multiple records + this.db.exec( + 'CREATE INDEX IF NOT EXISTS session_worker_worker_idx ON session_worker (worker)', + ); + } + + private validateSchema() { + new Sqlite3SchemaValidation(SCHEMA, this.db).validate(); + } +} diff --git a/src/structures/server.dto.ts b/src/structures/server.dto.ts index 46daade2..a3265583 100644 --- a/src/structures/server.dto.ts +++ b/src/structures/server.dto.ts @@ -36,6 +36,14 @@ export class StopResponse { stopping: boolean = true; } +export class WorkerInfo { + @ApiProperty({ + example: 'waha', + description: 'The worker ID.', + }) + id: string; +} + export class ServerStatusResponse { @ApiProperty({ example: 1723788847247, @@ -48,4 +56,6 @@ export class ServerStatusResponse { description: 'The uptime of the server in milliseconds.', }) uptime: number; + + worker: WorkerInfo; } diff --git a/src/structures/sessions.dto.ts b/src/structures/sessions.dto.ts index b8278319..1aeee6c4 100644 --- a/src/structures/sessions.dto.ts +++ b/src/structures/sessions.dto.ts @@ -159,6 +159,7 @@ export class MeInfo { export class SessionInfo extends SessionDTO { me?: MeInfo; + assignedWorker?: string; } export class SessionDetailedInfo extends SessionInfo {