1 | "use strict";
|
2 | Object.defineProperty(exports, "__esModule", { value: true });
|
3 | exports.ClientRMQ = void 0;
|
4 | const logger_service_1 = require("@nestjs/common/services/logger.service");
|
5 | const load_package_util_1 = require("@nestjs/common/utils/load-package.util");
|
6 | const random_string_generator_util_1 = require("@nestjs/common/utils/random-string-generator.util");
|
7 | const events_1 = require("events");
|
8 | const rxjs_1 = require("rxjs");
|
9 | const operators_1 = require("rxjs/operators");
|
10 | const constants_1 = require("../constants");
|
11 | const rmq_record_serializer_1 = require("../serializers/rmq-record.serializer");
|
12 | const client_proxy_1 = require("./client-proxy");
|
13 | let rqmPackage = {};
|
14 | const REPLY_QUEUE = 'amq.rabbitmq.reply-to';
|
15 | class 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 | }
|
163 | exports.ClientRMQ = ClientRMQ;
|