summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-08-23 19:15:12 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-08-23 19:15:12 +0200
commitd5f24344b7ae5d43d81cf9fecd7efee9515f1806 (patch)
treebba9a2af27d8a2c71a3e47c4de4fb79c186ed176
parent8ac499fdbfdc1f1f784f499b77e1c2d17b9040ea (diff)
downloadprocessor-d5f24344b7ae5d43d81cf9fecd7efee9515f1806.tar.gz
processor-d5f24344b7ae5d43d81cf9fecd7efee9515f1806.zip
send analyze_pending
-rw-r--r--worker.js11
1 files changed, 8 insertions, 3 deletions
diff --git a/worker.js b/worker.js
index b95a12a..e585cc7 100644
--- a/worker.js
+++ b/worker.js
@@ -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)) {