1 | "use strict";
|
2 | Object.defineProperty(exports, "__esModule", { value: true });
|
3 | exports.ClientTCP = void 0;
|
4 | const common_1 = require("@nestjs/common");
|
5 | const net = require("net");
|
6 | const rxjs_1 = require("rxjs");
|
7 | const operators_1 = require("rxjs/operators");
|
8 | const constants_1 = require("../constants");
|
9 | const helpers_1 = require("../helpers");
|
10 | const client_proxy_1 = require("./client-proxy");
|
11 | class ClientTCP extends client_proxy_1.ClientProxy {
|
12 | constructor(options) {
|
13 | super();
|
14 | this.logger = new common_1.Logger(ClientTCP.name);
|
15 | this.isConnected = false;
|
16 | this.port = this.getOptionsProp(options, 'port') || constants_1.TCP_DEFAULT_PORT;
|
17 | this.host = this.getOptionsProp(options, 'host') || constants_1.TCP_DEFAULT_HOST;
|
18 | this.socketClass =
|
19 | this.getOptionsProp(options, 'socketClass') || helpers_1.JsonSocket;
|
20 | this.initializeSerializer(options);
|
21 | this.initializeDeserializer(options);
|
22 | }
|
23 | connect() {
|
24 | if (this.connection) {
|
25 | return this.connection;
|
26 | }
|
27 | this.socket = this.createSocket();
|
28 | this.bindEvents(this.socket);
|
29 | const source$ = this.connect$(this.socket.netSocket).pipe((0, operators_1.tap)(() => {
|
30 | this.isConnected = true;
|
31 | this.socket.on(constants_1.MESSAGE_EVENT, (buffer) => this.handleResponse(buffer));
|
32 | }), (0, operators_1.share)());
|
33 | this.socket.connect(this.port, this.host);
|
34 | this.connection = (0, rxjs_1.lastValueFrom)(source$).catch(err => {
|
35 | if (err instanceof rxjs_1.EmptyError) {
|
36 | return;
|
37 | }
|
38 | throw err;
|
39 | });
|
40 | return this.connection;
|
41 | }
|
42 | async handleResponse(buffer) {
|
43 | const { err, response, isDisposed, id } = await this.deserializer.deserialize(buffer);
|
44 | const callback = this.routingMap.get(id);
|
45 | if (!callback) {
|
46 | return undefined;
|
47 | }
|
48 | if (isDisposed || err) {
|
49 | return callback({
|
50 | err,
|
51 | response,
|
52 | isDisposed: true,
|
53 | });
|
54 | }
|
55 | callback({
|
56 | err,
|
57 | response,
|
58 | });
|
59 | }
|
60 | createSocket() {
|
61 | return new this.socketClass(new net.Socket());
|
62 | }
|
63 | close() {
|
64 | this.socket && this.socket.end();
|
65 | this.handleClose();
|
66 | }
|
67 | bindEvents(socket) {
|
68 | socket.on(constants_1.ERROR_EVENT, (err) => err.code !== constants_1.ECONNREFUSED && this.handleError(err));
|
69 | socket.on(constants_1.CLOSE_EVENT, () => this.handleClose());
|
70 | }
|
71 | handleError(err) {
|
72 | this.logger.error(err);
|
73 | }
|
74 | handleClose() {
|
75 | this.isConnected = false;
|
76 | this.socket = null;
|
77 | this.connection = undefined;
|
78 | }
|
79 | publish(partialPacket, callback) {
|
80 | try {
|
81 | const packet = this.assignPacketId(partialPacket);
|
82 | const serializedPacket = this.serializer.serialize(packet);
|
83 | this.routingMap.set(packet.id, callback);
|
84 | this.socket.sendMessage(serializedPacket);
|
85 | return () => this.routingMap.delete(packet.id);
|
86 | }
|
87 | catch (err) {
|
88 | callback({ err });
|
89 | }
|
90 | }
|
91 | async dispatchEvent(packet) {
|
92 | const pattern = this.normalizePattern(packet.pattern);
|
93 | const serializedPacket = this.serializer.serialize(Object.assign(Object.assign({}, packet), { pattern }));
|
94 | return this.socket.sendMessage(serializedPacket);
|
95 | }
|
96 | }
|
97 | exports.ClientTCP = ClientTCP;
|