diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-08-25 23:06:27 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-08-25 23:06:27 +0200 |
| commit | 93aa94dc157db3c722e1f1921df983bde9c4bd00 (patch) | |
| tree | 2a50815728ba68c88bd18d2b901053cf9176ac93 | |
| parent | d5f24344b7ae5d43d81cf9fecd7efee9515f1806 (diff) | |
| download | processor-93aa94dc157db3c722e1f1921df983bde9c4bd00.tar.gz processor-93aa94dc157db3c722e1f1921df983bde9c4bd00.zip | |
fix exceptions when forwarding to failed queue
| -rw-r--r-- | worker.js | 11 |
1 files changed, 7 insertions, 4 deletions
@@ -195,7 +195,8 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { // it is really `"null"`. // reject invalid matches (handling API bugs) await ch.nack(msg, false, false); - await ch.sendToQueue(QUEUE + "_failed", msg.content, { + await ch.sendToQueue(QUEUE + "_failed", + new Buffer(JSON.stringify(msg.content)), { persistent: true, type: msg.properties.type, headers: msg.properties.headers @@ -273,7 +274,7 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { // notify follow up services if (DOANALYZEMATCH) { await Promise.each(match_objects, async (msg) => { - await ch.publish("amq.topic", m.properties.headers.notify, + await ch.publish("amq.topic", msg.properties.headers.notify, new Buffer("analyze_pending")); await ch.sendToQueue(ANALYZE_QUEUE, new Buffer(msg.content.id), { @@ -295,7 +296,8 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { // log, move to error queue and NACK logger.error(err); await Promise.map(player_objects, async (m) => { - await ch.sendToQueue(QUEUE + "_failed", m.content, { + await ch.sendToQueue(QUEUE + "_failed", + new Buffer(JSON.stringify(m.content)), { persistent: true, type: m.properties.type, headers: m.properties.headers @@ -303,7 +305,8 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { await ch.nack(m, false, false); }); await Promise.map(match_objects, async (m) => { - await ch.sendToQueue(QUEUE + "_failed", m.content, { + await ch.sendToQueue(QUEUE + "_failed", + new Buffer(JSON.stringify(m.content)), { persistent: true, type: m.properties.type, headers: m.properties.headers |
