From f8bcaebf4a245f9349a38f2c6ffc9a68ebbe6eb0 Mon Sep 17 00:00:00 2001 From: devlikepro Date: Fri, 12 Dec 2025 12:15:31 +0700 Subject: [PATCH] [core] GOWS - make grpc pool for clients fix #1715 --- .../engines/gows/GowsEventStreamObservable.ts | 9 +--- src/core/engines/gows/clients.ts | 32 +++++++++++---- src/core/engines/gows/pools.ts | 41 +++++++++++++++++++ src/core/engines/gows/session.gows.core.ts | 1 - 4 files changed, 66 insertions(+), 17 deletions(-) create mode 100644 src/core/engines/gows/pools.ts diff --git a/src/core/engines/gows/GowsEventStreamObservable.ts b/src/core/engines/gows/GowsEventStreamObservable.ts index e6ea5753..a8c1d0de 100644 --- a/src/core/engines/gows/GowsEventStreamObservable.ts +++ b/src/core/engines/gows/GowsEventStreamObservable.ts @@ -22,7 +22,7 @@ export class GowsEventStreamObservable extends Observable { }, ) { super((subscriber) => { - logger.debug('Creating grpc client and stream...'); + logger.debug('Creating grpc stream...'); logger.setBindings({ id: rand() }); const { client, stream } = factory(); this._client = client; @@ -41,13 +41,6 @@ export class GowsEventStreamObservable extends Observable { logger.warn({ err }, 'Failed to cancel gRPC stream'); } - logger.debug({ reason }, 'Closing gRPC client...'); - try { - client.close(); - } catch (err) { - logger.warn({ err }, 'Failed to close gRPC client'); - } - await sleep(this.CLIENT_CLOSE_TIMEOUT); }; diff --git a/src/core/engines/gows/clients.ts b/src/core/engines/gows/clients.ts index 54a8cafa..fa040c16 100644 --- a/src/core/engines/gows/clients.ts +++ b/src/core/engines/gows/clients.ts @@ -1,15 +1,25 @@ import * as grpc from '@grpc/grpc-js'; import { messages } from '@waha/core/engines/gows/grpc/gows'; +import { Pool, SizedPool } from '@waha/core/engines/gows/pools'; + +let CLIENTS: Pool | null = null; +let STREAM_CLIENTS: Pool | null = null; export const GetMessageServiceClient = ( session: string, address: string, credentials: grpc.ChannelCredentials, ): messages.MessageServiceClient => { - return new messages.MessageServiceClient(address, credentials, { - 'grpc.max_send_message_length': 128 * 1024 * 1024, - 'grpc.max_receive_message_length': 128 * 1024 * 1024, - }); + if (!CLIENTS) { + const factory = () => { + return new messages.MessageServiceClient(address, credentials, { + 'grpc.max_send_message_length': 128 * 1024 * 1024, + 'grpc.max_receive_message_length': 128 * 1024 * 1024, + }); + }; + CLIENTS = new SizedPool(16, factory); + } + return CLIENTS.get(session); }; export const GetEventStreamClient = ( @@ -17,8 +27,14 @@ export const GetEventStreamClient = ( address: string, credentials: grpc.ChannelCredentials, ): messages.EventStreamClient => { - return new messages.EventStreamClient(address, credentials, { - 'grpc.max_send_message_length': 128 * 1024 * 1024, - 'grpc.max_receive_message_length': 128 * 1024 * 1024, - }); + if (!STREAM_CLIENTS) { + const factory = () => { + return new messages.EventStreamClient(address, credentials, { + 'grpc.max_send_message_length': 128 * 1024 * 1024, + 'grpc.max_receive_message_length': 128 * 1024 * 1024, + }); + }; + STREAM_CLIENTS = new SizedPool(16, factory); + } + return STREAM_CLIENTS.get(session); }; diff --git a/src/core/engines/gows/pools.ts b/src/core/engines/gows/pools.ts new file mode 100644 index 00000000..5219cf6f --- /dev/null +++ b/src/core/engines/gows/pools.ts @@ -0,0 +1,41 @@ +import * as crypto from 'crypto'; + +export class Pool { + private instances: Map = new Map(); + + constructor(private factory: () => Client) {} + + protected key(name: string): any { + return name; + } + + get(name: string): Client { + const key = this.key(name); + if (this.instances.has(key)) { + return this.instances.get(key); + } + const client = this.factory(); + this.instances.set(key, client); + return client; + } +} + +export class SizedPool extends Pool { + constructor( + protected size: number, + factory: () => Client, + ) { + if (!Number.isInteger(size) || size <= 0) { + throw new Error('size must be a positive integer'); + } + + super(factory); + } + + protected key(name: string): any { + const hash = crypto.createHash('sha256').update(name).digest(); + const num = hash.readUInt32BE(0); + const bucket = num % this.size; + return Number(bucket); + } +} diff --git a/src/core/engines/gows/session.gows.core.ts b/src/core/engines/gows/session.gows.core.ts index e026b750..4b493079 100644 --- a/src/core/engines/gows/session.gows.core.ts +++ b/src/core/engines/gows/session.gows.core.ts @@ -765,7 +765,6 @@ export class WhatsappSessionGoWSCore extends WhatsappSession { this.status = WAHASessionStatus.STOPPED; this.events?.stop(); this.stopEvents(); - this.client?.close(); this.mediaManager.close(); }