summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-05-25 14:29:43 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-05-25 14:29:43 +0200
commit506ad32a251f95fd950596a3efb0c595485613b7 (patch)
treef05993938175026bcd8f17e00fb597976d635431
parentf2bdd5c9c41796f1887eea94576de67af6bc88e0 (diff)
downloadcruncher-506ad32a251f95fd950596a3efb0c595485613b7.tar.gz
cruncher-506ad32a251f95fd950596a3efb0c595485613b7.zip
allow queue configuration via env var
-rw-r--r--worker.js7
1 files changed, 4 insertions, 3 deletions
diff --git a/worker.js b/worker.js
index c40b216..0934392 100644
--- a/worker.js
+++ b/worker.js
@@ -20,6 +20,7 @@ const amqp = require("amqplib"),
const RABBITMQ_URI = process.env.RABBITMQ_URI,
DATABASE_URI = process.env.DATABASE_URI,
+ QUEUE = process.env.QUEUE || "crunch",
LOGGLY_TOKEN = process.env.LOGGLY_TOKEN,
// size of connection pool
MAXCONNS = parseInt(process.env.MAXCONNS) || 3,
@@ -41,7 +42,7 @@ if (LOGGLY_TOKEN)
logger.add(winston.transports.Loggly, {
inputToken: LOGGLY_TOKEN,
subdomain: "kvahuja",
- tags: ["backend", "cruncher"],
+ tags: ["backend", "cruncher", QUEUE],
json: true
});
@@ -58,7 +59,7 @@ if (LOGGLY_TOKEN)
}),
rabbit = await amqp.connect(RABBITMQ_URI, { heartbeat: 320 }),
ch = await rabbit.createChannel();
- await ch.assertQueue("crunch", {durable: true});
+ await ch.assertQueue(QUEUE, {durable: true});
break;
} catch (err) {
logger.error("Error connecting", err);
@@ -83,7 +84,7 @@ if (LOGGLY_TOKEN)
// set maximum allowed number of unacked msgs
// TODO maybe split queues by type
await ch.prefetch(BATCHSIZE);
- ch.consume("crunch", (msg) => {
+ ch.consume(QUEUE, (msg) => {
const api_id = msg.content.toString();
if (msg.properties.type == "global")
participants_global.add(api_id);