UNPKG

7.83 kBJavaScriptView Raw
1"use strict";
2Object.defineProperty(exports, "__esModule", { value: true });
3exports.ClientRMQ = void 0;
4const logger_service_1 = require("@nestjs/common/services/logger.service");
5const load_package_util_1 = require("@nestjs/common/utils/load-package.util");
6const random_string_generator_util_1 = require("@nestjs/common/utils/random-string-generator.util");
7const events_1 = require("events");
8const rxjs_1 = require("rxjs");
9const operators_1 = require("rxjs/operators");
10const constants_1 = require("../constants");
11const rmq_record_serializer_1 = require("../serializers/rmq-record.serializer");
12const client_proxy_1 = require("./client-proxy");
13let rqmPackage = {};
14const REPLY_QUEUE = 'amq.rabbitmq.reply-to';
15class ClientRMQ extends client_proxy_1.ClientProxy {
16 constructor(options) {
17 super();
18 this.options = options;
19 this.logger = new logger_service_1.Logger(client_proxy_1.ClientProxy.name);
20 this.client = null;
21 this.channel = null;
22 this.urls = this.getOptionsProp(this.options, 'urls') || [constants_1.RQM_DEFAULT_URL];
23 this.queue =
24 this.getOptionsProp(this.options, 'queue') || constants_1.RQM_DEFAULT_QUEUE;
25 this.queueOptions =
26 this.getOptionsProp(this.options, 'queueOptions') ||
27 constants_1.RQM_DEFAULT_QUEUE_OPTIONS;
28 this.replyQueue =
29 this.getOptionsProp(this.options, 'replyQueue') || REPLY_QUEUE;
30 this.persistent =
31 this.getOptionsProp(this.options, 'persistent') || constants_1.RQM_DEFAULT_PERSISTENT;
32 (0, load_package_util_1.loadPackage)('amqplib', ClientRMQ.name, () => require('amqplib'));
33 rqmPackage = (0, load_package_util_1.loadPackage)('amqp-connection-manager', ClientRMQ.name, () => require('amqp-connection-manager'));
34 this.initializeSerializer(options);
35 this.initializeDeserializer(options);
36 }
37 close() {
38 this.channel && this.channel.close();
39 this.client && this.client.close();
40 this.channel = null;
41 this.client = null;
42 }
43 connect() {
44 if (this.client) {
45 return this.connection;
46 }
47 this.client = this.createClient();
48 this.handleError(this.client);
49 this.handleDisconnectError(this.client);
50 const connect$ = this.connect$(this.client);
51 this.connection = (0, rxjs_1.lastValueFrom)(this.mergeDisconnectEvent(this.client, connect$).pipe((0, operators_1.switchMap)(() => this.createChannel()), (0, operators_1.share)())).catch(err => {
52 if (err instanceof rxjs_1.EmptyError) {
53 return;
54 }
55 throw err;
56 });
57 return this.connection;
58 }
59 createChannel() {
60 return new Promise(resolve => {
61 this.channel = this.client.createChannel({
62 json: false,
63 setup: (channel) => this.setupChannel(channel, resolve),
64 });
65 });
66 }
67 createClient() {
68 const socketOptions = this.getOptionsProp(this.options, 'socketOptions');
69 return rqmPackage.connect(this.urls, {
70 connectionOptions: socketOptions,
71 });
72 }
73 mergeDisconnectEvent(instance, source$) {
74 const eventToError = (eventType) => (0, rxjs_1.fromEvent)(instance, eventType).pipe((0, operators_1.map)((err) => {
75 throw err;
76 }));
77 const disconnect$ = eventToError(constants_1.DISCONNECT_EVENT);
78 const urls = this.getOptionsProp(this.options, 'urls', []);
79 const connectFailed$ = eventToError(constants_1.CONNECT_FAILED_EVENT).pipe((0, operators_1.retryWhen)(e => e.pipe((0, operators_1.scan)((errorCount, error) => {
80 if (urls.indexOf(error.url) >= urls.length - 1) {
81 throw error;
82 }
83 return errorCount + 1;
84 }, 0))));
85 return (0, rxjs_1.merge)(source$, disconnect$, connectFailed$).pipe((0, operators_1.first)());
86 }
87 async setupChannel(channel, resolve) {
88 const prefetchCount = this.getOptionsProp(this.options, 'prefetchCount') ||
89 constants_1.RQM_DEFAULT_PREFETCH_COUNT;
90 const isGlobalPrefetchCount = this.getOptionsProp(this.options, 'isGlobalPrefetchCount') ||
91 constants_1.RQM_DEFAULT_IS_GLOBAL_PREFETCH_COUNT;
92 await channel.assertQueue(this.queue, this.queueOptions);
93 await channel.prefetch(prefetchCount, isGlobalPrefetchCount);
94 this.responseEmitter = new events_1.EventEmitter();
95 this.responseEmitter.setMaxListeners(0);
96 await this.consumeChannel(channel);
97 resolve();
98 }
99 async consumeChannel(channel) {
100 const noAck = this.getOptionsProp(this.options, 'noAck', constants_1.RQM_DEFAULT_NOACK);
101 await channel.consume(this.replyQueue, (msg) => this.responseEmitter.emit(msg.properties.correlationId, msg), {
102 noAck,
103 });
104 }
105 handleError(client) {
106 client.addListener(constants_1.ERROR_EVENT, (err) => this.logger.error(err));
107 }
108 handleDisconnectError(client) {
109 client.addListener(constants_1.DISCONNECT_EVENT, (err) => {
110 this.logger.error(constants_1.DISCONNECTED_RMQ_MESSAGE);
111 this.logger.error(err);
112 this.close();
113 });
114 }
115 async handleMessage(packet, callback) {
116 const { err, response, isDisposed } = await this.deserializer.deserialize(packet);
117 if (isDisposed || err) {
118 callback({
119 err,
120 response,
121 isDisposed: true,
122 });
123 }
124 callback({
125 err,
126 response,
127 });
128 }
129 publish(message, callback) {
130 try {
131 const correlationId = (0, random_string_generator_util_1.randomStringGenerator)();
132 const listener = ({ content }) => this.handleMessage(JSON.parse(content.toString()), callback);
133 Object.assign(message, { id: correlationId });
134 const serializedPacket = this.serializer.serialize(message);
135 const options = serializedPacket.options;
136 delete serializedPacket.options;
137 this.responseEmitter.on(correlationId, listener);
138 this.channel.sendToQueue(this.queue, Buffer.from(JSON.stringify(serializedPacket)), Object.assign(Object.assign({ replyTo: this.replyQueue, persistent: this.persistent }, options), { headers: this.mergeHeaders(options === null || options === void 0 ? void 0 : options.headers), correlationId }));
139 return () => this.responseEmitter.removeListener(correlationId, listener);
140 }
141 catch (err) {
142 callback({ err });
143 }
144 }
145 dispatchEvent(packet) {
146 const serializedPacket = this.serializer.serialize(packet);
147 const options = serializedPacket.options;
148 delete serializedPacket.options;
149 return new Promise((resolve, reject) => this.channel.sendToQueue(this.queue, Buffer.from(JSON.stringify(serializedPacket)), Object.assign(Object.assign({ persistent: this.persistent }, options), { headers: this.mergeHeaders(options === null || options === void 0 ? void 0 : options.headers) }), (err) => (err ? reject(err) : resolve())));
150 }
151 initializeSerializer(options) {
152 var _a;
153 this.serializer = (_a = options === null || options === void 0 ? void 0 : options.serializer) !== null && _a !== void 0 ? _a : new rmq_record_serializer_1.RmqRecordSerializer();
154 }
155 mergeHeaders(requestHeaders) {
156 var _a, _b;
157 if (!requestHeaders && !((_a = this.options) === null || _a === void 0 ? void 0 : _a.headers)) {
158 return undefined;
159 }
160 return Object.assign(Object.assign({}, (_b = this.options) === null || _b === void 0 ? void 0 : _b.headers), requestHeaders);
161 }
162}
163exports.ClientRMQ = ClientRMQ;