summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-04-04 20:30:33 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-04-04 20:30:33 +0200
commitde9bd96b5ccda3b9a354cf3e02d49012bdc177d4 (patch)
treefb6bbefc594053c06add419cc9a8961af662ca37
parentd7fcdb2ee11fa2f5a10f911b6916741ea0672734 (diff)
downloadcompiler-de9bd96b5ccda3b9a354cf3e02d49012bdc177d4.tar.gz
compiler-de9bd96b5ccda3b9a354cf3e02d49012bdc177d4.zip
send jobs to compiler
-rw-r--r--worker.js9
1 files changed, 9 insertions, 0 deletions
diff --git a/worker.js b/worker.js
index 6e4e4cb..d0d6ca0 100644
--- a/worker.js
+++ b/worker.js
@@ -20,6 +20,7 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
rabbit = await amqp.connect(RABBITMQ_URI),
ch = await rabbit.createChannel();
await ch.assertQueue("compile", {durable: true});
+ await ch.assertQueue("analyze", {durable: true});
break;
} catch (err) {
console.error(err);
@@ -212,6 +213,14 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
console.log("acking batch");
await ch.ack(msgs.pop(), true); // ack all messages until the last
+ // notify analyzer
+ Promise.all(participant_ext_records.map(async (p) =>
+ await ch.sendToQueue("analyze", new Buffer(p.participant_api_id), {
+ persistent: true,
+ type: "participant"
+ })
+ ));
+
// notify web
await Promise.all([
Promise.all(players.map(async (p) => await ch.publish("amq.topic", "player." + p.name,