1 parent
935a2b515f
commit
f8bcaebf4a
4 files changed
+66
-17
No files matched your search
@@ -22,7 +22,7 @@ export class GowsEventStreamObservable extends Observable<EnginePayload> {
|
||||
},
|
||||
) {
|
||||
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<EnginePayload> {
|
||||
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);
|
||||
};
|
||||
|
||||
|
||||
@@ -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<messages.MessageServiceClient> | null = null;
|
||||
let STREAM_CLIENTS: Pool<messages.EventStreamClient> | 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);
|
||||
};
|
||||
@@ -0,0 +1,41 @@
|
||||
import * as crypto from 'crypto';
|
||||
|
||||
export class Pool<Client> {
|
||||
private instances: Map<any, Client> = 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<Client> extends Pool<Client> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -765,7 +765,6 @@ export class WhatsappSessionGoWSCore extends WhatsappSession {
|
||||
this.status = WAHASessionStatus.STOPPED;
|
||||
this.events?.stop();
|
||||
this.stopEvents();
|
||||
this.client?.close();
|
||||
this.mediaManager.close();
|
||||
}
|
||||
|
||||
|
||||
Reference in new issue
Block a user