import * as grpc from 'grpc'; import * as path from 'path'; import { from, Subject } from "rxjs"; import { GrpcRequest } from "../src/request"; import { TestServiceClient } from "../fixtures/generated/protos/test_grpc_pb"; import { TestMessage } from "../fixtures/generated/protos/test_pb"; import { mergeMap, reduce } from "rxjs/operators"; import { Client } from "../src/client"; const {createMockServer} = require("grpc-mock"); const portfinder = require('portfinder'); const killPort = require('kill-port'); function mockRules(otherRules: Array = []) { return [ { method: "Unary", input: { name: "test" }, output: { name: "unary response" } }, { method: "ClientStream", streamType: "client", stream: [ { input: { name: "test" } }, ], output: { name: "client stream response" } }, { method: "ServerStream", streamType: "server", stream: [ { output: { name: "server stream response" } }, { output: { name: "server stream response" } }, ], input: { name: "test" } }, { method: "BidiStream", streamType: "mutual", stream: [ { input: { name: "test-1" }, output: { name: "bidi outputPath 1" } }, { input: { name: "test-2" }, output: { name: "bidi outputPath 2" } }, ] }, ...otherRules, ]; } describe("Construct and configure client correctly", () => { let port = 0; let grpcAddress = ""; beforeEach(async () => { port = await portfinder.getPortPromise() as number; grpcAddress = `0.0.0.0:${port}`; }); it("Will set configurations to a client", () => { const credentials = grpc.credentials.createInsecure() const client = Client.create(TestServiceClient, { address: grpcAddress, credentials: credentials, size: 1, options: { interceptors: [], 'grpc.keepalive_time_ms': 10 * 1000, 'grpc.keepalive_permit_without_calls': 1, 'grpc.http2.min_time_between_pings_ms': 5 * 1000, 'grpc.http2.max_pings_without_data': 0, 'grpc.initial_reconnect_backoff_ms': 300, 'grpc.min_reconnect_backoff_ms': 1000, 'grpc.max_reconnect_backoff_ms': 5000, } }) expect(client.clientOptions()).toEqual({ "address": "0.0.0.0:8000", "credentials": credentials, size: 1, options: { interceptors: [], 'grpc.keepalive_time_ms': 10 * 1000, 'grpc.keepalive_permit_without_calls': 1, 'grpc.http2.min_time_between_pings_ms': 5 * 1000, 'grpc.http2.max_pings_without_data': 0, 'grpc.initial_reconnect_backoff_ms': 300, 'grpc.min_reconnect_backoff_ms': 1000, 'grpc.max_reconnect_backoff_ms': 5000, }, }) }) it("Will use default client configurations", () => { const credentials = grpc.credentials.createInsecure() Client.setClientDefaults({ options: { interceptors: [], 'grpc.keepalive_time_ms': 10 * 1000, 'grpc.keepalive_permit_without_calls': 1, 'grpc.http2.min_time_between_pings_ms': 5 * 1000, 'grpc.http2.max_pings_without_data': 0, 'grpc.initial_reconnect_backoff_ms': 300, 'grpc.min_reconnect_backoff_ms': 1000, 'grpc.max_reconnect_backoff_ms': 5000, } }) const client = Client.create(TestServiceClient, { address: grpcAddress, "credentials": credentials, }) expect(client.clientOptions()).toEqual({ "address": "0.0.0.0:8000", "credentials": credentials, size: 1, options: { interceptors: [], 'grpc.keepalive_time_ms': 10 * 1000, 'grpc.keepalive_permit_without_calls': 1, 'grpc.http2.min_time_between_pings_ms': 5 * 1000, 'grpc.http2.max_pings_without_data': 0, 'grpc.initial_reconnect_backoff_ms': 300, 'grpc.min_reconnect_backoff_ms': 1000, 'grpc.max_reconnect_backoff_ms': 5000, }, }) }) }) describe("grpc reactive requests", () => { let port = 0; let grpcAddress = ""; beforeEach(async () => { port = await portfinder.getPortPromise() as number; grpcAddress = `0.0.0.0:${port}`; }); afterEach(() => { return killPort(port); }); it("Will send a grpc unary call", async () => { const mockServer = createMockServer({ protoPath: path.resolve(__dirname, '..', 'fixtures', 'protos', 'test.proto'), packageName: "test", serviceName: "TestService", rules: mockRules(), }); mockServer.listen(grpcAddress); const plainClient = new TestServiceClient(grpcAddress, grpc.credentials.createInsecure()); const grpcRequest = GrpcRequest.create(plainClient); const message = new TestMessage(); message.setName("test"); const response = await grpcRequest .withMetadata({ unary: "metadata" }) .unary(message) .toPromise(); expect(response.getName()).toEqual("unary response"); return mockServer.close(true); }); it ("Will send multiple grpc unary calls", async (done) => { const mockServer = createMockServer({ protoPath: path.resolve(__dirname, '..', 'fixtures', 'protos', 'test.proto'), packageName: "test", serviceName: "TestService", rules: mockRules(), }); mockServer.listen(grpcAddress); const plainClient = new TestServiceClient(grpcAddress, grpc.credentials.createInsecure()); const grpcRequest = GrpcRequest.create(plainClient); const message = new TestMessage(); message.setName("test"); // Promise based const responsePromise = await from([1,2]) .pipe(mergeMap(() => grpcRequest.unary(message))) .pipe(reduce((acc,val) => acc.concat(val), [] as TestMessage[])) .toPromise(); expect(responsePromise.map(message => message.getName())) .toEqual(["unary response", "unary response"]); const spy = jest.fn(); from([1,2]) .pipe(mergeMap(() => grpcRequest.unary(message))) .subscribe({ next(message) { spy(message.getName()); }, complete() { expect(spy).toBeCalledTimes(2); expect(spy).toBeCalledWith("unary response"); mockServer.close(true); done(); } }); }); it ("will send server streaming call", async (done) => { const mockServer = createMockServer({ protoPath: path.resolve(__dirname, '..', 'fixtures', 'protos', 'test.proto'), packageName: "test", serviceName: "TestService", rules: mockRules(), }); mockServer.listen(grpcAddress); const plainClient = new TestServiceClient(grpcAddress, grpc.credentials.createInsecure()); const grpcRequest = GrpcRequest.create(plainClient); const message = new TestMessage(); message.setName("test"); const spy = jest.fn(); grpcRequest.serverStream(message) .subscribe({ next(c) { spy(c.getName()); }, complete() { expect(spy).toBeCalledTimes(2); expect(spy).toBeCalledWith("server stream response"); mockServer.close(true); done(); } }); }); it ("will send client streaming call", async (done) => { const mockServer = createMockServer({ protoPath: path.resolve(__dirname, '..', 'fixtures', 'protos', 'test.proto'), packageName: "test", serviceName: "TestService", rules: mockRules(), }); mockServer.listen(grpcAddress); const plainClient = new TestServiceClient(grpcAddress, grpc.credentials.createInsecure()); const grpcRequest = GrpcRequest.create(plainClient); const message = new TestMessage(); message.setName("test"); const streamClient = new Subject(); const spy = jest.fn(); grpcRequest.clientStream(streamClient) .subscribe({ next(c) { spy(c.getName()); }, complete() { mockServer.close(true); expect(spy).toBeCalledTimes(1); expect(spy).toBeCalledWith("client stream response"); done(); } }); streamClient.next(message); }); it ("will send a bi-directional streaming call", async (done) => { const mockServer = createMockServer({ protoPath: path.resolve(__dirname, '..', 'fixtures', 'protos', 'test.proto'), packageName: "test", serviceName: "TestService", rules: mockRules(), }); mockServer.listen(grpcAddress); const plainClient = new TestServiceClient(grpcAddress, grpc.credentials.createInsecure()); const grpcRequest = GrpcRequest.create(plainClient); const message1 = new TestMessage(); message1.setName("test-1"); const message2 = new TestMessage(); message2.setName("test-2"); const streamClient = new Subject(); const spy = jest.fn(); grpcRequest.bidiStream(streamClient) .subscribe({ next(c) { spy(c.getName()); }, complete() { mockServer.close(true); expect(spy).toBeCalledTimes(2); expect(spy).toBeCalledWith("bidi outputPath 1"); expect(spy).toBeCalledWith("bidi outputPath 2"); done(); } }); streamClient.next(message1); streamClient.next(message2); }); it ("Will return an error when the server returns one", async (done) => { const mockServer = createMockServer({ protoPath: path.resolve(__dirname, '..', 'fixtures', 'protos', 'test.proto'), packageName: "test", serviceName: "TestService", rules: mockRules([ { method: "Unary", input: { name: "test-error" }, error: { code: 3, message: "ohoh!"}, } ]), }); mockServer.listen(grpcAddress); const plainClient = new TestServiceClient(grpcAddress, grpc.credentials.createInsecure()); const grpcRequest = GrpcRequest.create(plainClient); const message = new TestMessage(); message.setName("test-error"); try { await from([1,2]) .pipe(mergeMap(() => grpcRequest.unary(message))) .pipe(reduce((acc,val) => acc.concat(val), [] as TestMessage[])) .toPromise(); } catch (e) { expect(e.message).toEqual('3 INVALID_ARGUMENT: ohoh!'); } const spy = jest.fn(); from([1,2]) .pipe(mergeMap(() => grpcRequest.unary(message))) .subscribe({ error(err) { expect(err.message).toEqual('3 INVALID_ARGUMENT: ohoh!'); expect(spy).toBeCalledTimes(0); mockServer.close(true); done(); }, next(message) { spy(message.getName()); }, }); }); it ("will send server streaming call with async iterator", async () => { const mockServer = createMockServer({ protoPath: path.resolve(__dirname, '..', 'fixtures', 'protos', 'test.proto'), packageName: "test", serviceName: "TestService", rules: mockRules(), }); mockServer.listen(grpcAddress); const plainClient = new TestServiceClient(grpcAddress, grpc.credentials.createInsecure()); const grpcRequest = GrpcRequest.create(plainClient); const message = new TestMessage(); message.setName("test"); const spy = jest.fn(); const messages = grpcRequest .serverStream(message) .toIterator(); for await (const message of messages) { spy(message.getName()); } expect(spy).toBeCalledTimes(2); expect(spy).toBeCalledWith("server stream response"); mockServer.close(true); }); });