diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-04-04 12:06:26 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-04-04 12:06:26 +0200 |
| commit | 3406a3e76754859cbee257a61f150160597b73fa (patch) | |
| tree | fb744436fad1af0f342aacd2409bf7d9604b482a | |
| parent | b7b7528a0639e0ad4cd1bdac0d021752ddbf0401 (diff) | |
| download | compiler-3406a3e76754859cbee257a61f150160597b73fa.tar.gz compiler-3406a3e76754859cbee257a61f150160597b73fa.zip | |
wait for rabbit and db to start
| -rw-r--r-- | package.json | 3 | ||||
| -rw-r--r-- | worker.js | 23 |
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": { @@ -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); |
