diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2017-08-23 19:15:12 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2017-08-23 19:15:12 +0200 |
| commit | d5f24344b7ae5d43d81cf9fecd7efee9515f1806 (patch) | |
| tree | bba9a2af27d8a2c71a3e47c4de4fb79c186ed176 | |
| parent | 8ac499fdbfdc1f1f784f499b77e1c2d17b9040ea (diff) | |
| download | processor-d5f24344b7ae5d43d81cf9fecd7efee9515f1806.tar.gz processor-d5f24344b7ae5d43d81cf9fecd7efee9515f1806.zip | |
send analyze_pending
| -rw-r--r-- | worker.js | 11 |
1 files changed, 8 insertions, 3 deletions
@@ -271,12 +271,17 @@ amqp.connect(RABBITMQ_URI).then(async (rabbit) => { if (match_objects.length > 0) await ch.publish("amq.topic", "global", new Buffer("matches_update")); // notify follow up services - if (DOANALYZEMATCH) // TODO + if (DOANALYZEMATCH) { await Promise.each(match_objects, async (msg) => { + await ch.publish("amq.topic", m.properties.headers.notify, + new Buffer("analyze_pending")); await ch.sendToQueue(ANALYZE_QUEUE, - new Buffer(msg.content.id), - { persistent: true }); + new Buffer(msg.content.id), { + persistent: true, + headers: { notify: msg.properties.headers.notify } + }); }); + } } catch (err) { if (err instanceof Seq.TimeoutError || (err instanceof Seq.DatabaseError && err.errno == 1213)) { |
