|
1 | 1 | import amqp from 'amqplib'; |
| 2 | +import queues from '../queues/index.mjs'; |
| 3 | +import logger from '../utils/logger.mjs'; |
| 4 | + |
| 5 | +let connection; |
| 6 | +let channel; |
| 7 | + |
| 8 | +// connect to RabbitMQ server. |
| 9 | +const connect = async () => { |
| 10 | + logger.info('[AMQP] Connecting...'); |
| 11 | + connection = await amqp.connect(`${process.env.RABBITMQ_HOST}:${process.env.RABBITMQ_PORT}`); |
| 12 | + connection.on('error', err => { |
| 13 | + if (err.message !== 'Connection closing') { |
| 14 | + logger.error(`[AMQP] Connection error: ${err.message}`); |
| 15 | + return setTimeout(connect, 1000); |
| 16 | + } |
| 17 | + }); |
| 18 | + connection.on('close', () => { |
| 19 | + logger.debug('[AMQP] Reconnecting'); |
| 20 | + return setTimeout(connect, 1000); |
| 21 | + }); |
| 22 | + channel = await connection.createChannel(); |
| 23 | + channel.prefetch(1); |
| 24 | + queues(); |
| 25 | +}; |
| 26 | + |
| 27 | +// publish message to a queue |
| 28 | +const sendMessage = async (queue, data, options = {}) => { |
| 29 | + channel.sendToQueue(queue, Buffer.from(JSON.stringify(data)), options); |
| 30 | +}; |
| 31 | + |
| 32 | +// only called from test-suite ./test/queue.mjs |
| 33 | +const cancelChannelConsume = consumerTag => { |
| 34 | + channel.cancel(consumerTag); |
| 35 | +}; |
| 36 | + |
| 37 | +const sendACK = message => { |
| 38 | + channel.ack(message); |
| 39 | +}; |
| 40 | + |
| 41 | +// only called from test-suite ./test/queue.mjs |
| 42 | +const sendNACK = message => { |
| 43 | + channel.nack(message); |
| 44 | +}; |
| 45 | + |
| 46 | +/* |
| 47 | + * Consumer: receive message from a queue |
| 48 | + */ |
| 49 | +const receiveMessage = async (queue, callback) => { |
| 50 | + await channel.assertQueue(queue); |
| 51 | + channel.consume(queue, callback, { noAck: false }); |
| 52 | +}; |
| 53 | + |
| 54 | +// only called from test-suite ./test/queue.mjs |
| 55 | +const listenToReplyQueue = async (queue, correlationId, callback) => { |
| 56 | + logger.info(`[AMQP] Listening to reply queue: ${queue}`); |
| 57 | + receiveMessage(queue, message => { |
| 58 | + if (message.properties.correlationId !== correlationId) { |
| 59 | + logger.debug(`[AMQP] Sending NACK due to different correlation id: ${JSON.stringify({ message: message.properties.correlationId, param: correlationId })}`); |
| 60 | + return sendNACK(message); |
| 61 | + } |
| 62 | + cancelChannelConsume(message.fields.consumerTag); |
| 63 | + sendACK(message); |
| 64 | + |
| 65 | + const response = JSON.parse(message.content.toString()); |
| 66 | + response.type = message.properties.type; |
| 67 | + |
| 68 | + return callback(response); |
| 69 | + }); |
| 70 | +}; |
| 71 | + |
| 72 | +// only called from test-suite ./test/queue.mjs |
| 73 | +const close = async () => { |
| 74 | + await channel.close(); |
| 75 | + await connection.close(); |
| 76 | +}; |
2 | 77 |
|
3 | 78 | export default { |
4 | | - // connect to RabbitMQ server. |
5 | | - async connect() { |
6 | | - this.connection = await amqp.connect( |
7 | | - `${process.env.RABBITMQ_HOST}:${process.env.RABBITMQ_PORT}`, |
8 | | - ); |
9 | | - this.channel = await this.connection.createChannel(); |
10 | | - this.channel.prefetch(1); |
11 | | - }, |
12 | | - |
13 | | - // publish message to a queue |
14 | | - async sendMessage(queue, data, options = {}) { |
15 | | - this.channel.sendToQueue(queue, Buffer.from(JSON.stringify(data)), options); |
16 | | - }, |
17 | | - |
18 | | - // only called from test-suite ./test/queue.mjs |
19 | | - cancelChannelConsume(consumerTag) { |
20 | | - this.channel.cancel(consumerTag); |
21 | | - }, |
22 | | - |
23 | | - sendACK(message) { |
24 | | - this.channel.ack(message); |
25 | | - }, |
26 | | - |
27 | | - // only called from test-suite ./test/queue.mjs |
28 | | - sendNACK(message) { |
29 | | - this.channel.nack(message); |
30 | | - }, |
31 | | - |
32 | | - /* |
33 | | - * Consumer: receive message from a queue |
34 | | - */ |
35 | | - async receiveMessage(queue, callback) { |
36 | | - await this.channel.assertQueue(queue); |
37 | | - this.channel.consume(queue, callback); |
38 | | - }, |
39 | | - |
40 | | - // only called from test-suite ./test/queue.mjs |
41 | | - listenToReplyQueue(queue, correlationId, callback) { |
42 | | - this.receiveMessage(queue, message => { |
43 | | - if (message.properties.correlationId !== correlationId) { |
44 | | - return this.sendNACK(message); |
45 | | - } |
46 | | - this.cancelChannelConsume(message.fields.consumerTag); |
47 | | - this.sendACK(message); |
48 | | - |
49 | | - const response = JSON.parse(message.content.toString()); |
50 | | - response.type = message.properties.type; |
51 | | - |
52 | | - return callback(response); |
53 | | - }); |
54 | | - }, |
55 | | - |
56 | | - // only called from test-suite ./test/queue.mjs |
57 | | - async close() { |
58 | | - await this.channel.close(); |
59 | | - await this.connection.close(); |
60 | | - }, |
| 79 | + connection, |
| 80 | + channel, |
| 81 | + connect, |
| 82 | + sendMessage, |
| 83 | + cancelChannelConsume, |
| 84 | + sendACK, |
| 85 | + sendNACK, |
| 86 | + receiveMessage, |
| 87 | + listenToReplyQueue, |
| 88 | + close, |
61 | 89 | }; |
0 commit comments