diff --git a/pxiii_bec/device_configs/x06da_device_config.yaml b/pxiii_bec/device_configs/x06da_device_config.yaml index cb1464e..7d79bde 100644 --- a/pxiii_bec/device_configs/x06da_device_config.yaml +++ b/pxiii_bec/device_configs/x06da_device_config.yaml @@ -394,8 +394,8 @@ xbox_diode: readoutPriority: monitored readOnly: true softwareTrigger: false -samdist: - description: Sample distance +gonpos: + description: Sample sensor distance deviceClass: ophyd.EpicsSignalRO deviceConfig: {read_pv: 'X06DA-ES-DF1:CBOX-USER1', auto_monitor: true} onFailure: buffer @@ -403,7 +403,7 @@ samdist: readoutPriority: monitored readOnly: true softwareTrigger: false -samrange: +gonvalid: description: Sample in valid distance deviceClass: ophyd.EpicsSignalRO deviceConfig: {read_pv: 'X06DA-ES-DF1:CBOX-CMP1', auto_monitor: true} @@ -430,6 +430,17 @@ samcam: readoutPriority: monitored readOnly: false softwareTrigger: false +samimg: + description: Sample camera ZMQ stream + deviceClass: pxiii_bec.devices.StdDaqPreviewDetector + deviceConfig: + url: 'tcp://129.129.110.12:9089' + deviceTags: + - detector + enabled: true + readoutPriority: async + readOnly: false + softwareTrigger: false bstop_pneum: @@ -544,6 +555,15 @@ gmy: readoutPriority: monitored readOnly: false softwareTrigger: false +gmz: + description: ABR axial stage + deviceClass: pxiii_bec.devices.A3200Axis + deviceConfig: {prefix: 'X06DA-ES-DF1:GMZ', base_pv: 'X06DA-ES'} + onFailure: buffer + enabled: true + readoutPriority: monitored + readOnly: false + softwareTrigger: false omega: description: ABR rotation stage deviceClass: pxiii_bec.devices.A3200Axis @@ -611,13 +631,13 @@ phi: -samimg: - description: Sample camera image - deviceClass: ophyd_devices.devices.areadetector.plugins.ImagePlugin_V35 - deviceConfig: {prefix: 'X06DA-SAMCAM:image1:'} - onFailure: buffer - enabled: false - readoutPriority: monitored - readOnly: true - softwareTrigger: false +# samimgs: +# description: Sample camera image +# deviceClass: ophyd_devices.devices.areadetector.plugins.ImagePlugin_V35 +# deviceConfig: {prefix: 'X06DA-SAMCAM:image1:', foo: 'bar'} +# onFailure: buffer +# enabled: false +# readoutPriority: monitored +# readOnly: true +# softwareTrigger: false diff --git a/pxiii_bec/devices/StdDaqPreview.py b/pxiii_bec/devices/StdDaqPreview.py new file mode 100644 index 0000000..2e27801 --- /dev/null +++ b/pxiii_bec/devices/StdDaqPreview.py @@ -0,0 +1,194 @@ +# -*- coding: utf-8 -*- +""" +Standard DAQ preview image stream module + +Created on Thu Jun 27 17:28:43 2024 + +@author: mohacsi_i +""" +import json +import enum +from time import sleep, time +from threading import Thread +import zmq +import numpy as np +from ophyd import Device, Signal, Component, Kind, DeviceStatus +from ophyd_devices.interfaces.base_classes.psi_detector_base import ( + CustomDetectorMixin, + PSIDetectorBase, +) + +from bec_lib import bec_logger +logger = bec_logger.logger +ZMQ_TOPIC_FILTER = b'' + + +class StdDaqPreviewState(enum.IntEnum): + """Standard DAQ ophyd device states""" + UNKNOWN = 0 + DETACHED = 1 + MONITORING = 2 + + +class StdDaqPreviewMixin(CustomDetectorMixin): + """Setup class for the standard DAQ preview stream + + Parent class: CustomDetectorMixin + """ + _mon = None + + def on_stage(self): + """Start listening for preview data stream""" + if self._mon is not None: + self.parent.unstage() + sleep(0.5) + + self.parent.connect() + self._stop_polling = False + self._mon = Thread(target=self.poll, daemon=True) + self._mon.start() + + def on_unstage(self): + """Stop a running preview""" + if self._mon is not None: + self._stop_polling = True + # Might hang on recv_multipart + self._mon.join(timeout=1) + # So also disconnect the socket + self.parent._socket.disconnect(self.parent.url.get()) + + def on_stop(self): + """Stop a running preview""" + self.on_unstage() + + def poll(self): + """Collect streamed updates""" + self.parent.status.set(StdDaqPreviewState.MONITORING, force=True) + try: + t_last = time() + while True: + try: + # Exit loop and finish monitoring + if self._stop_polling: + logger.info(f"[{self.parent.name}]\tDetaching monitor") + break + + # pylint: disable=no-member + r = self.parent._socket.recv_multipart(flags=zmq.NOBLOCK) + + # Length and throtling checks + if len(r) != 2: + logger.warning( + f"[{self.parent.name}] Received malformed array of length {len(r)}") + t_curr = time() + t_elapsed = t_curr - t_last + if t_elapsed < self.parent.throttle.get(): + sleep(0.1) + continue + + # Unpack the Array V1 reply to metadata and array data + meta, data = r + + # Update image and update subscribers + header = json.loads(meta) + if header["type"] == "uint16": + image = np.frombuffer(data, dtype=np.uint16) + if header["type"] == "uint8": + image = np.frombuffer(data, dtype=np.uint8) + if image.size != np.prod(header['shape']): + err = f"Unexpected array size of {image.size} for header: {header}" + raise ValueError(err) + image = image.reshape(header['shape']) + + # Update image and update subscribers + self.parent.frameno.put(header['frame'], force=True) + self.parent.image_shape.put(header['shape'], force=True) + self.parent.image.put(image, force=True) + self.parent._last_image = image + self.parent._run_subs(sub_type=self.parent.SUB_MONITOR, value=image) + t_last = t_curr + logger.debug( + f"[{self.parent.name}] Updated frame {header['frame']}\t" + f"Shape: {header['shape']}\tMean: {np.mean(image):.3f}" + ) + except ValueError: + # Happens when ZMQ partially delivers the multipart message + pass + except zmq.error.Again: + # Happens when receive queue is empty + sleep(0.1) + except Exception as ex: + logger.info(f"[{self.parent.name}]\t{str(ex)}") + raise + finally: + self._mon = None + self.parent.status.set(StdDaqPreviewState.DETACHED, force=True) + logger.info(f"[{self.parent.name}]\tDetaching monitor") + + +class StdDaqPreviewDetector(PSIDetectorBase): + """Detector wrapper class around the StdDaq preview image stream. + + This was meant to provide live image stream directly from the StdDAQ + but also works with other ARRAY v1 streamers, like the AreaDetector + ZMQ plugin. + Note that the preview stream must be already throtled in order to cope + with the incoming data and the python class might throttle it further. + + You can add a preview widget to the dock by: + cam_widget = gui.add_dock('cam_dock1').add_widget('BECFigure').image('daq_stream1') + """ + # Subscriptions for plotting image + USER_ACCESS = ["get_image"] + SUB_MONITOR = "device_monitor_2d" + _default_sub = SUB_MONITOR + + custom_prepare_cls = StdDaqPreviewMixin + + # Status attributes + url = Component(Signal, kind=Kind.config, metadata={"write_access": False}) + throttle = Component(Signal, value=0.25, kind=Kind.config) + status = Component(Signal, value=StdDaqPreviewState.UNKNOWN, kind=Kind.omitted, metadata={"write_access": False}) + frameno = Component(Signal, kind=Kind.hinted, metadata={"write_access": False}) + image_shape = Component(Signal, kind=Kind.normal, metadata={"write_access": False}) + # FIXME: The BEC client caches the read()s from the last 50 scans + image = Component(Signal, kind=Kind.omitted, metadata={"write_access": False}) + _last_image = None + + def __init__( + self, *args, url: str = "tcp://129.129.95.38:20000", parent: Device = None, **kwargs + ) -> None: + super().__init__(*args, parent=parent, **kwargs) + self.url.set(url, force=True).wait() + # Connect to the DAQ + self.connect() + + def connect(self): + """Connect to te StDAQs PUB-SUB streaming interface + + StdDAQ may reject connection for a few seconds when it restarts, + so if it fails, wait a bit and try to connect again. + """ + # pylint: disable=no-member + # Socket to talk to server + context = zmq.Context() + self._socket = context.socket(zmq.SUB) + self._socket.setsockopt(zmq.SUBSCRIBE, ZMQ_TOPIC_FILTER) + try: + self._socket.connect(self.url.get()) + except ConnectionRefusedError: + sleep(1) + self._socket.connect(self.url.get()) + + def get_image(self): + """ + Gets the last image as an attribute in case image must be abandoned + due to some caching on the BEC. + """ + return self._last_image + + +# Automatically connect to MicroSAXS testbench if directly invoked +if __name__ == "__main__": + daq = StdDaqPreviewDetector(url="tcp://129.129.95.111:20000", name="preview") + daq.wait_for_connection() diff --git a/pxiii_bec/devices/__init__.py b/pxiii_bec/devices/__init__.py index 05082c8..03d87ca 100644 --- a/pxiii_bec/devices/__init__.py +++ b/pxiii_bec/devices/__init__.py @@ -7,3 +7,4 @@ Ophyd devices for the PX III beamline, including the MX specific Aerotech A3200 from .A3200 import AerotechAbrStage from .A3200utils import A3200Axis from .SmarGon import SmarGonAxis +from .StdDaqPreview import StdDaqPreviewDetector \ No newline at end of file