From f2c9557610ab2915be4537205f94a5252f552f1e Mon Sep 17 00:00:00 2001
From: ddamiron <ddamiron@sii.fr>
Date: Fri, 28 Jun 2019 15:14:33 +0200
Subject: [PATCH] update rabbit indentifiaction

---
 workers/doc-enricher.py       | 7 ++++++-
 workers/doc-indexer.py        | 7 ++++++-
 workers/doc-processor.py      | 7 ++++++-
 workers/metadata-processor.py | 8 +++++++-
 workers/reindexer.py          | 7 ++++++-
 workers/sample-generator.py   | 7 ++++++-
 6 files changed, 37 insertions(+), 6 deletions(-)

diff --git a/workers/doc-enricher.py b/workers/doc-enricher.py
index 92bc148..af44a6a 100644
--- a/workers/doc-enricher.py
+++ b/workers/doc-enricher.py
@@ -314,8 +314,13 @@ def enrich_docs( channel, method, properties, body ):
 def main(cfg):
 
     #from lib.close_connection import on_timeout
+    credentials = pika.PlainCredentials(username=cfg['rabbitmq']['user'], password=cfg['rabbitmq']['password'])
 
-    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
+    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'],
+                                                                   port=cfg['rabbitmq_port'],
+                                                                   credentials=credentials))
+
+    # connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
     #timeout = 5
     #connection.add_timeout(timeout, on_timeout(connection))
 
diff --git a/workers/doc-indexer.py b/workers/doc-indexer.py
index 4e302a7..1143c99 100644
--- a/workers/doc-indexer.py
+++ b/workers/doc-indexer.py
@@ -225,8 +225,13 @@ def index_docs(channel, method, properties, body):
 def main(cfg):
 
     #from lib.close_connection import on_timeout
+    credentials = pika.PlainCredentials(username=cfg['rabbitmq']['user'], password=cfg['rabbitmq']['password'])
 
-    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
+    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'],
+                                                                   port=cfg['rabbitmq_port'],
+                                                                   credentials=credentials))
+
+    # connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
     # timeout = 5
     # connection.add_timeout(timeout, on_timeout(connection))
 
diff --git a/workers/doc-processor.py b/workers/doc-processor.py
index ba1e9dd..43dedf7 100644
--- a/workers/doc-processor.py
+++ b/workers/doc-processor.py
@@ -184,8 +184,13 @@ def process_docs( channel, method, properties, body ):
 def main(cfg):
 
     #from lib.close_connection import on_timeout
+    credentials = pika.PlainCredentials(username=cfg['rabbitmq']['user'], password=cfg['rabbitmq']['password'])
 
-    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'],port=cfg['rabbitmq_port']))
+    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'],
+                                                                   port=cfg['rabbitmq_port'],
+                                                                   credentials=credentials))
+
+    # connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
     #timeout = 5
     #connection.add_timeout(timeout, on_timeout(connection))
 
diff --git a/workers/metadata-processor.py b/workers/metadata-processor.py
index 4528b23..c4c4fd0 100644
--- a/workers/metadata-processor.py
+++ b/workers/metadata-processor.py
@@ -518,7 +518,13 @@ def main(cfg):
 
     # logging.debug(cfg)
     #global connection
-    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
+    credentials = pika.PlainCredentials(username=cfg['rabbitmq']['user'], password=cfg['rabbitmq']['password'])
+
+    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'],
+                                                                   port=cfg['rabbitmq_port'],
+                                                                   credentials=credentials))
+
+    # connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
     # timeout = 5
     # connection.add_timeout(timeout, on_timeout(connection))
     channel = connection.channel()
diff --git a/workers/reindexer.py b/workers/reindexer.py
index 2bcfb66..f20ad02 100644
--- a/workers/reindexer.py
+++ b/workers/reindexer.py
@@ -239,8 +239,13 @@ def on_msg_callback(channel, method, properties, body):
 def main(cfg):
 
     #from lib.close_connection import on_timeout
+    credentials = pika.PlainCredentials(username=cfg['rabbitmq']['user'], password=cfg['rabbitmq']['password'])
 
-    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
+    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'],
+                                                                   port=cfg['rabbitmq_port'],
+                                                                   credentials=credentials))
+
+    # connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
     #timeout = 5
     #connection.add_timeout(timeout, on_timeout(connection))
     channel = connection.channel()
diff --git a/workers/sample-generator.py b/workers/sample-generator.py
index b13fd14..c949dd8 100644
--- a/workers/sample-generator.py
+++ b/workers/sample-generator.py
@@ -281,8 +281,13 @@ def main(cfg):
     # es_logger.setLevel(logging.INFO)
 
     #from lib.close_connection import on_timeout
+    credentials = pika.PlainCredentials(username=cfg['rabbitmq']['user'], password=cfg['rabbitmq']['password'])
 
-    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
+    connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'],
+                                                                   port=cfg['rabbitmq_port'],
+                                                                   credentials=credentials))
+
+    # connection = pika.BlockingConnection(pika.ConnectionParameters(host=cfg['rabbitmq_host'], port=cfg['rabbitmq_port']))
     #timeout = 5
     #connection.add_timeout(timeout, on_timeout(connection))
 
-- 
GitLab