diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-07-26 21:15:29 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-07-26 21:15:29 +0200 |
| commit | 02a30adac0f92366ef9e6118581e39602b6d0d39 (patch) | |
| tree | c91fbd6d96c859b8259fc59037ffc3576c59f0e6 | |
| parent | 78a114a0702207d8bfd5412bf59dea67df63c006 (diff) | |
| download | processor-02a30adac0f92366ef9e6118581e39602b6d0d39.tar.gz processor-02a30adac0f92366ef9e6118581e39602b6d0d39.zip | |
catch deadlocks, fix failed msg passing
| -rw-r--r-- | worker.js | 6 |
1 files changed, 4 insertions, 2 deletions
@@ -225,6 +225,7 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { profiler = undefined; logger.info("processing batch", { + messages: msgs.size, players: player_data.size, matches: match_data.size }); @@ -282,7 +283,8 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { await ch.sendToQueue(ANALYZE_QUEUE, new Buffer(m.id), { persistent: true })); } catch (err) { - if (err instanceof Seq.TimeoutError) { + if (err instanceof Seq.TimeoutError || + (err instanceof Seq.DatabaseError && err.code == 1213)) { // deadlocks / timeout logger.error("SQL error", err); await Promise.map(msgs, async (m) => @@ -293,7 +295,7 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { await Promise.map(msgs, async (m) => { await ch.sendToQueue(QUEUE + "_failed", m.content, { persistent: true, - type: msg.properties.type, + type: m.properties.type, headers: m.properties.headers }); await ch.nack(m, false, false); |
