mirror of
https://github.com/paulscherrerinstitute/sf_daq_broker.git
synced 2026-09-01 02:50:43 +02:00
Modify message broker client
This commit is contained in:
@@ -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="",
|
||||
|
||||
Reference in New Issue
Block a user