From 8675e05c640a9cbfdfec0de0ae63b27cdf174d50 Mon Sep 17 00:00:00 2001 From: Andrej Babic Date: Tue, 17 Nov 2020 13:52:06 +0100 Subject: [PATCH] Refactor image assembly to expose socket creation --- sf_daq_broker/detector/image_assembler.py | 17 +++++++++++------ 1 file changed, 11 insertions(+), 6 deletions(-) diff --git a/sf_daq_broker/detector/image_assembler.py b/sf_daq_broker/detector/image_assembler.py index bb94eef..0a25423 100644 --- a/sf_daq_broker/detector/image_assembler.py +++ b/sf_daq_broker/detector/image_assembler.py @@ -23,12 +23,7 @@ class ImageAssembler(object): raise ValueError("RamBuffer must have at least 5 slots.") for module_id in range(self.n_modules): - # We use the PUSH/PULL mechanism to moderate disk throughput. - receiver = self.zmq_context.socket(zmq.PULL) - # No buffering on send side - receiver dictates reading speed. - receiver.setsockopt(zmq.RCVHWM, zmq_rcv_hwm) - # If in 1 second nothing was received we have a problem in reading. - receiver.setsockopt(zmq.RCVTIMEO, 1000) + receiver = get_pull_receiver(self.zmq_context, zmq_rcv_hwm=zmq_rcv_hwm) receiver.connect("inproc://%s" % module_id) @@ -97,3 +92,13 @@ class ImageAssembler(object): self.receivers.clear() + +def get_pull_receiver(context, zmq_rcv_hwm): + # We use the PUSH/PULL mechanism to moderate disk throughput. + receiver = context.socket(zmq.PULL) + # No buffering on send side - receiver dictates reading speed. + receiver.setsockopt(zmq.RCVHWM, zmq_rcv_hwm) + # If in 1 second nothing was received we have a problem in reading. + receiver.setsockopt(zmq.RCVTIMEO, 1000) + + return receiver