summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-04-04 12:06:26 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-04-04 12:06:26 +0200
commit3406a3e76754859cbee257a61f150160597b73fa (patch)
treefb744436fad1af0f342aacd2409bf7d9604b482a
parentb7b7528a0639e0ad4cd1bdac0d021752ddbf0401 (diff)
downloadcompiler-3406a3e76754859cbee257a61f150160597b73fa.tar.gz
compiler-3406a3e76754859cbee257a61f150160597b73fa.zip
wait for rabbit and db to start
-rw-r--r--package.json3
-rw-r--r--worker.js23
2 files changed, 19 insertions, 7 deletions
diff --git a/package.json b/package.json
index 0f84dca..66afc71 100644
--- a/package.json
+++ b/package.json
@@ -6,7 +6,8 @@
"dependencies": {
"amqplib": "^0.5.1",
"mysql": "^2.13.0",
- "sequelize": "^3.30.4"
+ "sequelize": "^3.30.4",
+ "sleep-promise": "^2.0.0"
},
"devDependencies": {},
"scripts": {
diff --git a/worker.js b/worker.js
index 11da36f..a31f97d 100644
--- a/worker.js
+++ b/worker.js
@@ -3,7 +3,8 @@
'use strict';
var amqp = require("amqplib"),
- Seq = require("sequelize");
+ Seq = require("sequelize"),
+ sleep = require("sleep-promise");
var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
DATABASE_URI = process.env.DATABASE_URI || "sqlite:///db.sqlite",
@@ -11,17 +12,27 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
IDLE_TIMEOUT = process.env.PROCESSOR_IDLETIMEOUT || 500; // ms
(async () => {
- let seq = new Seq(DATABASE_URI, { logging: () => {} }),
- model = require("../orm/model")(seq, Seq),
- rabbit = await amqp.connect(RABBITMQ_URI),
- ch = await rabbit.createChannel();
+ let seq, model, rabbit, ch;
+
+ while (true) {
+ try {
+ seq = new Seq(DATABASE_URI, { logging: () => {} }),
+ rabbit = await amqp.connect(RABBITMQ_URI),
+ ch = await rabbit.createChannel();
+ await ch.assertQueue("compile", {durable: true});
+ break;
+ } catch (err) {
+ console.error(err);
+ await sleep(5000);
+ }
+ }
+ model = require("../orm/model")(seq, Seq);
let queue = [],
timer = undefined;
await seq.sync();
- await ch.assertQueue("compile", {durable: true});
// as long as the queue is filled, msg are not ACKed
// server sends as long as there are less than `prefetch` unACKed
await ch.prefetch(BATCHSIZE);