summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-04-04 12:06:18 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-04-04 12:06:18 +0200
commit63b41da0ba9e6e14bd4b5f8c72992d72e156c596 (patch)
treed31e6897abffc977f42a73d219b22d204e1c3fa5
parent8955e221c9656a07f3a5da4ddd4cfe2860ed0590 (diff)
downloadprocessor-63b41da0ba9e6e14bd4b5f8c72992d72e156c596.tar.gz
processor-63b41da0ba9e6e14bd4b5f8c72992d72e156c596.zip
wait for rabbit and db to start
-rw-r--r--package.json1
-rw-r--r--worker.js26
2 files changed, 20 insertions, 7 deletions
diff --git a/package.json b/package.json
index aeeb1e7..1c85324 100644
--- a/package.json
+++ b/package.json
@@ -7,6 +7,7 @@
"amqplib": "^0.5.1",
"mysql": "^2.13.0",
"sequelize": "^3.30.4",
+ "sleep-promise": "^2.0.0",
"snakecase-keys": "^1.1.0"
},
"devDependencies": {},
diff --git a/worker.js b/worker.js
index c5c0582..a2287b9 100644
--- a/worker.js
+++ b/worker.js
@@ -5,7 +5,8 @@
var amqp = require("amqplib"),
Seq = require("sequelize"),
snakeCaseKeys = require("snakecase-keys"),
- item_name_map = require("../orm/items");
+ item_name_map = require("../orm/items"),
+ sleep = require("sleep-promise");
var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
DATABASE_URI = process.env.DATABASE_URI || "sqlite:///db.sqlite",
@@ -13,10 +14,23 @@ 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("process", {durable: true});
+ 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;
@@ -32,8 +46,6 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
await model.Item.findAll()
.map((item) => item_db_map[item.name] = item.id);
- await ch.assertQueue("process", {durable: true});
- 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);