diff --git a/bec_server/bec_server/device_server/devices/devicemanager.py b/bec_server/bec_server/device_server/devices/devicemanager.py index 3a0b4c003..061fe0dc0 100644 --- a/bec_server/bec_server/device_server/devices/devicemanager.py +++ b/bec_server/bec_server/device_server/devices/devicemanager.py @@ -13,13 +13,11 @@ from collections import deque from typing import TYPE_CHECKING, Callable -import numpy as np import ophyd import ophyd_devices as opd from ophyd.ophydobj import OphydObject from ophyd.signal import EpicsSignalBase from ophyd_devices.utils.bec_signals import BECMessageSignal -from typeguard import typechecked from bec_lib import messages, plugin_helper from bec_lib.alarm_handler import Alarms @@ -545,7 +543,6 @@ def initialize_device(self, dev: dict, config: dict, obj: OphydObject) -> DSDevi # Add subscriptions to device events and signal if supported by the device if hasattr(obj, "event_types"): self._subscribe_to_device_events(obj, opaas_obj) - self._subscribe_to_bec_device_events(obj) self._subscribe_to_auto_monitors(obj) self._subscribe_to_limit_updates(obj) self._subscribe_to_bec_signals(obj) @@ -577,35 +574,6 @@ def _subscribe_to_device_events(self, obj: OphydObject, opaas_obj: DSDevice): if hasattr(obj, "motor_is_moving"): obj.motor_is_moving.subscribe(self._obj_callback_is_moving, run=opaas_obj.enabled) # type: ignore - def _subscribe_to_bec_device_events(self, obj: OphydObject): - """ - Subscribe to BEC device events, such as device_monitor_2d, device_monitor_1d, - file_event, done_moving, flyer, and progress. - - These events are deprecated and will be removed in the future. Use the - _subscribe_to_bec_signals method instead. - - Args: - obj (OphydObject): Ophyd object to subscribe to BEC device events - - """ - if "device_monitor_2d" in obj.event_types: - obj.subscribe( - self._obj_callback_device_monitor_2d, event_type="device_monitor_2d", run=False - ) - if "device_monitor_1d" in obj.event_types: - obj.subscribe( - self._obj_callback_device_monitor_1d, event_type="device_monitor_1d", run=False - ) - if "file_event" in obj.event_types: - obj.subscribe(self._obj_callback_file_event, event_type="file_event", run=False) - if "done_moving" in obj.event_types: - obj.subscribe(self._obj_callback_done_moving, event_type="done_moving", run=False) - if "flyer" in obj.event_types: - obj.subscribe(self._obj_flyer_callback, event_type="flyer", run=False) - if "progress" in obj.event_types: - obj.subscribe(self._obj_callback_progress, event_type="progress", run=False) - def _subscribe_to_auto_monitors(self, obj: OphydObject): """ If the component has set the _auto_monitor attribute to True, @@ -783,98 +751,6 @@ def _obj_callback_configuration(self, *_args, obj: OphydObject, **kwargs): ) pipe.execute() - @typechecked - def _obj_callback_device_monitor_2d( - self, *_args, obj: OphydObject, value: np.ndarray, timestamp: float | None = None, **kwargs - ): - """ - DEPRECATED: Use _obj_callback_preview instead. - - Callback for ophyd monitor events. Sends the data to redis. - Introduces a check of the data size, and incorporates a limit which is defined in max_size (in MB) - - Args: - obj (OphydObject): ophyd object - value (np.ndarray): data from ophyd device - - """ - # Convert sizes from bytes to MB - dsize = len(value.tobytes()) / 1e6 - max_size = 1000 - if dsize > max_size: - logger.warning( - f"Data size of single message is too large to send, current max_size {max_size}." - ) - return - if obj.connected: - name = obj.root.name - metadata = self.devices[name].metadata - msg = messages.DeviceMonitor2DMessage( - device=name, - data=value, - metadata=metadata, - timestamp=timestamp if timestamp else time.time(), - ) - stream_msg = {"data": msg} - self.connector.xadd( - MessageEndpoints.device_monitor_2d(name), - stream_msg, - max_size=min(100, int(max_size // dsize)), - expire=3600, - ) - - def _obj_callback_device_monitor_1d( - self, *_args, obj: OphydObject, value: np.ndarray, timestamp: float | None = None, **kwargs - ): - """ - DEPRECATED: Use _obj_callback_preview instead. - - Callback for ophyd monitor events. Sends the data to redis. - Introduces a check of the data size, and incorporates a limit which is defined in max_size (in MB) - - Args: - obj (OphydObject): ophyd object - value (np.ndarray): data from ophyd device - - """ - # Convert sizes from bytes to MB - dsize = len(value.tobytes()) / 1e6 - max_size = 1000 - if dsize > max_size: - logger.warning( - f"Data size of single message is too large to send, current max_size {max_size}." - ) - return - if obj.connected: - name = obj.root.name - metadata = self.devices[name].metadata - msg = messages.DeviceMonitor1DMessage( - device=name, - data=value, - metadata=metadata, - timestamp=timestamp if timestamp else time.time(), - ) - stream_msg = {"data": msg} - self.connector.xadd( - MessageEndpoints.device_monitor_1d(name), - stream_msg, - max_size=min(100, int(max_size // dsize)), - expire=3600, - ) - - def _obj_callback_acq_done(self, *_args, **kwargs): - device = kwargs["obj"].root.name - status = 0 - metadata = self.devices[device].metadata - self.connector.set( - MessageEndpoints.device_status(device), - messages.DeviceStatusMessage(device=device, status=status, metadata=metadata), - ) - - def _obj_callback_done_moving(self, *args, **kwargs): - self._obj_callback_readback(*args, **kwargs) - # self._obj_callback_acq_done(*args, **kwargs) - def _obj_callback_is_moving(self, *_args, **kwargs): device = kwargs["obj"].root.name status = int(kwargs.get("value")) @@ -884,110 +760,6 @@ def _obj_callback_is_moving(self, *_args, **kwargs): messages.DeviceStatusMessage(device=device, status=status, metadata=metadata), ) - def _obj_flyer_callback(self, *_args, **kwargs): - obj = kwargs["obj"] - logger.warning( - f"Flyer callback will be deprecated in future, please refactor your device {obj.root.name} in favor of an async devices as soon as possible." - ) - data = kwargs["value"].get("data") - ds_obj = self.devices[obj.root.name] - metadata = ds_obj.metadata - if "scan_id" not in metadata: - return - - if not hasattr(ds_obj, "emitted_points"): - ds_obj.emitted_points = {} - - emitted_points = ds_obj.emitted_points.get(metadata["scan_id"], 0) - - # make sure all arrays are of equal length - max_points = min(len(d) for d in data.values()) - - pipe = self.connector.pipeline() - for ii in range(emitted_points, max_points): - timestamp = time.time() - signals = {} - for key, val in data.items(): - signals[key] = {"value": val[ii], "timestamp": timestamp} - msg = messages.DeviceMessage(signals=signals, metadata={"point_id": ii, **metadata}) - self.connector.set_and_publish( - MessageEndpoints.device_read(obj.root.name), msg, pipe=pipe - ) - - ds_obj.emitted_points[metadata["scan_id"]] = max_points - msg = messages.DeviceStatusMessage( - device=obj.root.name, status=max_points, metadata=metadata - ) - self.connector.set(MessageEndpoints.device_status(obj.root.name), msg, pipe=pipe) - pipe.execute() - - def _obj_callback_progress(self, *_args, obj, value, max_value, done, **kwargs): - """ - DEPRECATED: Use _obj_callback_progress_signal instead. - - Callback for progress events. Sends the data to redis. - """ - metadata = self.devices[obj.root.name].metadata - msg = messages.ProgressMessage( - value=value, max_value=max_value, done=done, metadata=metadata - ) - self.connector.set_and_publish( - MessageEndpoints.device_progress(obj.root.name), msg, expire=3600 - ) - - def _obj_callback_file_event( - self, - *_args, - obj, - file_path: str, - done: bool, - successful: bool, - file_type: str = "h5", - hinted_h5_entries: dict[str, str] | None = None, - **kwargs, - ): - """ - DEPRECATED: Use _obj_callback_file_event_signal instead. - - Callback for file events on devices. This callback set and publishes - a file message to the file_event and public_file endpoints in Redis to inform - the file writer and other services about externally created files. - - Args: - obj (OphydObject): ophyd object - file_path (str): file path to the created file - done (bool): if the file is done - successful (bool): if the file was created successfully - file_type (str): Optional, file type. Default is h5. - hinted_h5_entry (dict[str, str] | None): Optional, hinted h5 entry. Please check FileMessage for more details - """ - device_name = obj.root.name - metadata = self.devices[device_name].metadata - if kwargs.get("metadata") is not None: - metadata.update(kwargs.get("metadata")) - scan_id = metadata.get("scan_id") - msg = messages.FileMessage( - file_path=file_path, - done=done, - successful=successful, - file_type=file_type, - device_name=device_name, - is_master_file=False, - hinted_h5_entries=hinted_h5_entries, - metadata=metadata, - ) - pipe = self.connector.pipeline() - self.connector.set_and_publish( - MessageEndpoints.file_event(device_name), msg, pipe=pipe, expire=3600 - ) - self.connector.set_and_publish( - MessageEndpoints.public_file(scan_id=scan_id, name=device_name), - msg, - pipe=pipe, - expire=3600, - ) - pipe.execute() - def _obj_callback_bec_message_signal( self, *_args, obj: OphydObject, value: messages.BECMessage, **kwargs ): diff --git a/bec_server/tests/tests_device_server/test_device_manager_ds.py b/bec_server/tests/tests_device_server/test_device_manager_ds.py index 721985613..d4b206713 100644 --- a/bec_server/tests/tests_device_server/test_device_manager_ds.py +++ b/bec_server/tests/tests_device_server/test_device_manager_ds.py @@ -140,75 +140,6 @@ def mocked_failed_connection(obj, **kwargs): ) -@pytest.mark.parametrize("device_manager_class", [DeviceManagerDS]) -def test_flyer_event_callback(dm_with_devices, connected_connector): - device_manager = dm_with_devices - samx = device_manager.devices.samx - samx.metadata = {"scan_id": "12345"} - # Use here fake redis connector to avoid complications with PipelineMock - device_manager.connector = connected_connector - device_manager._obj_flyer_callback( - obj=samx.obj, - value={"data": {"idata": np.random.rand(20), "edata": np.random.rand(20)}}, - metadata={"scan_id": "test_scan_id"}, - ) - msg = connected_connector.get(MessageEndpoints.device_read("samx")) - assert "signals" in msg.content - assert "idata" in msg.content["signals"] - assert "edata" in msg.content["signals"] - msg = connected_connector.get(MessageEndpoints.device_status("samx")) - assert msg.metadata["scan_id"] == "12345" - assert msg.content["device"] == "samx" - assert msg.content["status"] == 20 - - -@pytest.mark.parametrize("device_manager_class", [DeviceManagerDS]) -def test_obj_callback_progress(dm_with_devices): - device_manager = dm_with_devices - samx = device_manager.devices.samx - samx.metadata = {"scan_id": "12345"} - - with mock.patch.object(device_manager, "connector") as mock_connector: - device_manager._obj_callback_progress(obj=samx.obj, value=1, max_value=2, done=False) - mock_connector.set_and_publish.assert_called_once_with( - MessageEndpoints.device_progress("samx"), - messages.ProgressMessage( - value=1, max_value=2, done=False, metadata={"scan_id": "12345"} - ), - expire=3600, - ) - - -@pytest.mark.parametrize( - "value", [np.empty(shape=(10, 10)), np.empty(shape=(100, 100)), np.empty(shape=(1000, 1000))] -) -@pytest.mark.parametrize("device_manager_class", [DeviceManagerDS]) -def test_obj_device_monitor_2d_callback(dm_with_devices, value): - device_manager = dm_with_devices - eiger = device_manager.devices.eiger - eiger.metadata = {"scan_id": "12345"} - value_size = len(value.tobytes()) / 1e6 # MB - max_size = 1000 - timestamp = time.time() - with mock.patch.object(device_manager, "connector") as mock_connector: - device_manager._obj_callback_device_monitor_2d( - obj=eiger.obj, value=value, timestamp=timestamp - ) - stream_msg = { - "data": messages.DeviceMonitor2DMessage( - device=eiger.name, data=value, metadata={"scan_id": "12345"}, timestamp=timestamp - ) - } - - assert mock_connector.xadd.call_count == 1 - assert mock_connector.xadd.call_args == mock.call( - MessageEndpoints.device_monitor_2d(eiger.name), - stream_msg, - max_size=min(100, int(max_size // value_size)), - expire=3600, - ) - - @pytest.mark.parametrize("device_manager_class", [DeviceManagerDS]) def test_device_manager_ds_reset_config(dm_with_devices): with mock.patch.object(dm_with_devices, "connector") as mock_connector: @@ -224,78 +155,19 @@ def test_device_manager_ds_reset_config(dm_with_devices): ) -@pytest.mark.parametrize("device_manager_class", [DeviceManagerDS]) -def test_obj_callback_file_event(dm_with_devices, connected_connector): - device_manager = dm_with_devices - eiger = device_manager.devices.eiger - eiger.metadata = {"scan_id": "12345"} - # Use here fake redis connector, pipe is used and checks pydantic models - device_manager.connector = connected_connector - device_manager._obj_callback_file_event( - obj=eiger.obj, - file_path="test_file_path", - done=True, - successful=True, - hinted_h5_entries={"my_entry": "entry/data/data"}, - metadata={"user_info": "my_info"}, - ) - msg = connected_connector.get(MessageEndpoints.file_event(name="eiger")) - msg2 = connected_connector.get(MessageEndpoints.public_file(scan_id="12345", name="eiger")) - assert msg == msg2 - assert msg.content["file_path"] == "test_file_path" - assert msg.content["done"] is True - assert msg.content["successful"] is True - assert msg.content["hinted_h5_entries"] == {"my_entry": "entry/data/data"} - assert msg.content["file_type"] == "h5" - assert msg.metadata == {"scan_id": "12345", "user_info": "my_info"} - assert msg.content["is_master_file"] is False - - @pytest.mark.parametrize("device_manager_class", [DeviceManagerDS]) def test_subscribe_to_device_events(dm_with_devices): opaas_obj = mock.MagicMock() opaas_obj.enabled = False obj = mock.MagicMock() - # Test 2 event types together - obj.event_types = ("file_event", "device_monitor_1d") - with mock.patch.object(dm_with_devices, "_obj_callback_file_event") as mock_callback_file_event: - with mock.patch.object( - dm_with_devices, "_obj_callback_device_monitor_1d" - ) as mock_callback_device_monitor_1d: - dm_with_devices._subscribe_to_device_events(obj=obj, opaas_obj=opaas_obj) - assert obj.subscribe.call_count == 0 - dm_with_devices._subscribe_to_bec_device_events(obj=obj) - assert obj.subscribe.call_count == 2 - assert ( - mock.call(mock_callback_file_event, event_type="file_event", run=False) - in obj.subscribe.call_args_list - ) - assert ( - mock.call( - mock_callback_device_monitor_1d, event_type="device_monitor_1d", run=False - ) - in obj.subscribe.call_args_list - ) - - # Test all event types - for ii, event_type in enumerate( - [ - "readback", - "value", - "device_monitor_1d", - "device_monitor_2d", - "file_event", - "done_moving", - "progress", - ] - ): + for event_type in ["readback", "value"]: + obj.reset_mock() obj.event_types = (event_type,) callback_name = ( f"_obj_callback_{event_type}" if event_type != "value" else "_obj_callback_readback" ) with mock.patch.object(dm_with_devices, callback_name) as mock_callback: dm_with_devices._subscribe_to_device_events(obj=obj, opaas_obj=opaas_obj) - dm_with_devices._subscribe_to_bec_device_events(obj=obj) assert obj.subscribe.call_args == mock.call( mock_callback, event_type=event_type, run=False )