import * as grpc from 'grpc'; import { Interceptor } from "./interceptors/chain"; import { GrpcMonitor } from "./monitor"; import { Sentinel } from "@nimbus/operational"; export interface PoolOptions { address: string; credentials: grpc.ChannelCredentials; size?: number; options?: { [x: string]: any; interceptors?: Interceptor[]; }, } export enum PoolStatus { AllDisconnected, AllConnected, SomeConnected, } export type ClientConstructor = new (address: string, credentials: grpc.ChannelCredentials, options?: object) => Client export class Pool { clients: Client[] = []; provider: IterableIterator; status: PoolStatus = PoolStatus.AllDisconnected; _monitor?: GrpcMonitor; constructor(clientProto: ClientConstructor, options: PoolOptions) { const size = options.size || 1; for (let i = 1; i <= size; i++) { this.clients.push(this.make(i, clientProto, options)); } this.provider = this.startProvider(); this.watchConnections(); } private watchConnections(): void { const result: boolean[] = []; this.clients.forEach((client, index) => { const channel = grpc.getClientChannel(client); setInterval(() => { const connectivityState = channel.getConnectivityState(true); result[index] = connectivityState === grpc.connectivityState.READY; }, 2000); }); setInterval( () => { if (result.every(s => s === true)) { this.status = PoolStatus.AllConnected; } else if (result.every(s => s === false)) { this.status = PoolStatus.AllDisconnected; } else { this.status = PoolStatus.SomeConnected; } }, 2000); } protected make(index: number, clientProto: ClientConstructor, options: PoolOptions): Client { const grpcOptions = { ...(options.options ? options.options : {}), 'connection.index': index, }; return new clientProto(options.address, options.credentials, grpcOptions); } private *startProvider(): IterableIterator { let index = 0; while (true) { yield this.clients[index]; index = (index + 1) % this.clients.length; } } get(): Client { return this.provider.next().value; } monitor(sentinel: Sentinel) { this._monitor = new GrpcMonitor(this); this._monitor.reportHealth(sentinel); return this; } }