import * as grpc from "grpc"; export function retryInterceptorFactory( maxRetries = 1, codeMap: grpc.status[] = [grpc.status.UNAVAILABLE, grpc.status.UNKNOWN], ) { return (options: grpc.CallOptions, mainChainNext: Function): grpc.InterceptingCall => { let originMetadata: grpc.Metadata; let originMessage: any; let lastReceivedMessage: any; let receivedMessageCallback: Function; let lastReceivedStatus: grpc.StatusObject; return new grpc.InterceptingCall(mainChainNext(options), { sendMessage(message, next) { originMessage = message; next(message); }, start(metadata, listener, next) { originMetadata = metadata; const newListener: grpc.Listener = { onReceiveMessage(message, onReceiveNext) { lastReceivedMessage = message; receivedMessageCallback = onReceiveNext; }, onReceiveStatus(status, statusChainNext) { let currentRetries = 0; function retry(): void { if (currentRetries >= maxRetries) { receivedMessageCallback(lastReceivedMessage); statusChainNext(lastReceivedStatus); return; // no more retries allowed } currentRetries++; // trigger the interceptor chain manually for retry const retryCall = mainChainNext(options); retryCall.start(originMetadata, { onReceiveMessage(message: any) { lastReceivedMessage = message; }, onReceiveStatus(retryStatus: grpc.StatusObject) { lastReceivedStatus = retryStatus; if (codeMap.includes(retryStatus.code)) { retry(); } else { receivedMessageCallback(lastReceivedMessage); statusChainNext(lastReceivedStatus); } }, }); retryCall.sendMessage(originMessage); retryCall.halfClose(); } lastReceivedStatus = status; if (codeMap.includes(status.code)) { retry(); } else { // no need to retry, just return what was originally cached receivedMessageCallback(lastReceivedMessage); statusChainNext(lastReceivedStatus); } }, }; next(metadata, newListener); }, } as grpc.Requester); }; }