Change how detector metadate is stored - now in redis - so status updates are faster.
This commit is contained in:
@@ -1101,6 +1101,68 @@ class BeamlineConfig:
|
||||
"smargon": self.simulate_smargon,
|
||||
}
|
||||
|
||||
@staticmethod
|
||||
def _coerce_optional_float(value) -> float | None:
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
if isinstance(value, np.ndarray):
|
||||
if value.size == 0:
|
||||
return None
|
||||
value = value.flatten().tolist()
|
||||
|
||||
if isinstance(value, (list, tuple)):
|
||||
if len(value) == 0:
|
||||
return None
|
||||
value = value[0]
|
||||
|
||||
try:
|
||||
return float(value)
|
||||
except (TypeError, ValueError):
|
||||
logger.warning(f"Failed to coerce cached numeric value to float: {value!r}")
|
||||
return None
|
||||
|
||||
def _detector_metadata_key(self) -> str:
|
||||
return f"{self.__bl}:detector_metadata"
|
||||
|
||||
def get_detector_metadata(self) -> dict:
|
||||
try:
|
||||
raw = self.__client.get(self._detector_metadata_key())
|
||||
if raw in (None, "", b""):
|
||||
return {}
|
||||
|
||||
if isinstance(raw, bytes):
|
||||
raw = raw.decode("utf-8")
|
||||
|
||||
payload = json.loads(str(raw))
|
||||
return payload if isinstance(payload, dict) else {}
|
||||
except Exception as e:
|
||||
logger.warning(f"Failed to read detector metadata from Redis: {e}")
|
||||
return {}
|
||||
|
||||
def set_detector_metadata(self, payload: dict) -> dict:
|
||||
safe_payload = dict(payload or {})
|
||||
safe_payload["dtz_low"] = self._coerce_optional_float(safe_payload.get("dtz_low"))
|
||||
safe_payload["dtz_high"] = self._coerce_optional_float(safe_payload.get("dtz_high"))
|
||||
safe_payload["pixel_size_mm"] = self._coerce_optional_float(safe_payload.get("pixel_size_mm"))
|
||||
safe_payload["updated_at"] = datetime.now().isoformat(timespec="seconds")
|
||||
self.__client.set(self._detector_metadata_key(), json.dumps(safe_payload))
|
||||
return safe_payload
|
||||
|
||||
@property
|
||||
def cached_detector_metadata(self) -> dict:
|
||||
return self.get_detector_metadata()
|
||||
|
||||
@property
|
||||
def cached_dtz_low(self) -> float | None:
|
||||
value = self.get_detector_metadata().get("dtz_low")
|
||||
return self._coerce_optional_float(value)
|
||||
|
||||
@property
|
||||
def cached_dtz_high(self) -> float | None:
|
||||
value = self.get_detector_metadata().get("dtz_high")
|
||||
return self._coerce_optional_float(value)
|
||||
|
||||
@property
|
||||
def local_contact_links(self) -> dict[str, str | None]:
|
||||
detector_frontend = None
|
||||
|
||||
+86
-52
@@ -70,7 +70,6 @@ from aare.daq.operations.common.runtime import (
|
||||
)
|
||||
from aare.devices.area_detector import AutoEnum
|
||||
from aare.devices.jfjoch import JFJochWrapper
|
||||
from aare.devices.mx_lib import clean_filename
|
||||
|
||||
from aare.common.exception_handler import (
|
||||
TransformationInvalidException,
|
||||
@@ -241,6 +240,36 @@ class AareDAQ:
|
||||
pgroup_provider=_DAQPGroupProvider(self),
|
||||
)
|
||||
|
||||
def _cached_detector_metadata(self) -> dict:
|
||||
return self.__cfg.cached_detector_metadata
|
||||
|
||||
def refresh_detector_metadata_cache(self) -> dict[str, object]:
|
||||
det_cfg = self.__jfjoch.detector()
|
||||
dtz_low = self.__cfg._coerce_optional_float(self.__devs.dtz_low)
|
||||
dtz_high = self.__cfg._coerce_optional_float(self.__devs.dtz_high)
|
||||
|
||||
payload = self.__cfg.set_detector_metadata(
|
||||
{
|
||||
"detector_description": det_cfg.description,
|
||||
"detector_serial_number": det_cfg.serial_number,
|
||||
"detector_width": det_cfg.width,
|
||||
"detector_height": det_cfg.height,
|
||||
"pixel_size_mm": det_cfg.pixel_size_mm,
|
||||
"dtz_low": dtz_low,
|
||||
"dtz_high": dtz_high,
|
||||
}
|
||||
)
|
||||
logger.info(
|
||||
"Refreshed hardware metadata cache",
|
||||
extra={
|
||||
"detector_description": payload.get("detector_description"),
|
||||
"detector_serial_number": payload.get("detector_serial_number"),
|
||||
"dtz_low": payload.get("dtz_low"),
|
||||
"dtz_high": payload.get("dtz_high"),
|
||||
},
|
||||
)
|
||||
return payload
|
||||
|
||||
def get_runtime_simulation_state(self) -> dict[str, bool]:
|
||||
return self.__cfg.runtime_simulation_state
|
||||
|
||||
@@ -2083,7 +2112,7 @@ class AareDAQ:
|
||||
def dtz(self) -> float:
|
||||
tmp = self.__cfg.dtz
|
||||
if tmp is None:
|
||||
return 150.0
|
||||
return cfg_get('daq.data_collection_settings.default_raster_scan_settings.dtz', 200)
|
||||
else:
|
||||
return tmp
|
||||
|
||||
@@ -2092,15 +2121,33 @@ class AareDAQ:
|
||||
self.__cfg.try_set_busy(timeout=360)
|
||||
state = self.__cfg.state
|
||||
|
||||
if val < self.__devs.dtz_low or val > self.__devs.dtz_high:
|
||||
dtz_low = self.__cfg.cached_dtz_low
|
||||
dtz_high = self.__cfg.cached_dtz_high
|
||||
|
||||
if dtz_low is None or dtz_high is None:
|
||||
logger.warning("DTZ limits not found in cache, refreshing hardware metadata")
|
||||
try:
|
||||
self.refresh_detector_metadata_cache()
|
||||
except Exception as e:
|
||||
self.__cfg.state_busy = False
|
||||
raise RuntimeError(f"DTZ limits unavailable and refresh failed: {e}") from e
|
||||
dtz_low = self.__cfg.cached_dtz_low
|
||||
dtz_high = self.__cfg.cached_dtz_high
|
||||
|
||||
if dtz_low is None or dtz_high is None:
|
||||
self.__cfg.state_busy = False
|
||||
raise RuntimeError(f"dtz={val} outside limits {self.__devs.dtz_low} to {self.__devs.dtz_high}")
|
||||
raise RuntimeError("DTZ limits are unavailable")
|
||||
|
||||
if val < dtz_low or val > dtz_high:
|
||||
self.__cfg.state_busy = False
|
||||
raise RuntimeError(f"dtz={val} outside limits {dtz_low} to {dtz_high}")
|
||||
|
||||
if state == BeamlineStateEnum.DataCollection:
|
||||
self.__cfg.state_busy = False
|
||||
raise RuntimeError("Cannot set dtz during data collection")
|
||||
elif state == BeamlineStateEnum.SampleAlignment:
|
||||
self.__devs.set_dtz(val, wait=False)
|
||||
|
||||
self.__cfg.dtz = val
|
||||
self.__cfg.state_busy = False
|
||||
|
||||
@@ -2962,30 +3009,39 @@ class AareDAQ:
|
||||
@property
|
||||
def diffraction_geometry(self) -> DiffractionGeometry:
|
||||
try:
|
||||
#start = time.perf_counter()
|
||||
det_cfg = self.__jfjoch.detector()
|
||||
#logger.debug(f"Safe diffraction geometry call took {time.perf_counter() - start:.3f}s")
|
||||
metadata = self._cached_detector_metadata()
|
||||
width = int(metadata.get("detector_width", 1))
|
||||
height = int(metadata.get("detector_height", 1))
|
||||
pixel_size_mm = float(metadata.get("pixel_size_mm", 0.15))
|
||||
detector_description = str(metadata.get("detector_description", "unavailable"))
|
||||
detector_serial_number = str(metadata.get("detector_serial_number", "unavailable"))
|
||||
energy=self.__devs.energy_kev
|
||||
dtz=self.__devs.dtz
|
||||
beam_center=self.__cfg.beam_center
|
||||
return DiffractionGeometry(
|
||||
energy_keV=self.__devs.energy_kev,
|
||||
dtz_mm=self.__devs.dtz,
|
||||
detector_size_pxl=(det_cfg.width, det_cfg.height),
|
||||
pixel_size_mm=det_cfg.pixel_size_mm,
|
||||
beam_center_pxl=self.__cfg.beam_center,
|
||||
detector_description=det_cfg.description,
|
||||
detector_serial_number=det_cfg.serial_number,
|
||||
energy_keV=energy,
|
||||
dtz_mm=dtz,
|
||||
detector_size_pxl=(width, height),
|
||||
pixel_size_mm=pixel_size_mm,
|
||||
beam_center_pxl=beam_center,
|
||||
detector_description=detector_description,
|
||||
detector_serial_number=detector_serial_number,
|
||||
poni_rot1_rad=-0.001396263,
|
||||
poni_rot2_rad=-0.003839724,
|
||||
)
|
||||
except JFJochCommunicationError as e:
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
f"Falling back to default diffraction geometry because detector metadata is unavailable: {e}"
|
||||
f"Falling back to default diffraction geometry because cached detector metadata is unavailable: {e}"
|
||||
)
|
||||
energy=self.__devs.energy_kev
|
||||
dtz=self.__devs.dtz
|
||||
beam_center=self.__cfg.beam_center
|
||||
return DiffractionGeometry(
|
||||
energy_keV=self.__devs.energy_kev,
|
||||
dtz_mm=self.__devs.dtz,
|
||||
energy_keV=energy,
|
||||
dtz_mm=dtz,
|
||||
detector_size_pxl=(1, 1),
|
||||
pixel_size_mm=0.15,
|
||||
beam_center_pxl=self.__cfg.beam_center,
|
||||
beam_center_pxl=beam_center,
|
||||
detector_description="unavailable",
|
||||
detector_serial_number="unavailable",
|
||||
poni_rot1_rad=-0.001396263,
|
||||
@@ -2995,49 +3051,26 @@ class AareDAQ:
|
||||
@property
|
||||
def beamline_status(self) -> BeamlineStatus:
|
||||
try:
|
||||
# start = time.perf_counter()
|
||||
ring_current = self.__devs.ring_current
|
||||
# logger.debug(f"Safe ring current call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
front_light = self.front_light
|
||||
# logger.debug(f"Safe front light call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
back_light = self.back_light
|
||||
# logger.debug(f"Safe back light call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
cryojet_temp = self.__devs.cryojet_temp
|
||||
# logger.debug(f"Safe cryojet temp call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
shutter_open = self.__devs.shutter
|
||||
# logger.debug(f"Safe shutter call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
exp_shutter_open = self.__devs.exp_shutter.state()
|
||||
# logger.debug(f"Safe exp shutter call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
flux = self.__devs.full_flux
|
||||
# logger.debug(f"Safe flux call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
samcam_settings = self.__devs.samcam_settings
|
||||
# logger.debug(f"Safe samcam call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
bl = self.__bl
|
||||
# logger.debug(f"Safe bl call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
transmission = self.__devs.transmission
|
||||
# logger.debug(f"Safe transmission call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
zoom = self.__devs.zoom
|
||||
# logger.debug(f"Safe zoom call took {time.perf_counter() - start:.3f}s")
|
||||
# start = time.perf_counter()
|
||||
commisioning_mode = self.__cfg.commissioning_mode
|
||||
# logger.debug(f"Safe commissioning_mode call took {time.perf_counter() - start:.3f}s")
|
||||
start = time.perf_counter()
|
||||
#TODO work this one out, do it on startup!
|
||||
dtz_min = self.__devs.dtz_low
|
||||
logger.debug(f"Safe dtz_min call took {time.perf_counter() - start:.3f}s")
|
||||
start = time.perf_counter()
|
||||
dtz_max = self.__devs.dtz_high
|
||||
logger.debug(f"Safe dtz_max call took {time.perf_counter() - start:.3f}s")
|
||||
|
||||
dtz_min = self.__cfg.cached_dtz_low
|
||||
dtz_max = self.__cfg.cached_dtz_high
|
||||
|
||||
if dtz_min is None or dtz_max is None:
|
||||
logger.warning("DTZ limits missing from cache, using conservative defaults in beamline_status")
|
||||
dtz_min = 20.0
|
||||
dtz_max = 1000.0
|
||||
return BeamlineStatus(
|
||||
ring_current_mA=ring_current,
|
||||
front_light=front_light,
|
||||
@@ -3051,9 +3084,10 @@ class AareDAQ:
|
||||
transmission=transmission,
|
||||
zoom=zoom,
|
||||
commissioning_mode=commisioning_mode,
|
||||
dtz_min=20,
|
||||
dtz_max=1000,
|
||||
dtz_min=dtz_min,
|
||||
dtz_max=dtz_max,
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
logger.exception(f"Failed to retrieve beamline status: {e}")
|
||||
raise
|
||||
|
||||
+26
-5
@@ -95,6 +95,11 @@ async def lifespan(application: FastAPI):
|
||||
except Exception as e:
|
||||
logger.warning(f"Failed to reset automation progress Redis keys: {e}")
|
||||
|
||||
try:
|
||||
daq.refresh_detector_metadata_cache()
|
||||
except Exception as e:
|
||||
logger.warning(f"Initial hardware metadata refresh failed: {e}")
|
||||
|
||||
# ── Initial TELL sync ──
|
||||
try:
|
||||
daq.sync_current_sample_from_tell(force=True)
|
||||
@@ -331,7 +336,6 @@ async def status(token: str = Depends(oauth2_scheme)) -> DAQStatusModel:
|
||||
else:
|
||||
own_gui = cfg.get_gui_session(data.session)
|
||||
full.open_guis = [own_gui] if own_gui is not None else []
|
||||
|
||||
return full
|
||||
@app.get("/beamline/geometry")
|
||||
async def sample_geometry(
|
||||
@@ -488,6 +492,7 @@ async def smargon(val: SmargonCoordinate, token: str = Depends(oauth2_scheme)):
|
||||
logger.debug(f"Setting smargon to {val}")
|
||||
auth.check_jwt_rw(cfg, auth.parse_token(token))
|
||||
daq.smargon = val
|
||||
logger.debug(f"smargon set to {daq.smargon}")
|
||||
return "OK"
|
||||
|
||||
|
||||
@@ -656,6 +661,7 @@ async def local_contact_simulation_state(token: str = Depends(oauth2_scheme)) ->
|
||||
auth.check_jwt_staff_only(data)
|
||||
return daq.get_runtime_simulation_state()
|
||||
|
||||
|
||||
@app.get("/local_contact/device_state")
|
||||
async def local_contact_device_state(token: str = Depends(oauth2_scheme)) -> dict:
|
||||
"""
|
||||
@@ -720,6 +726,23 @@ async def local_contact_restart_device(
|
||||
result["message"] = f"{device} backend restarted."
|
||||
return result
|
||||
|
||||
@app.post("/local_contact/resync/hardware_metadata")
|
||||
async def local_contact_resync_hardware_metadata(
|
||||
token: str = Depends(oauth2_scheme),
|
||||
) -> dict:
|
||||
"""
|
||||
Refresh cached detector metadata and DTZ limits. Staff only.
|
||||
"""
|
||||
data = auth.parse_token(token)
|
||||
auth.check_jwt_staff_only(data)
|
||||
|
||||
payload = daq.refresh_detector_metadata_cache()
|
||||
return {
|
||||
"ok": True,
|
||||
"message": "Hardware metadata cache resynced.",
|
||||
"payload": payload,
|
||||
}
|
||||
|
||||
@app.post("/beamline/goto_abr_meas_pos")
|
||||
async def goto_abr_meas_pos(token: str = Depends(oauth2_scheme)):
|
||||
"""
|
||||
@@ -1577,8 +1600,6 @@ async def set_smart_params(p: SimpleScanParameters, token: str = Depends(oauth2_
|
||||
Returns:
|
||||
"OK" on success.
|
||||
"""
|
||||
#token_data = auth.parse_token(token)
|
||||
#logger.debug(f"{token_data.session} Try to set smart params: {p}")
|
||||
auth.check_jwt_rw(cfg, auth.parse_token(token))
|
||||
cfg.auto_params = p
|
||||
return "OK"
|
||||
@@ -2410,10 +2431,10 @@ async def maintenance(token: str = Depends(oauth2_scheme)) -> str:
|
||||
|
||||
def main():
|
||||
# Remove in production!
|
||||
urllib3.disable_warnings()
|
||||
#urllib3.disable_warnings()
|
||||
|
||||
# Run the application using uvicorn
|
||||
uvicorn.run("aare.daq.server:app", host="127.0.0.1", port=5210, workers=2, proxy_headers=False, log_config=get_uvicorn_logging_config())
|
||||
uvicorn.run("aare.daq.server:app", host="127.0.0.1", port=5210, workers=4, proxy_headers=False, log_config=get_uvicorn_logging_config())
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
@@ -435,6 +435,11 @@ class LocalContactPanel(QFrame):
|
||||
row,
|
||||
1,
|
||||
)
|
||||
grid.addWidget(self._make_button(
|
||||
"Resync detector/DTZ hardware cache",
|
||||
self._daq.resync_local_contact_detector_metadata,
|
||||
"Resyncing detector metadata and DTZ limits cache.",
|
||||
))
|
||||
|
||||
wrapper = QGroupBox("Detector actions", tab)
|
||||
wrapper_layout = QVBoxLayout(wrapper)
|
||||
|
||||
@@ -93,6 +93,7 @@ class DAQWorker(QObject):
|
||||
fluorimeter_update = Signal(list, list, int)
|
||||
fluorimeter_spectrum_update = Signal(FluorescenceSpectrumOutputModel)
|
||||
sample_resync_completed = Signal(str)
|
||||
detector_metadata_resync_completed = Signal(str)
|
||||
|
||||
automation_critical_failure = Signal(str)
|
||||
manual_collection_critical_failure = Signal(str)
|
||||
@@ -636,6 +637,20 @@ class DAQWorker(QObject):
|
||||
logger.error(f"Sample resync failed: {e}")
|
||||
self.http_error.emit(str(e))
|
||||
|
||||
def _handle_detector_metadata_resync_response(self, reply: QNetworkReply):
|
||||
try:
|
||||
response_data = self.handle_response(reply)
|
||||
payload = json.loads(response_data) if response_data else {}
|
||||
message = str(payload.get("message") or "Hardware metadata cache resynced.")
|
||||
logger.info(message)
|
||||
self.status_message.emit(message, False)
|
||||
self.detector_metadata_resync_completed.emit(message)
|
||||
self.send_status_request()
|
||||
self.load_local_contact_device_state()
|
||||
except Exception as e:
|
||||
logger.error(f"Hardware metadata resync failed: {e}")
|
||||
self.http_error.emit(str(e))
|
||||
|
||||
def _handle_recovery_action_response(self, reply: QNetworkReply, default_message: str):
|
||||
try:
|
||||
response_data = self.handle_response(reply)
|
||||
@@ -1344,6 +1359,18 @@ class DAQWorker(QObject):
|
||||
reply = self.__net_manager.post(request, QByteArray(b""))
|
||||
reply.finished.connect(lambda: self._handle_sample_resync_response(reply))
|
||||
|
||||
@Slot()
|
||||
def resync_local_contact_detector_metadata(self):
|
||||
if self.__base_url is None:
|
||||
logger.info("POST /local_contact/resync/detector_metadata")
|
||||
return
|
||||
|
||||
request = QNetworkRequest(QUrl(f"{self.__base_url}/local_contact/resync/detector_metadata"))
|
||||
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
|
||||
request.setRawHeader(b"Content-Type", b"application/json")
|
||||
reply = self.__net_manager.post(request, QByteArray(b""))
|
||||
reply.finished.connect(lambda: self._handle_detector_metadata_resync_response(reply))
|
||||
|
||||
@Slot(float)
|
||||
def anneal(self, time_s: float):
|
||||
self.generic_post(f"beamline/anneal?time_s={time_s:.1f}")
|
||||
|
||||
Reference in New Issue
Block a user