From f2f5f362e65c4d4dcc81e576151e3734ec83e9ce Mon Sep 17 00:00:00 2001 From: Andrej Babic Date: Thu, 20 Aug 2020 14:47:39 +0200 Subject: [PATCH] Modify message broker client --- sf_daq_broker/rabbitmq/msg_broker_client.py | 43 +++++---------------- 1 file changed, 10 insertions(+), 33 deletions(-) diff --git a/sf_daq_broker/rabbitmq/msg_broker_client.py b/sf_daq_broker/rabbitmq/msg_broker_client.py index e74f75c..88e4b24 100644 --- a/sf_daq_broker/rabbitmq/msg_broker_client.py +++ b/sf_daq_broker/rabbitmq/msg_broker_client.py @@ -1,55 +1,32 @@ import json +import sf_daq_broker.rabbitmq.config as broker_config + from pika import BlockingConnection, ConnectionParameters, BasicProperties class RabbitMqClient(object): - REQUEST_EXCHANGE = "request" - STATUS_EXCHANGE = "status" - DEFAULT_BROKER_URL = "127.0.0.1" - - def __init__(self, broker_url=DEFAULT_BROKER_URL): + def __init__(self, broker_url=broker_config.DEFAULT_BROKER_URL): self.connection = BlockingConnection(ConnectionParameters(broker_url)) self.channel = self.connection.channel() - self.channel.exchange_declare(exchange=self.REQUEST_EXCHANGE, + self.channel.exchange_declare(exchange=broker_config.REQUEST_EXCHANGE, exchange_type="topic") - self.channel.exchange_declare(exchange=self.STATUS_EXCHANGE, + self.channel.exchange_declare(exchange=broker_config.STATUS_EXCHANGE, exchange_type="fanout") def close(self): self.connection.close() - def request_write(self, - output_prefix, - metadata=None, - detectors=None, - bsread_channels=None, - epics_pvs=None): + def send(self, tag, write_request): - routing_key = "." + routing_key = "." + tag + "." - if detectors: - for detector in detectors: - routing_key += detector + "." + body_bytes = json.dumps(write_request).encode() - if bsread_channels: - routing_key += "bsread" + "." - - if epics_pvs: - routing_key += "epics" + "." - - body_bytes = json.dumps({ - "output_prefix": output_prefix, - "metadata": metadata, - "detectors": detectors, - "bsread_channels": bsread_channels, - "epics_pvs": epics_pvs - }).encode() - - self.channel.basic_publish(exchange=self.REQUEST_EXCHANGE, + self.channel.basic_publish(exchange=broker_config.REQUEST_EXCHANGE, routing_key=routing_key, body=body_bytes) @@ -59,7 +36,7 @@ class RabbitMqClient(object): "routing_key": routing_key } - self.channel.basic_publish(exchange=self.STATUS_EXCHANGE, + self.channel.basic_publish(exchange=broker_config.STATUS_EXCHANGE, properties=BasicProperties( headers=status_header), routing_key="",