|
8 | 8 | from bec_lib import messages |
9 | 9 | from bec_lib.alarm_handler import AlarmBase |
10 | 10 | from bec_lib.bl_states import DeviceWithinLimitsStateConfig |
| 11 | +from bec_lib.data_api import DataAPI |
11 | 12 | from bec_lib.devicemanager import DeviceConfigError |
12 | 13 | from bec_lib.endpoints import MessageEndpoints |
13 | 14 | from bec_lib.logger import bec_logger |
@@ -116,6 +117,78 @@ def dummy_callback(data, metadata): |
116 | 117 | assert msg.content == reference_container["data"][ii] |
117 | 118 |
|
118 | 119 |
|
| 120 | +@pytest.mark.timeout(100) |
| 121 | +def test_data_api_bundles_monitored_grid_scan_with_monitored_async_signal_lib(bec_client_lib): |
| 122 | + bec = bec_client_lib |
| 123 | + scans = bec.scans |
| 124 | + dev = bec.device_manager.devices |
| 125 | + callbacks = [] |
| 126 | + |
| 127 | + def callback(data, metadata): |
| 128 | + callbacks.append((data, metadata)) |
| 129 | + |
| 130 | + DataAPI.clear_instance() |
| 131 | + try: |
| 132 | + bec.metadata.update({"unit_test": "test_data_api_grid_scan_monitored_async_bundle"}) |
| 133 | + scans.umv(dev.samx, 0, dev.samy, 0, relative=False) |
| 134 | + dev.waveform.sim.select_model("ConstantModel") |
| 135 | + dev.waveform.async_update.set("add") |
| 136 | + |
| 137 | + data_api = DataAPI(bec) |
| 138 | + with data_api.create_subscription() as subscription: |
| 139 | + subscription.add_device("samx", "samx") |
| 140 | + subscription.add_device("samy", "samy") |
| 141 | + subscription.add_device("waveform", "waveform_waveform_0d") |
| 142 | + subscription.set_callback(callback) |
| 143 | + |
| 144 | + status = scans.grid_scan( |
| 145 | + dev.samx, -5, 5, 10, dev.samy, -5, 5, 10, exp_time=0.05, relative=True |
| 146 | + ) |
| 147 | + scan_id = status.queue_item.scan_ids[0] |
| 148 | + |
| 149 | + deadline = time.time() + 10 |
| 150 | + readout_priority = None |
| 151 | + while time.time() < deadline: |
| 152 | + scan_item = bec.queue.scan_storage.find_scan_by_ID(scan_id) |
| 153 | + status_message = getattr(scan_item, "status_message", None) |
| 154 | + readout_priority = getattr(status_message, "readout_priority", None) |
| 155 | + if readout_priority and "samx" in readout_priority.get("monitored", []): |
| 156 | + break |
| 157 | + time.sleep(0.05) |
| 158 | + |
| 159 | + assert readout_priority is not None |
| 160 | + assert "samx" in readout_priority.get("monitored", []) |
| 161 | + assert "samy" in readout_priority.get("monitored", []) |
| 162 | + |
| 163 | + status.wait(num_points=True, file_written=True) |
| 164 | + expected_points = status.scan.num_points |
| 165 | + |
| 166 | + deadline = time.time() + 15 |
| 167 | + while time.time() < deadline and len(callbacks) < expected_points: |
| 168 | + time.sleep(0.1) |
| 169 | + |
| 170 | + assert expected_points == 100 |
| 171 | + assert len(callbacks) == expected_points |
| 172 | + |
| 173 | + for idx, (data, metadata) in enumerate(callbacks, start=1): |
| 174 | + assert set(data) == {"samx", "samy", "waveform"} |
| 175 | + assert "samx" in data["samx"] |
| 176 | + assert "samy" in data["samy"] |
| 177 | + assert "waveform_waveform_0d" in data["waveform"] |
| 178 | + assert metadata["scan_id"] == scan_id |
| 179 | + assert metadata["async_update"]["type"] == "add" |
| 180 | + waveform_value = data["waveform"]["waveform_waveform_0d"]["value"] |
| 181 | + if idx == 1: |
| 182 | + assert isinstance(waveform_value, (int, np.integer)) |
| 183 | + else: |
| 184 | + assert isinstance(waveform_value, list) |
| 185 | + assert len(waveform_value) == idx |
| 186 | + assert not isinstance(data["samx"]["samx"]["value"], list) |
| 187 | + assert not isinstance(data["samy"]["samy"]["value"], list) |
| 188 | + finally: |
| 189 | + DataAPI.clear_instance() |
| 190 | + |
| 191 | + |
119 | 192 | @pytest.mark.timeout(100) |
120 | 193 | def test_rpc_call_in_event_callback(bec_client_lib): |
121 | 194 | scans = bec_client_lib.scans |
|
0 commit comments