import * as grpc from "grpc"; import { Metadata, MethodDefinition } from "grpc"; import { Observable, Subject, Observer } from 'rxjs'; import { Readable, Writable } from "stream"; import { Interceptor } from "./interceptors/chain"; import './observable.iterator'; export type ObservableGrpcCall< Client extends grpc.Client, Property extends keyof Client, Unary = ReturnType< & ((...args: any) => any) & Client[Property]>, ClientStream = ReturnType<& ((...args: any) => any) & Client[Property]>, ServerStream = ReturnType<& ((...args: any) => any) & Client[Property]>, Duplex = ReturnType<& ((...args: any) => any) & Client[Property]>, > = Duplex extends grpc.ClientDuplexStream ? ( subjectStream: Subject, metadata?: Parameters<(() => void) & Client[Property]>[0], grpcOptions?: Parameters<(() => void) & Client[Property]>[1], ) => Observable : ClientStream extends grpc.ClientWritableStream ? ( subjectStream: Subject, metadata?: Parameters<(() => void) & Client[Property]>[0], grpcOptions?: Parameters<(() => void) & Client[Property]>[1], ) => Observable : ServerStream extends grpc.ClientReadableStream ? ( request: (Parameters<(() => void) & Client[Property]>[0]), metadata?: Parameters<(() => void) & Client[Property]>[1], grpcOptions?: Parameters<(() => void) & Client[Property]>[2], ) => Observable : Unary extends grpc.ClientUnaryCall ? ( request: (Parameters<(() => void) & Client[Property]>[0]), metadata?: Parameters<(() => void) & Client[Property]>[1], grpcOptions?: Parameters<(() => void) & Client[Property]>[2], ) => Observable void) & Client[Property]>[3]>[1] > : Client[Property]; type ObservableGrpc< Client extends grpc.Client, ClientMethodNames extends keyof Client = keyof Client, > = { [P in ClientMethodNames]: ObservableGrpcCall; } /** * GrpcRequest */ export class GrpcRequest { protected constructor( protected readonly _client: Client & Record, protected readonly _metadata: Metadata = new Metadata(), protected readonly _options?: grpc.CallOptions ) { return new Proxy(this as (ObservableGrpc & this & Record), { get: (obj, prop) => { if (prop in obj) { return obj[prop as string]; } if (!(prop in _client)) { return; } return (...args: any[]) => { return this.callReactiveGrpcMethod(prop, ...args); } }, }); } static create(client: Client, metadata?: Metadata, options?: grpc.CallOptions) { return new GrpcRequest(client, metadata, options) as (GrpcRequest & ObservableGrpc); } get client() { return this._client; } withMetadata(metadata: Record) { for (let key in metadata) { if (metadata.hasOwnProperty(key)) { this._metadata.set(key, metadata[key]); } } return this; } protected callReactiveGrpcMethod(methodName: string | number | symbol, ...args: any[]) { return new Observable((observer) => { const method = this._client[methodName as string] as MethodDefinition; // when the request is a client-stream, we expect a Subject // to be passed in. Although that is a parameter that is // required only by our grpc wrapper, for that reason we'll extract it // and save it. let clientStreamingRequest: Subject | undefined; if (method.requestStream) { clientStreamingRequest = args[0]; args = args.slice(0, args.length - 1); } args = this.gRPCCallArguments(...args); // use a callback handler for unary calls. if (!method.responseStream) { args.push(this.unaryHandler(observer)); } // Make the grpc call const call: grpc.Call = this._client[methodName as string](...args); const onData = (data: any) => { observer.next(data); }; const onError = (error: Error) => { observer.error(error); }; const onEnd = () => { call.removeListener('error', onError); observer.complete(); }; const originalUnsub = observer.unsubscribe; observer.unsubscribe = function(...args) { // Remove observable error listener call.removeListener('error', onError); // Add silent error handler. Avoids errors if stream responses have already come in // and we are delaying observable response, like when using delay() call.on('error', () => {}); // Like end, but does not throw error // @see https://nodejs.org/api/stream.html#stream_readable_streams // note GRPC stream extends readable_stream if (call instanceof Readable) { call.destroy(); } originalUnsub.call(observer, ...args); call.cancel(); }; if (method.responseStream) { call.on("data", onData); call.on("error", onError); call.on("end", onEnd); } if (method.requestStream) { const observable = clientStreamingRequest; if (!(observable instanceof Subject)) { return observer.error( new Error('Subject required to push client data') ); } observable.subscribe({ next: data => { if (call instanceof Writable) { call.write(data); } }, error: () => { call.cancel(); }, complete: () => { if (call instanceof Writable) { call.end(); } }, }); } }); } /** * Unary call handler * @param observer */ protected unaryHandler(observer: Observer) { return (err: Error, response: T) => { if (err) { observer.error(err); } else { observer.next(response); } observer.complete(); } } /** * get the appropriate grpc call arguments * @param args */ protected gRPCCallArguments(...args: any[]): any[] { // remove passed handler. if (args.length > 0 && typeof args[args.length - 1] === "function") { args = args.slice(0, args.length - 1); } // detect passed metadata and add it to the default const argMetadataIndex = args.findIndex(arg => arg instanceof Metadata); if (argMetadataIndex > 0) { const passedMetadata = args[argMetadataIndex]; const valueMap = passedMetadata.getMap(); Object.keys(valueMap).forEach(key => { this._metadata.add(key, valueMap[key]); }); args[argMetadataIndex] = this._metadata; } else { // no metadata passed i'll use the default metadata of this class // shifting the parameters. Metadata could only be present // within the 1st or 2nd param const newArgs = []; if (args.length === 0) { newArgs.push(this._metadata); } if (args.length === 1) { newArgs.push(...args, this._metadata); } // metadata and maybe options if (args.length > 1) { let restParams = args.slice(1, args.length - 1); // grpc Options if (restParams.length == 2) { const options = restParams[1]; restParams[1] = { ...(this._options ? this._options : {}), ...(options ? options : {}), }; } else { restParams.push(this._options); } newArgs.push(args[0], this._metadata, ...restParams); } // assign the arguments args = newArgs; } return args; } }