import { Cache } from '../cache/Cache' import { Transport, TransportType } from '../transports/Transport' import { CallRequestDTO } from '../CallRequestDTO' import { CallResponseDTO } from '../CallResponseDTO' import { CallRequestError } from '../errors/CallRequestError' import { InterceptorError } from '../errors/InterceptorError' import { CallResponseInterceptor } from '../CallResponseInterceptor' import { CallRequestInterceptor } from '../CallRequestInterceptor' import { makeCallRequestKey } from '../utils/cache-utils' import { RPCClientIdentity } from '../RPCClientIdentity' import { cloneDeep } from 'lodash-es' import { GetDebugLogger } from '../utils/debug-logger' import { UnsubscribeCallback } from '../types' import { CallOptions, RPCClientOptions } from '../RPCClient' import { CallRequestTransportError, CallRequestTransportTimeoutError, } from '../transports/errors' const debug = GetDebugLogger('rpc:ClientManager') interface InFlightEntry { promise: Promise createdAt: number } const MAX_INFLIGHT_CALLS = 100 const STALE_ENTRY_MULTIPLIER = 2 export interface ClientManagerOptions { cache?: Cache deadlineMs: number onTransportRequestError?: ( error: Error | CallRequestTransportError | CallRequestTransportTimeoutError, transportType: TransportType, request: CallRequestDTO, ) => void requestInterceptor?: CallRequestInterceptor responseInterceptor?: CallResponseInterceptor transports: Transport[] transportOptions?: RPCClientOptions['transportOptions'] } export interface ClientManager { addResponseInterceptor( interceptor: CallResponseInterceptor, ): UnsubscribeCallback addRequestInterceptor( interceptor: CallRequestInterceptor, ): UnsubscribeCallback clearCache(request?: CallRequestDTO): void getCachedResponse(request: CallRequestDTO): CallResponseDTO | undefined getInFlightCallCount(): number manageClientRequest( request: CallRequestDTO, opts?: CallOptions, ): Promise setIdentity(identity?: RPCClientIdentity): void } export class ClientManager implements ClientManager { private cache?: Cache private cleanupIntervalId?: ReturnType private inflightCallsByKey: Map = new Map() private interceptors: { request: CallRequestInterceptor[] response: CallResponseInterceptor[] } = { request: [], response: [], } private transports: Transport[] constructor(readonly options: ClientManagerOptions) { this.cache = options.cache if (options.requestInterceptor) { this.interceptors.request.push(options.requestInterceptor) } if (options.responseInterceptor) { this.interceptors.response.push(options.responseInterceptor) } this.transports = options.transports // Start periodic cleanup of stale in-flight entries this.cleanupIntervalId = setInterval( () => this.cleanupStaleEntries(), this.options.deadlineMs, ) } private cleanupStaleEntries = () => { const now = Date.now() const staleThreshold = this.options.deadlineMs * STALE_ENTRY_MULTIPLIER for (const [key, entry] of this.inflightCallsByKey) { if (now - entry.createdAt > staleThreshold) { debug('cleanupStaleEntries removing stale entry:', key) this.inflightCallsByKey.delete(key) } } } public addResponseInterceptor = (interceptor: CallResponseInterceptor) => { if (typeof interceptor !== 'function') throw new Error('cannot add interceptor that is not a function') this.interceptors.response.push(interceptor) return () => { this.interceptors.response = this.interceptors.response.filter( (cb) => cb !== interceptor, ) } } public addRequestInterceptor = (interceptor: CallRequestInterceptor) => { if (typeof interceptor !== 'function') throw new Error('cannot add interceptor that is not a function') this.interceptors.request.push(interceptor) return () => { this.interceptors.request = this.interceptors.request.filter( (cb) => cb !== interceptor, ) } } public clearCache = (request?: CallRequestDTO) => { debug('clearCache') if (!this.cache) return this.cache.clearCache(request) // Also clear matching in-flight entries to prevent stale data from // being returned after cache is cleared if (request) { const key = makeCallRequestKey(request) this.inflightCallsByKey.delete(key) } else { this.inflightCallsByKey.clear() } } public getCachedResponse = (request: CallRequestDTO) => { debug('getCachedResponse', request) return this?.cache?.getCachedResponse(request) } public getInFlightCallCount = () => { debug('getInFlightCallCount') return this.inflightCallsByKey.size } public manageClientRequest = async ( originalRequest: CallRequestDTO, opts?: CallOptions, ) => { debug('manageClientRequest', originalRequest) const request = await this.interceptRequestMutator(originalRequest) const key = makeCallRequestKey(request) const existingEntry = this.inflightCallsByKey.get(key) if (existingEntry) { debug('manageClientRequest using an existing in-flight call', key) return existingEntry.promise.then((res: CallResponseDTO) => { const response = cloneDeep(res) if (originalRequest.correlationId) { response.correlationId = originalRequest.correlationId } return response }) } // Check max size before adding new entry if (this.inflightCallsByKey.size >= MAX_INFLIGHT_CALLS) { console.warn( `RPCClient - MAX_INFLIGHT_CALLS` ) this.cleanupStaleEntries() // If still at max after cleanup, remove oldest entry if (this.inflightCallsByKey.size >= MAX_INFLIGHT_CALLS) { const oldestKey = this.inflightCallsByKey.keys().next().value if (oldestKey) { debug('manageClientRequest removing oldest entry due to max size:', oldestKey) this.inflightCallsByKey.delete(oldestKey) } } } // Create placeholder entry FIRST (synchronously) to prevent race condition // where concurrent calls could slip through between checking and setting const entry: InFlightEntry = { promise: null as any, // Will be replaced immediately below createdAt: Date.now(), } this.inflightCallsByKey.set(key, entry) const requestPromiseWrapper = async () => { const preferredTransport = this.options.transportOptions?.preferredTransport const useTransport = opts?.transport ?? preferredTransport const transportResponse = await this.sendRequestWithTransport(request, { ...opts, transport: useTransport, }) // Cache the RAW transport response BEFORE running interceptors // This prevents user-specific interceptor mutations from being cached if (this.cache && transportResponse) { this.cache.setCachedResponse(request, transportResponse) } const response = await this.interceptResponseMutator( transportResponse, request, ) return response } const requestPromise = requestPromiseWrapper().finally(() => { this.inflightCallsByKey.delete(key) }) // Update the placeholder with the actual promise entry.promise = requestPromise return await requestPromise } public setIdentity = (identity?: RPCClientIdentity) => { debug('setIdentity', identity) this.transports.forEach((transport) => transport.setIdentity(identity)) } private interceptRequestMutator = async ( originalRequest: CallRequestDTO, ): Promise => { let request = originalRequest if (this.interceptors.request.length) { debug('interceptRequestMutator(), original request:', originalRequest) let interceptorIndex = 0 for (const interceptor of this.interceptors.request) { try { request = await interceptor(request) } catch (e: any) { debug( 'caught request interceptor error, interceptor index:', interceptorIndex, 'request:', request, 'error:', e, ) throw new InterceptorError({ type: 'request', originalError: e instanceof Error ? e : new Error(String(e)), request: originalRequest, interceptorIndex, }) } interceptorIndex++ } } return request } private interceptResponseMutator = async ( originalResponse: CallResponseDTO, request: CallRequestDTO, ): Promise => { let response = originalResponse if (this.interceptors.response.length) { debug('interceptResponseMutator', request, originalResponse) let interceptorIndex = 0 for (const interceptor of this.interceptors.response) { try { response = await interceptor(response, request) } catch (e: any) { debug( 'caught response interceptor error, interceptor index:', interceptorIndex, 'request:', request, 'original response:', originalResponse, 'mutated response:', response, 'error:', e, ) throw new InterceptorError({ type: 'response', originalError: e instanceof Error ? e : new Error(String(e)), request, response: originalResponse, interceptorIndex, }) } interceptorIndex++ } } return response } private sendRequestWithTransport = async ( request: CallRequestDTO, opts?: CallOptions, ) => { debug('sendRequestWithTransport', request) let response let transportsIndex = 0 let transports = this.transports if (opts?.transport) { const transportToUse = this.transports.find( (transport) => transport.type === opts.transport, ) if (!transportToUse) { throw new CallRequestError( `Specified Transport "${opts.transport}" not available`, ) } transports = [transportToUse] } while (!response && transports[transportsIndex]) { const transport = transports[transportsIndex] if (!transport.isConnected()) { transportsIndex++ debug( `sendRequestWithTransport transport ${transport.name} not connected, trying next transport "${transports[transportsIndex]?.name ?? 'no next transport'}"`, ) continue } debug(`sendRequestWithTransport trying ${transport.name}`, request) try { let timeout: ReturnType | undefined let isSettled = false const transportPromise = transport.sendRequest(request, opts || {}) const deadlinePromise = new Promise((_, reject) => { timeout = setTimeout(() => { if (!isSettled) { reject( new CallRequestTransportTimeoutError( `Procedure ${request.procedure}`, ), ) } }, opts?.timeout || this.options.deadlineMs) }) // Handle the race with proper cleanup const raceResult = await Promise.race([ deadlinePromise, transportPromise.then((result) => { isSettled = true return result }), ]).finally(() => { if (timeout) { clearTimeout(timeout) timeout = undefined } }) // Handle the losing promise to prevent unhandled rejections // The transport promise will complete eventually, we just ignore its result transportPromise.catch((err) => { if (isSettled) { debug( 'sendRequestWithTransport: transport promise rejected after timeout, ignoring:', err, ) } }) response = raceResult } catch (e: any) { this.options.onTransportRequestError?.( e, transports[transportsIndex].type, request, ) if ( e instanceof CallRequestTransportTimeoutError || e instanceof CallRequestTransportError ) { console.warn( `RPCClient ClientManager#sendRequestWithTransport() sending request with ${ transports[transportsIndex].name } failed ${ transports[transportsIndex + 1] ? `trying next transport "${ transports[transportsIndex + 1].name }"` : 'no more transports to try' }`, ) transportsIndex++ } else { throw e } } } if (!response && !this.transports[transportsIndex]) { console.error( `RPCClient ClientManager#sendRequestWithTransport() ${request.scope}::${request.procedure}::${request.version} did not get a response from any transports`, ) throw new CallRequestError( `Procedure ${request.procedure} did not get a response from any transports`, ) } else if (!response) { throw new CallRequestError( `Procedure ${request.procedure} did not get a response from the remote`, ) } return response } public destroy = () => { debug('destroy') if (this.cleanupIntervalId) { clearInterval(this.cleanupIntervalId) this.cleanupIntervalId = undefined } this.inflightCallsByKey.clear() this.interceptors.request = [] this.interceptors.response = [] } }