|
| 1 | +'use strict'; |
| 2 | + |
| 3 | +const amqp = require('amqplib/callback_api'); |
| 4 | +const Ably = require('ably'); |
| 5 | + |
| 6 | +/* Start the worker that consumes from the AMQP QUEUE */ |
| 7 | +exports.start = function(apiKey, answersChannelName, queueName, queueEndpoint) { |
| 8 | + const appId = apiKey.split('.')[0]; |
| 9 | + const queue = appId + ":" + queueName; |
| 10 | + const endpoint = queueEndpoint; |
| 11 | + const url = 'amqps://' + apiKey + '@' + endpoint; |
| 12 | + const rest = new Ably.Rest({ key: apiKey }); |
| 13 | + const answersChannel = rest.channels.get(answersChannelName); |
| 14 | + |
| 15 | + /* Connect to Ably queue */ |
| 16 | + amqp.connect(url, (err, conn) => { |
| 17 | + if (err) { |
| 18 | + console.error('worker:', 'Queue error!', err); |
| 19 | + return; |
| 20 | + } |
| 21 | + console.log('worker:', 'Connected to AMQP endpoint', endpoint); |
| 22 | + |
| 23 | + /* Create a communication channel */ |
| 24 | + conn.createChannel((err, ch) => { |
| 25 | + if (err) { |
| 26 | + console.error('worker:', 'Queue error!', err); |
| 27 | + return; |
| 28 | + } |
| 29 | + console.log('worker:', 'Waiting for messages'); |
| 30 | + |
| 31 | + /* Wait for messages published to the Ably Reactor queue */ |
| 32 | + ch.consume(queue, (item) => { |
| 33 | + const decodedEnvelope = JSON.parse(item.content); |
| 34 | + |
| 35 | + const messages = Ably.Realtime.Message.fromEncodedArray(decodedEnvelope.messages); |
| 36 | + messages.forEach(function(message) { |
| 37 | + console.log('worker:', 'Received question', message.data, ' - about to ask Wolfram'); |
| 38 | + answersChannel.publish('answer', 'Placeholder until Wolfram connected', function(err) { |
| 39 | + if (err) { |
| 40 | + console.error('worker:', 'Failed to publish question', message.data, ' - err:', JSON.stringify(err)); |
| 41 | + } |
| 42 | + }) |
| 43 | + }); |
| 44 | + |
| 45 | + /* Remove message from queue */ |
| 46 | + ch.ack(item); |
| 47 | + }); |
| 48 | + }); |
| 49 | + |
| 50 | + conn.on('error', function(err) { console.error('worker:', 'Connection error!', err); }); |
| 51 | + }); |
| 52 | +}; |
0 commit comments