Files
AareDAQ/src/aare/gui/threads/daq_worker.py
T
appleb_mandClaude Opus 4.8 47e05dc28c GUI/DAQ: hutch-safety gating, error pop-ups, and automation pause-and-wait
Fix the GUI pop-up path and add personnel-safety-system (PSS) gating so
door-open / beam-down / shutter-closed conditions are surfaced and acted on.

- exception pop-ups: connect the previously-orphaned http_error signal;
  failed user operations now raise a modal dialog, background/polling errors
  a non-modal banner.
- PSS device (devices/pss_state.py) reading EH1-PSYS PROHIBITED-STATE /
  ALARM-STATE; new critical DoorSafetyError + DOOR_SAFETY_ERROR code.
- mounting service blocks mount/unmount when the hutch is not prohibited or
  an alarm is active; /status now publishes pss_prohibited / pss_alarm.
- GUI blocks manual mount/unmount and the automation Run button immediately
  (pop-up) on door-open, and shows a warning banner while an alarm is active.
- centralise per-action precondition checks (ring current, safety shutter,
  hutch door) into one combined "continue?" dialog with a session-global
  "don't ask again for 1 hour" snooze, applied to all data-collection buttons.
- live automation pauses and auto-resumes on bad conditions (beam, shutter,
  door, robot) with continue-now / stop overrides, gated by a default-on
  "Pause on bad conditions" checkbox replacing the dead CHECK_ENABLED constant.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-24 15:30:08 +02:00

2405 lines
96 KiB
Python

import copy
import os
import random
import re
import time
import json
import logging
from datetime import datetime
from dataclasses import dataclass
from collections import deque
from typing import cast, Literal
from PySide6.QtCore import Signal, QUrl, Slot, QTimer, QObject, QByteArray
from PySide6.QtNetwork import QNetworkAccessManager, QNetworkRequest, QNetworkReply, QSslError
from jfjoch_client import ScanResult, ScanResultImagesInner
from aare.common.auth_models import BatonStatus
from aare.common.coordinate import SmargonCoordinate, AerotechCoordinate
from aare.common.error_codes import export_error_codes, AareErrorCode, AuthErrorCode
from aare.common.models import (
DAQStatusModel,
SampleShortInfoList,
SampleShortInfo,
SampleCameraSettings,
AutofocusSettings,
SimpleScanParameters,
FluorescenceSpectrumParameterModel,
FluorescenceSpectrumOutputModel,
OpenGuiSessionInfo,
)
from aare.common.automation_models import (
AutomationProgress,
LogEvent,
StepState,
StepStatus,
WorkflowStateKind,
)
from aare.common.recurrence_watcher import (
DEFAULT_WATCHERS,
WatcherTrip,
create_default_watchers,
load_watcher_threshold_overrides,
redis_key_to_env_var,
resolve_exception_class,
)
from aare.common.raster_grid import RasterGridRequest, CompletedRasterGrid, CompletedRasterGridElem
from aare.common.rotation_scan import RotationScanRequest, CompletedRotationScan
from aare.common.logger_config import setup_logger
logger = setup_logger("aareGUI")
SPREADHSEET_FREQUENCY = 25 # Every 5 seconds
@dataclass(frozen=True)
class ErrorInfo:
critical: bool
code: str | None
exception_class: str | None
message: str
context: dict
class DAQWorker(QObject):
update = Signal(DAQStatusModel)
spreadsheet = Signal(SampleShortInfoList)
reference_tools = Signal(SampleShortInfoList)
http_error = Signal(str)
status_message = Signal(str, bool)
automation_progress = Signal(object)
gui_sessions_loaded = Signal(list)
gui_close_requested = Signal(int, int, str)
recovery_action_completed = Signal(str)
bec_user_macros_loaded = Signal(list)
bec_devices_loaded = Signal(list)
local_contact_simulation_state_loaded = Signal(dict)
local_contact_device_state_loaded = Signal(dict)
local_contact_links_loaded = Signal(dict)
local_contact_transfer_error = Signal(str)
polled_devices_status = Signal(str, bool) # (message, is_error)
detector_error = Signal(str, bool) # (message, is_error)
auth_error = Signal()
sample_missing = Signal(str)
# Failure of a user-triggered operation POST (mount, unmount, ...). Carries
# (title, message, critical) and is surfaced as a modal pop-up by the GUI.
operation_failed = Signal(str, str, bool)
# Hutch PSS alarm (ALARM-STATE != 0) became active/inactive. Edge-triggered
# from /status; surfaced as a non-modal warning banner.
pss_alarm_changed = Signal(bool)
automated_scan_done = Signal(int, bool, str) # sample ID, success
run_number_incremented = Signal()
raster_scan_completed = Signal(CompletedRasterGrid)
standard_scan_completed = Signal(CompletedRotationScan)
raster_generated_by_ml = Signal(RasterGridRequest)
face_detection_result = Signal(dict)
staff_pgroups_loaded = Signal(list)
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)
error_codes_loaded = Signal(dict)
last_error_payload_changed = Signal(dict)
last_error_payloads_changed = Signal(list)
baton_status_changed = Signal(BatonStatus)
baton_request_result = Signal(dict)
baton_response_result = Signal(dict)
baton_incoming_request = Signal(dict)
baton_timeout_checked = Signal(dict)
def __init__(self, base_url: str | None, token: str, parent=None):
"""
Initialize the DAQWorker.
Args:
base_url: The base URL of the DAQ server.
token: The authentication token.
parent: The parent QObject.
"""
super().__init__(parent)
self._active_status_error_key = None
self.__token = token
self.__base_url = base_url
self.__net_manager = QNetworkAccessManager()
self.__net_manager.sslErrors.connect(self._handle_ssl_errors)
self.__timer = QTimer()
self.__timer.setInterval(500)
self.__timer.timeout.connect(self.regular_update)
self.__timer.start()
self.__counter = 0
self._automation_progress_buffer = ""
self._cleanup_done = False
self._last_auth_error_log_ts = 0.0
self._auth_error_min_interval = 10.0
self._last_status_request_ts = 0.0
self._status_request_min_interval = 0.5
self._smargon_retry_interval_s = 2.0
self._smargon_log_min_interval_s = 10.0
self._last_status_can_read: bool | None = None
self._device_error_log_min_interval_s = 10.0
self._last_device_error_log_ts: dict[str, float] = {"tell": 0.0, "smargon": 0.0}
self._last_device_error_log_key: dict[str, str | None] = {"tell": None, "smargon": None}
self._last_status_error = None
self._smargon_error_active = False
self._last_smargon_log_ts = 0.0
self._last_smargon_log_key: str | None = None
self._last_error_payload: dict = {}
self._last_error_payloads = deque(maxlen=10)
self._last_tell_connected: bool | None = None
self._last_tell_error: str | None = None
self._last_smargon_connected: bool | None = None
self._last_smargon_error: str | None = None
self._last_aerotech_connected: bool | None = None
self._last_aerotech_error: str | None = None
# Edge-trigger the hutch PSS alarm banner; start False so a clear hutch
# at startup doesn't emit a spurious "cleared" notification.
self._last_pss_alarm: bool = False
self._server_connected: bool | None = None
self._last_server_error: str | None = None
self._has_seen_disconnection: bool = False
self._server_was_disconnected: bool = False
self._last_polled_msg: str | None = None
self._last_polled_is_error: bool | None = None
self._last_detector_msg: str | None = None
self._last_detector_is_error: bool | None = None
self._last_manual_collection_critical_msg: str | None = None
self._baton_stream_reply: QNetworkReply | None = None
self._last_baton_status: BatonStatus | None = None
self._baton_timeout_timer = QTimer(self)
self._baton_timeout_timer.setInterval(1000)
self._baton_timeout_timer.timeout.connect(self.check_baton_timeout)
self._face_detection_stream_reply: QNetworkReply | None = None
self._automation_progress_stream_reply: QNetworkReply | None = None
self._face_detection_stream_blocked_403 = False
self._automation_progress_stream_blocked_403 = False
self._last_seen_automation_event_ts: dict[int | None, datetime] = {}
self._recurrence_watchers = self._load_recurrence_watchers()
self._local_contact_metadata_poll_enabled = True
self._local_contact_metadata_error: str | None = None
if self.__base_url is not None:
self.start_face_detection_stream()
self.start_baton_stream()
self.start_automation_progress_stream()
def get_last_error_payload(self) -> dict:
return dict(self._last_error_payload or {})
def get_last_error_payloads(self) -> list[dict]:
return [dict(p or {}) for p in list(self._last_error_payloads)]
def _set_last_error_payload(self, payload: dict) -> None:
payload = payload or {}
self._last_error_payload = payload
self._last_error_payloads.append(payload)
self.last_error_payload_changed.emit(self.get_last_error_payload())
self.last_error_payloads_changed.emit(self.get_last_error_payloads())
def _emit_polled_device_message(self, msg: str, is_error: bool) -> None:
"""Emit polled device status message with deduplication."""
if (msg, is_error) == (self._last_polled_msg, self._last_polled_is_error):
return
self._last_polled_msg = msg
self._last_polled_is_error = is_error
self.polled_devices_status.emit(msg, is_error)
def _emit_detector_message(self, msg: str, is_error: bool) -> None:
"""Emit detector/request-time error message with deduplication."""
if (msg, is_error) == (self._last_detector_msg, self._last_detector_is_error):
return
self._last_detector_msg = msg
self._last_detector_is_error = is_error
self.detector_error.emit(msg, is_error)
def _emit_status_if_changed(self, key: str | None, message: str | None, is_error: bool) -> None:
"""
Emit polled device status to the primary alert banner.
Used for Server/Tell/Smargon/Aerotech connection status.
"""
if not message:
self._active_status_error_key = None
self._emit_polled_device_message("", False)
return
if is_error:
self._has_seen_disconnection = True
if self._active_status_error_key != key:
self._active_status_error_key = key
self._emit_polled_device_message(message, True)
return
if not self._has_seen_disconnection:
self._active_status_error_key = None
return
self._active_status_error_key = None
self._emit_polled_device_message(message, False)
def _log_smargon_throttled(self, *, endpoint: str | None, message: str) -> None:
"""
Log immediately if endpoint/message changed; otherwise at most every N seconds.
"""
key = f"{endpoint or ''}|{message or ''}"
now = time.monotonic()
if self._last_smargon_log_key != key:
self._last_smargon_log_key = key
self._last_smargon_log_ts = now
logger.error(f"Smargon connection error (endpoint={endpoint}): {message}")
return
if now - self._last_smargon_log_ts >= self._smargon_log_min_interval_s:
self._last_smargon_log_ts = now
logger.error(f"Smargon connection error (endpoint={endpoint}): {message}")
def _log_device_error_throttled(self, *, device: str, message: str | None) -> None:
"""
Log device error string immediately if it changes; otherwise at most every N seconds.
Intended for status-poll derived errors like tell_error / smargon_error.
"""
if not message:
return
device = (device or "unknown").lower()
now = time.monotonic()
key = str(message)
last_key = self._last_device_error_log_key.get(device)
last_ts = self._last_device_error_log_ts.get(device, 0.0)
if last_key != key:
self._last_device_error_log_key[device] = key
self._last_device_error_log_ts[device] = now
logger.error(f"{device.upper()} error: {message}")
return
if now - last_ts >= self._device_error_log_min_interval_s:
self._last_device_error_log_ts[device] = now
logger.error(f"{device.upper()} error: {message}")
def _restart_blocked_sse_streams_if_access_restored(self) -> None:
if self._last_status_can_read is not True:
return
if self._face_detection_stream_blocked_403:
self._face_detection_stream_blocked_403 = False
self.start_face_detection_stream()
if self._automation_progress_stream_blocked_403:
self._automation_progress_stream_blocked_403 = False
self.start_automation_progress_stream()
@Slot()
def regular_update(self):
"""
Periodically triggered slot to update spreadsheet and request status from server.
"""
if self.__counter % SPREADHSEET_FREQUENCY == 0:
self.load_spreadsheet()
self.load_reference_tools()
self.__counter = (self.__counter + 1) % SPREADHSEET_FREQUENCY
self.send_status_request()
def send_status_request(self):
"""
Send an asynchronous HTTP GET request to the /status endpoint.
"""
if self.__base_url is None:
return
now = time.monotonic()
if now - self._last_status_request_ts < self._status_request_min_interval:
return
self._last_status_request_ts = now
request = QNetworkRequest(QUrl(f"{self.__base_url}/status"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self.handle_status_response(reply))
_SELF_SIGNED_ERRORS = {
QSslError.SslError.SelfSignedCertificate,
QSslError.SslError.SelfSignedCertificateInChain,
}
@staticmethod
def _handle_ssl_errors(reply: QNetworkReply, errors: list):
self_signed = [e for e in errors if e.error() in DAQWorker._SELF_SIGNED_ERRORS]
if self_signed:
reply.ignoreSslErrors(self_signed)
@staticmethod
def handle_response(reply: QNetworkReply):
"""
Common handler for QNetworkReply to read data and handle errors.
Args:
reply: The QNetworkReply object.
Returns:
The response data as a UTF-8 string.
Raises:
RuntimeError: If the reply contains a network error.
"""
if reply.error() == QNetworkReply.NetworkError.NoError:
response_data = reply.readAll().data().decode("utf-8")
reply.deleteLater()
return response_data
else:
reply.deleteLater()
logger.error(f"Error in response: {reply.errorString()}")
raise RuntimeError(reply.errorString())
def _compose_device_status_message(
self,
*,
tell_conn: bool,
smargon_conn: bool,
aerotech_conn: bool,
tell_changed: bool,
smargon_changed: bool,
aerotech_changed: bool,
) -> tuple[str | None, str | None, bool]:
disconnected: list[str] = []
restored: list[str] = []
if not tell_conn:
disconnected.append("TELL")
elif tell_changed:
restored.append("TELL")
if not smargon_conn:
disconnected.append("Smargon")
elif smargon_changed:
restored.append("Smargon")
if not aerotech_conn:
disconnected.append("Aerotech")
elif aerotech_changed:
restored.append("Aerotech")
if disconnected:
if len(disconnected) == 3:
return (
"all-devices-down",
"TELL, Smargon, and Aerotech disconnected.",
True,
)
if len(disconnected) == 2:
return (
"+".join(sorted(d.lower() for d in disconnected)) + "-down",
f"{disconnected[0]} and {disconnected[1]} disconnected.",
True,
)
device = disconnected[0]
return (
f"{device.lower()}-down",
f"{device} disconnected.",
True,
)
if restored:
if len(restored) == 3:
return (None, "TELL, Smargon, and Aerotech reconnected.", False)
if len(restored) == 2:
return (None, f"{restored[0]} and {restored[1]} reconnected.", False)
return (None, f"{restored[0]} reconnected.", False)
return (None, None, False)
@Slot(QNetworkReply)
def handle_status_response(self, reply: QNetworkReply):
"""
Handle the server's response to a status request.
Args:
reply: The QNetworkReply object.
"""
try:
response_data = self.handle_response(reply)
parsed_response = DAQStatusModel.model_validate_json(response_data)
self.update.emit(parsed_response)
pss_alarm = bool(getattr(getattr(parsed_response, "bl", None), "pss_alarm", False))
if pss_alarm != self._last_pss_alarm:
self._last_pss_alarm = pss_alarm
self.pss_alarm_changed.emit(pss_alarm)
self._last_status_can_read = True
self._restart_blocked_sse_streams_if_access_restored()
# Handle server reconnection - show "Server reconnected" not device messages
if self._server_connected is False and self._has_seen_disconnection:
self._emit_status_if_changed(None, "Server reconnected.", False)
# Reset device states so we don't also emit device reconnection messages
self._last_tell_connected = None
self._last_smargon_connected = None
self._last_aerotech_connected = None
self._server_was_disconnected = False
self._server_connected = True
self._last_server_error = None
smargon_conn = bool(getattr(parsed_response, "smargon_connected", True))
smargon_err = getattr(parsed_response, "smargon_error", None)
tell_conn = bool(getattr(parsed_response, "tell_connected", True))
tell_err = getattr(parsed_response, "tell_error", None)
aerotech_conn = bool(getattr(parsed_response, "aerotech_connected", True))
aerotech_err = getattr(parsed_response, "aerotech_error", None)
tell_err_text = None if tell_err is None else str(tell_err).strip()
smargon_err_text = None if smargon_err is None else str(smargon_err).strip()
aerotech_err_text = None if aerotech_err is None else str(aerotech_err).strip()
# Skip device status processing if we just reconnected from server down
# (we already showed "Server reconnected")
if self._last_tell_connected is None and self._last_smargon_connected is None and self._last_aerotech_connected is None:
# First status after startup or server reconnect - just record states, don't emit
self._last_tell_connected = tell_conn
self._last_smargon_connected = smargon_conn
self._last_aerotech_connected = aerotech_conn
self._last_tell_error = tell_err_text
self._last_smargon_error = smargon_err_text
self._last_aerotech_error = aerotech_err_text
return
tell_changed = (
self._last_tell_connected != tell_conn
or self._last_tell_error != tell_err_text
)
smargon_changed = (
self._last_smargon_connected != smargon_conn
or self._last_smargon_error != smargon_err_text
)
aerotech_changed = (
self._last_aerotech_connected != aerotech_conn
or self._last_aerotech_error != aerotech_err_text
)
self._last_tell_connected = tell_conn
self._last_smargon_connected = smargon_conn
self._last_aerotech_connected = aerotech_conn
self._last_tell_error = tell_err_text
self._last_smargon_error = smargon_err_text
self._last_aerotech_error = aerotech_err_text
status_key, status_msg, is_error = self._compose_device_status_message(
tell_conn=tell_conn,
smargon_conn=smargon_conn,
aerotech_conn=aerotech_conn,
tell_changed=tell_changed,
smargon_changed=smargon_changed,
aerotech_changed=aerotech_changed,
)
self._emit_status_if_changed(status_key, status_msg, is_error)
if not tell_conn:
self._log_device_error_throttled(device="tell", message=tell_err)
if not smargon_conn:
self._log_device_error_throttled(device="smargon", message=smargon_err)
if not aerotech_conn:
self._log_device_error_throttled(device="aerotech", message=aerotech_err)
except Exception as e:
status=None
err_msg = str(e)
try:
status = reply.attribute(QNetworkRequest.Attribute.HttpStatusCodeAttribute)
raw_body = reply.readAll().data().decode("utf-8")
if raw_body:
body_json = json.loads(raw_body)
if isinstance(body_json, dict):
err_msg = body_json.get("message") or body_json.get("detail") or err_msg
if body_json.get("code") == "AEROTECH_UNAVAILABLE":
extra = body_json.get("extra") or {}
if isinstance(extra, dict):
aerotech_detail = extra.get("message")
if aerotech_detail:
err_msg = aerotech_detail
except Exception:
logger.error(f"Exception reading status response: {e}")
if status == 403:
self._last_status_can_read = False
elif status is not None:
self._last_status_can_read = None
if status == 503:
if self._server_connected is not False:
self._has_seen_disconnection = True
self._server_was_disconnected = True
self._emit_status_if_changed(
"server-down",
f"Aerotech unavailable: {err_msg}",
True,
)
else:
if self._server_connected is not False:
self._has_seen_disconnection = True
self._server_was_disconnected = True
self._emit_status_if_changed(
"server-down",
"Server disconnected. Reconnecting...",
True,
)
self._server_connected = False
self._last_server_error = err_msg
self._last_tell_connected = None
self._last_smargon_connected = None
self._last_aerotech_connected = None
logger.error(f"Exception from status response: {e}")
@Slot(QNetworkReply)
def handle_spreadsheet_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
parsed_response = SampleShortInfoList.model_validate_json(response_data)
self.spreadsheet.emit(parsed_response)
except Exception as e:
logger.error(f"Exception from spreadsheet response: {e}")
self.http_error.emit(str(e))
@Slot(QNetworkReply)
def handle_reference_tools_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
parsed_response = SampleShortInfoList.model_validate_json(response_data)
self.reference_tools.emit(parsed_response)
except Exception as e:
logger.error(f"Exception from reference tools response: {e}")
self.http_error.emit(str(e))
def handle_req_response(self, reply: QNetworkReply):
if reply.error() != QNetworkReply.NetworkError.NoError:
status, err_details, body_json = self._extract_reply_error_details(reply)
error_info = self._classify_reply_error(status, err_details, body_json)
raw_body = ""
try:
raw_body = reply.readAll().data().decode("utf-8")
except Exception:
pass
try:
url = reply.request().url().toString()
except Exception:
url = ""
net_err = reply.error()
net_err_name = getattr(net_err, "name", None)
net_err_value = getattr(net_err, "value", None)
self._set_last_error_payload({
"url": url,
"http_status": int(status) if status is not None else None,
"network_error": net_err_name or str(net_err),
"network_error_value": int(net_err_value) if isinstance(net_err_value, int) else None,
"error_string": str(reply.errorString()),
"body_raw": raw_body,
"body_json": body_json,
})
if self._is_auth_error_code(error_info.code) or status == 401:
now = time.monotonic()
if now - self._last_auth_error_log_ts > self._auth_error_min_interval:
logger.error(f"{error_info.message}: baton taken by another user")
self._last_auth_error_log_ts = now
self.auth_error.emit()
elif status in (404, 410, 417):
self.sample_missing.emit(error_info.message)
else:
logger.error(f"{error_info.message}")
title = self._operation_error_title(error_info.exception_class)
self.operation_failed.emit(title, error_info.message, error_info.critical)
reply.deleteLater()
@staticmethod
def _operation_error_title(exception_class: str | None) -> str:
"""Human-friendly dialog title from an exception class name.
``MountingFailed`` -> ``"Mounting Failed"``; falls back to
``"Operation Failed"`` when the class is unknown.
"""
if not exception_class:
return "Operation Failed"
spaced = re.sub(r"(?<=[a-z0-9])(?=[A-Z])", " ", exception_class)
spaced = spaced.replace("Exception", "Error").strip()
return spaced or "Operation Failed"
def _handle_sample_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 "TELL sample cache resynced.")
logger.info(message)
self.status_message.emit(message, False)
self.sample_resync_completed.emit(message)
self.send_status_request()
except Exception as e:
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)
payload = json.loads(response_data) if response_data else {}
message = str(payload.get("message") or default_message)
logger.info(message)
self.status_message.emit(message, False)
self.recovery_action_completed.emit(message)
self.send_status_request()
except Exception as e:
logger.error(f"Recovery action failed: {e}")
self.http_error.emit(str(e))
def _emit_placeholder_local_contact_action(self, action_name: str) -> None:
message = f"{action_name}: This hasn't been connected yet."
logger.info(message)
self.status_message.emit(message, False)
def generic_post(self, url: str, body: str = ""):
"""
Send a generic HTTP POST request to the server.
Args:
url: The relative URL (endpoint).
body: The request body string.
"""
if self.__base_url is None:
logger.info(f"POST /{url}: {body}")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/{url}"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
if str:
request.setRawHeader(b"Content-Type", b"application/json")
reply = self.__net_manager.post(request, QByteArray(body.encode("utf-8")))
reply.finished.connect(lambda: self.handle_req_response(reply))
def generic_put(self, url: str, body: str = ""):
"""
Send a generic HTTP PUT request to the server.
Args:
url: The relative URL (endpoint).
body: The request body string.
"""
if self.__base_url is None:
logger.info(f"PUT /{url}: {body}")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/{url}"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
if str:
request.setRawHeader(b"Content-Type", b"application/json")
reply = self.__net_manager.put(request, QByteArray(body.encode("utf-8")))
reply.finished.connect(lambda: self.handle_req_response(reply))
def generic_delete(self, url: str):
"""
Send a generic HTTP DELETE request to the server.
Args:
url: The relative URL (endpoint).
"""
if self.__base_url is None:
logger.info(f"DELETE /{url}")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/{url}"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.deleteResource(request)
reply.finished.connect(lambda: self.handle_req_response(reply))
@Slot(float)
def set_omega(self, f: float):
self.generic_put(f"beamline/omega?val={f:.3f}")
@Slot(float)
def set_omega_rel(self, f: float):
self.generic_put(f"beamline/omega_rel?val={f:.3f}")
@Slot(float)
def zoom(self, f: float):
self.generic_put(f"beamline/zoom?val={f:.3f}")
@Slot()
def mono_pitch_scan(self):
self.generic_post("beamline/mono_pitch_scan")
@Slot(float)
def change_energy(self, value: float):
self.generic_put(f"beamline/change_energy?value={value:.3f}")
@Slot(int)
def front_light(self, v: int):
self.generic_put(f"beamline/front_light?val={v:d}")
@Slot(int)
def back_light(self, v: int):
self.generic_put(f"beamline/back_light?val={v:d}")
@Slot()
def close_shutter(self):
self.generic_post(f"beamline/shutter?val=false")
@Slot()
def open_shutter(self):
self.generic_post(f"beamline/shutter?val=true")
@Slot()
def center_loop(self):
self.generic_post("alc/center_loop")
@Slot()
def force_session(self):
self.generic_post("access/force_current_session")
@Slot()
def end_session(self):
self.generic_post("access/end_session")
@Slot()
def dewar_exchange(self):
self.generic_post("state/dewar_exchange")
@Slot()
def sample_exchange(self):
self.generic_post("state/sample_exchange")
@Slot()
def sample_alignment(self):
self.generic_post("state/sample_alignment")
@Slot()
def beam_location(self):
self.generic_post("state/beam_location")
@Slot()
def beamstop_alignment(self):
self.generic_post("state/beamstop_alignment")
@Slot()
def flux_measurement(self):
self.generic_post("state/flux_measurement")
@Slot()
def data_collection(self):
self.generic_post("state/data_collection")
@Slot()
def robot_sample_exchange(self):
self.generic_post("state/robot_sample_exchange")
@Slot()
def xray_fluorescence(self):
self.generic_post("state/xray_fluorescence")
@Slot()
def xtal_snapshot(self):
self.generic_post("state/xtal_snapshot")
@Slot(str)
def free_beamline(self, confirmation_code: str):
if self.__base_url is None:
logger.info("POST /state/free_beamline")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/state/free_beamline"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
request.setRawHeader(b"Content-Type", b"application/json")
body = json.dumps({"confirmation_code": confirmation_code})
reply = self.__net_manager.post(request, QByteArray(body.encode("utf-8")))
reply.finished.connect(
lambda: self._handle_recovery_action_response(reply, "Beamline busy flag cleared.")
)
@Slot(str)
def take_over_beamline(self, confirmation_code: str):
if self.__base_url is None:
logger.info("POST /access/take_over_beamline")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/access/take_over_beamline"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
request.setRawHeader(b"Content-Type", b"application/json")
body = json.dumps({"confirmation_code": confirmation_code})
reply = self.__net_manager.post(request, QByteArray(body.encode("utf-8")))
reply.finished.connect(
lambda: self._handle_recovery_action_response(reply, "Beamline session taken over.")
)
@Slot(str)
def recover_beamline(self, confirmation_code: str):
if self.__base_url is None:
logger.info("POST /recovery/recover_beamline")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/recovery/recover_beamline"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
request.setRawHeader(b"Content-Type", b"application/json")
body = json.dumps({"confirmation_code": confirmation_code})
reply = self.__net_manager.post(request, QByteArray(body.encode("utf-8")))
reply.finished.connect(
lambda: self._handle_recovery_action_response(reply, "Beamline recovered to Maintenance.")
)
@Slot(str)
def recovery_unmount_sample(self, confirmation_code: str):
if self.__base_url is None:
logger.info("POST /recovery/unmount_sample")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/recovery/unmount_sample"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
request.setRawHeader(b"Content-Type", b"application/json")
body = json.dumps({"confirmation_code": confirmation_code})
reply = self.__net_manager.post(request, QByteArray(body.encode("utf-8")))
reply.finished.connect(
lambda: self._handle_recovery_action_response(reply, "Recovery unmount completed.")
)
@Slot(str)
def set_pgroup(self, val: str):
if val == "":
self.generic_delete("access/pgroup")
else:
self.generic_put(f"access/pgroup?val={val}")
self._face_detection_stream_blocked_403 = False
self._automation_progress_stream_blocked_403 = False
self.send_status_request()
self.start_face_detection_stream()
self.start_automation_progress_stream()
@Slot(QNetworkReply)
def _handle_all_pgroups_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
import json
arr = json.loads(response_data) if response_data else []
if not isinstance(arr, list):
raise RuntimeError("Invalid all_pgroups payload")
# Ensure list[str]
out = [str(x) for x in arr if isinstance(x, (str, int))]
self.staff_pgroups_loaded.emit(out)
except Exception as e:
logger.error(f"Exception from all_pgroups response: {e}")
self.http_error.emit(str(e))
@Slot()
def get_all_pgroups(self):
logger.debug("PUT /access/all_pgroups")
if self.__base_url is None:
self.staff_pgroups_loaded.emit([])
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/access/all_pgroups"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
request.setRawHeader(b"Content-Type", b"application/json")
reply = self.__net_manager.put(request, QByteArray(b""))
reply.finished.connect(lambda: self._handle_all_pgroups_response(reply))
@Slot(SampleCameraSettings)
def samcam_settings(self, s: SampleCameraSettings):
self.generic_put("beamline/samcam", s.model_dump_json())
@Slot(AutofocusSettings)
def autofocus(self, f: AutofocusSettings):
self.generic_post("samcam/autofocus", f.model_dump_json())
@Slot(SmargonCoordinate)
def move_smargon(self, coord: SmargonCoordinate):
self.generic_put("beamline/smargon", coord.model_dump_json())
def handle_rotation_scan_response(self, reply: QNetworkReply):
try:
# Check for HTTP errors first
if reply.error() != QNetworkReply.NetworkError.NoError:
status, err_msg, body_json = self._extract_reply_error_details(reply)
if self._is_detector_state_failure_message(err_msg):
self._emit_detector_message(f"JFJoch: {err_msg}", is_error=True)
else:
code = body_json.get("code", "") if isinstance(body_json, dict) else ""
if code == AareErrorCode.JF_JOCH_COMMUNICATION_ERROR or status == 503:
self._emit_detector_message(f"JFJoch: {err_msg}", is_error=True)
if self._is_critical_detector_failure(status, body_json, err_msg):
critical_msg = (
"Manual collection stopped because there is an error with the detector. "
"Please call your local contact.\n\n"
f"Details: {err_msg}"
)
self._emit_detector_message(f"Detector error during manual collection: {err_msg}", is_error=True)
self._emit_manual_collection_critical_failure(critical_msg)
short_msg = err_msg.split("input':", 1)[0].strip() if "input':" in err_msg else err_msg
logger.error(f"Rotation scan failed: {short_msg}")
self.http_error.emit(short_msg)
reply.deleteLater()
return
response_data = self.handle_response(reply)
parsed_response = CompletedRotationScan.model_validate_json(response_data)
self.standard_scan_completed.emit(parsed_response)
except Exception as e:
logger.error(f"Exception from rotation scan response: {e}")
self.http_error.emit(str(e))
finally:
reply.deleteLater()
@Slot(RotationScanRequest)
def standard_scan(self, r: RotationScanRequest):
"""
Trigger a standard rotation scan.
Args:
r: The RotationScanRequest object.
"""
self.run_number_incremented.emit()
request = QNetworkRequest(QUrl(f"{self.__base_url}/scan/rotation"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
request.setRawHeader(b"Content-Type", b"application/json")
body = r.model_dump_json()
reply = self.__net_manager.post(request, QByteArray(body.encode("utf-8")))
reply.finished.connect(lambda: self.handle_rotation_scan_response(reply))
def handle_raster_scan_response(self, reply: QNetworkReply):
try:
# Check for HTTP errors first
if reply.error() != QNetworkReply.NetworkError.NoError:
status, err_msg, body_json = self._extract_reply_error_details(reply)
if self._is_detector_state_failure_message(err_msg):
self._emit_detector_message(f"JFJoch: {err_msg}", is_error=True)
else:
code = body_json.get("code", "") if isinstance(body_json, dict) else ""
if code == AareErrorCode.JF_JOCH_COMMUNICATION_ERROR or status == 503:
self._emit_detector_message(f"JFJoch: {err_msg}", is_error=True)
if self._is_critical_detector_failure(status, body_json, err_msg):
critical_msg = (
"Manual collection stopped because there is an error with the detector. "
"Please call your local contact.\n\n"
f"Details: {err_msg}"
)
self._emit_detector_message(f"Detector error during manual collection: {err_msg}", is_error=True)
self._emit_manual_collection_critical_failure(critical_msg)
logger.error(f"Raster scan failed: {err_msg}")
self.http_error.emit(err_msg)
reply.deleteLater()
return
response_data = self.handle_response(reply)
parsed_response = CompletedRasterGrid.model_validate_json(response_data)
self.raster_scan_completed.emit(parsed_response)
except Exception as e:
logger.error(f"Exception from raster scan response: {e}")
self.http_error.emit(str(e))
finally:
reply.deleteLater()
@Slot(RasterGridRequest)
def raster_scan(self, r: RasterGridRequest):
"""
Trigger a raster scan.
Args:
r: The RasterGridRequest object.
"""
self.run_number_incremented.emit()
if self.__base_url is None:
logger.info(f"POST /scan/raster: {r.model_dump_json()}")
image_number = r.get_image_number()
new_copy = copy.deepcopy(r)
images = []
for i in range(image_number):
images.append(ScanResultImagesInner(
number=i,
efficiency=1.0,
bkg = random.gauss(3.0, 0.1),
spots= random.randint(0, 250),
index= random.randint(0, 1),
b= random.uniform(15.0, 80.0)
))
raster_elem = CompletedRasterGridElem(
request=new_copy,
result=ScanResult(file_prefix=r.file_prefix, images=images),
centre_of_mass = None,
)
reply = CompletedRasterGrid(r=[raster_elem])
self.raster_scan_completed.emit(reply)
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/scan/raster?auto_center=false"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
request.setRawHeader(b"Content-Type", b"application/json")
body = r.model_dump_json()
reply = self.__net_manager.post(request, QByteArray(body.encode("utf-8")))
reply.finished.connect(lambda: self.handle_raster_scan_response(reply))
@Slot(RasterGridRequest)
def raster_scan_auto(self, r: RasterGridRequest):
"""
Trigger an automated raster scan.
Args:
r: The RasterGridRequest object.
"""
self.run_number_incremented.emit()
if self.__base_url is None:
logger.info(f"POST /scan/raster: {r.model_dump_json()}")
image_number = r.get_image_number()
new_copy = copy.deepcopy(r)
images = []
for i in range(image_number):
images.append(ScanResultImagesInner(
number=i,
efficiency=1.0,
bkg=random.gauss(3.0, 0.1),
spots=random.randint(0, 250),
index=random.randint(0, 1),
b=random.uniform(15.0, 80.0)
))
raster_elem = CompletedRasterGridElem(
request=new_copy,
result=ScanResult(file_prefix=r.file_prefix, images=images),
centre_of_mass = None,
)
reply = CompletedRasterGrid(r=[raster_elem])
self.raster_scan_completed.emit(reply)
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/scan/raster?auto_center=true"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
request.setRawHeader(b"Content-Type", b"application/json")
body = r.model_dump_json()
reply = self.__net_manager.post(request, QByteArray(body.encode("utf-8")))
reply.finished.connect(lambda: self.handle_raster_scan_response(reply))
@Slot()
def load_spreadsheet(self):
if self.__base_url is None:
logger.info(f"GET /sample/spreadsheet")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/sample/spreadsheet"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self.handle_spreadsheet_response(reply))
@Slot()
def load_reference_tools(self):
if self.__base_url is None:
logger.info(f"GET /sample/reference_tools")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/sample/reference_tools"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self.handle_reference_tools_response(reply))
def _emit_manual_collection_critical_failure(self, msg: str) -> None:
if not msg:
return
if msg == self._last_manual_collection_critical_msg:
return
self._last_manual_collection_critical_msg = msg
self.manual_collection_critical_failure.emit(msg)
@staticmethod
def _load_recurrence_watchers():
beamline = os.getenv("AARE_BEAMLINE", "default")
watcher_names = [name for name, *_ in DEFAULT_WATCHERS]
overrides = load_watcher_threshold_overrides(
beamline=beamline,
watcher_names=watcher_names,
get_value=lambda k: os.getenv(redis_key_to_env_var(k)),
)
return create_default_watchers(overrides)
def _reset_recurrence_watchers(self) -> None:
for watcher in self._recurrence_watchers:
watcher.reset()
def _observe_recurrence(self, exception_class_name: str | None) -> WatcherTrip | None:
exception_class = resolve_exception_class(exception_class_name)
for watcher in self._recurrence_watchers:
trip = watcher.maybe_trip(exception_class)
if trip is not None:
return trip
return None
def _emit_watcher_trip(self, trip: WatcherTrip) -> None:
message = trip.message()
logger.critical(message)
self.http_error.emit(message)
self.automation_critical_failure.emit(message)
@staticmethod
def _extract_reply_error_details(reply: QNetworkReply) -> tuple[int | None, str, dict | None]:
status = reply.attribute(QNetworkRequest.Attribute.HttpStatusCodeAttribute)
err_msg = reply.errorString()
body_json = None
try:
raw_body = reply.readAll().data().decode("utf-8")
if raw_body:
try:
parsed = json.loads(raw_body)
if isinstance(parsed, dict):
body_json = parsed
err_msg = (
parsed.get("message")
or parsed.get("detail")
or parsed.get("error")
or raw_body
)
else:
err_msg = raw_body
except Exception:
err_msg = raw_body
except Exception:
pass
return status, str(err_msg), body_json
@staticmethod
def _is_critical(body_json: dict | None) -> bool:
if not isinstance(body_json, dict):
return True
return bool(body_json.get("critical", True))
@staticmethod
def _classify_reply_error(
status: int | None,
err_msg: str,
body_json: dict | None,
) -> ErrorInfo:
code = body_json.get("code") if isinstance(body_json, dict) else None
exception_class = body_json.get("exception_class") if isinstance(body_json, dict) else None
context = body_json.get("context") if isinstance(body_json, dict) else {}
if not isinstance(context, dict):
context = {}
return ErrorInfo(
critical=DAQWorker._is_critical(body_json),
code=str(code) if code is not None else None,
exception_class=str(exception_class) if exception_class is not None else None,
message=str(err_msg or ""),
context=context,
)
@staticmethod
def _is_auth_error_code(code: str | None) -> bool:
if code is None:
return False
try:
return code in {item.value for item in AuthErrorCode}
except Exception:
return False
@staticmethod
def _is_detector_state_failure_message(message: str | None) -> bool:
text = str(message or "").lower()
return (
"daq state error" in text
or "must be idle to start measurement" in text
or "must be idle" in text
)
@staticmethod
def _is_critical_detector_failure(
status: int | None,
body_json: dict | None = None,
err_msg: str | None = None,
) -> bool:
try:
status_int = int(status) if status is not None else None
except Exception:
status_int = None
code = None
if isinstance(body_json, dict):
code = body_json.get("code")
text = str(err_msg or "").lower()
if status_int == 500:
return True
if code == AareErrorCode.JF_JOCH_COMMUNICATION_ERROR and status_int is not None and 500 <= status_int < 600:
return True
if "daq state error" in text or "must be idle to start measurement" in text or "must be idle" in text:
return True
return False
def handle_auto_scan_response(self, reply, sample_id: int):
if reply.error() == QNetworkReply.NetworkError.NoError:
resp = reply.readAll().data().decode("utf-8")
logger.info(f"Sample time {resp} s")
self._reset_recurrence_watchers()
self.automated_scan_done.emit(sample_id, True, "")
else:
status, err_str, body_json = self._extract_reply_error_details(reply)
error_info = self._classify_reply_error(status, err_str, body_json)
watcher_trip = self._observe_recurrence(error_info.exception_class)
if self._is_detector_state_failure_message(error_info.message):
self._emit_detector_message(f"JFJoch: {error_info.message}", is_error=True)
if self._is_auth_error_code(error_info.code) or status == 401:
logger.error(f"Error in auto scan: {error_info.message}")
self.automated_scan_done.emit(sample_id, False, "Authentication Error")
self.auth_error.emit()
elif status == 404:
self.sample_missing.emit(error_info.message)
self.automated_scan_done.emit(sample_id, False, "Missing")
elif status == 410:
self.sample_missing.emit(error_info.message)
self.automated_scan_done.emit(sample_id, False, "Warning")
elif status == 417:
self.sample_missing.emit(error_info.message)
self.automated_scan_done.emit(sample_id, False, "Critical")
elif watcher_trip is not None:
self.automated_scan_done.emit(sample_id, False, "Critical")
self._emit_watcher_trip(watcher_trip)
elif error_info.critical:
logger.critical(f"Critical automation failure: {error_info.message}")
self.http_error.emit(error_info.message)
self.automated_scan_done.emit(sample_id, False, "Critical")
self.automation_critical_failure.emit(error_info.message)
else:
logger.error(f"Error in auto scan: {error_info.message}")
self.http_error.emit(error_info.message)
self.automated_scan_done.emit(sample_id, False, error_info.message)
reply.deleteLater()
@Slot(SampleShortInfo)
def automated_scan(self, s: SampleShortInfo):
"""
Trigger a fully automated scan sequence for a sample.
Args:
s: The SampleShortInfo object.
"""
if self.__base_url is None:
logger.info(f"POST /scan/auto: {s.model_dump_json()}")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/scan/auto"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
request.setRawHeader(b"Content-Type", b"application/json")
body = s.model_dump_json()
reply = self.__net_manager.post(request, QByteArray(body.encode("utf-8")))
reply.finished.connect(lambda: self.handle_auto_scan_response(reply, s.db_id))
@Slot(SimpleScanParameters)
def smart_params(self, p: SimpleScanParameters):
if self.__base_url is None:
logger.info(f"POST /scan/smart_params: {p.model_dump_json()}")
return
self.generic_post("scan/smart_params", p.model_dump_json())
@Slot(AerotechCoordinate)
def abr_tweak(self, c: AerotechCoordinate):
self.generic_post("beamline/tweak_abr_meas_pos", c.model_dump_json())
@Slot()
def abr_save(self):
self.generic_post("beamline/save_abr_meas_pos")
@Slot()
def abr_goto_meas(self):
self.generic_post("beamline/goto_abr_meas_pos")
@Slot(float, float)
def beam_mark_add(self, x: float, y: float):
self.generic_post(f"beam_mark/add?x={x}&y={y}")
@Slot()
def beam_mark_clear(self):
self.generic_post(f"beam_mark/clear")
@Slot(float, float)
def beam_center(self, x: float, y: float):
self.generic_post(f"beamline/beam_center?x={x}&y={y}")
@Slot(float, float)
def beam_size_mm(self, x: float, y: float):
self.generic_post(f"beamline/beam_size_mm?x={x}&y={y}")
@Slot()
def resync_sample(self):
if self.__base_url is None:
logger.info("POST /sample/resync")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/sample/resync"))
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_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}")
#TODO combine dry and park and dry
@Slot()
def park_and_dry(self):
self.generic_post("tell/park_and_dry")
@Slot()
def tell_dry(self):
self.generic_post("tell/dry")
@Slot()
def tell_toggle_blower(self):
self.generic_post("tell/toggle_blower")
def _disable_local_contact_metadata_polling(self, message: str) -> None:
if self._local_contact_metadata_poll_enabled is False and self._local_contact_metadata_error == message:
return
self._local_contact_metadata_poll_enabled = False
self._local_contact_metadata_error = message
logger.error(f"Disabling Local Contact metadata polling: {message}")
self.local_contact_transfer_error.emit(message)
def _build_local_contact_error_message(self, context: str, reply: QNetworkReply, exc: Exception) -> str:
status = reply.attribute(QNetworkRequest.Attribute.HttpStatusCodeAttribute)
try:
url = reply.request().url().toString()
except Exception:
url = "unknown-url"
return (
f"{context}\n\n"
f"URL: {url}\n"
f"HTTP status: {status}\n"
f"Error: {exc}"
)
@Slot()
def load_local_contact_simulation_state(self):
if not self._local_contact_metadata_poll_enabled:
return
if self.__base_url is None:
self.local_contact_simulation_state_loaded.emit(
{"bec": False, "detector": False, "tell": False, "aerotech": False, "smargon": False}
)
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/local_contact/simulation_state"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self._handle_local_contact_simulation_state_response(reply))
def _handle_local_contact_simulation_state_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
payload = json.loads(response_data) if response_data else {}
if not isinstance(payload, dict):
raise RuntimeError("Invalid local contact simulation state payload")
self.local_contact_simulation_state_loaded.emit(payload)
except Exception as e:
message = self._build_local_contact_error_message(
"Error transferring information from DAQ while loading Local Contact simulation state.",
reply,
e,
)
self._disable_local_contact_metadata_polling(message)
@Slot()
def load_local_contact_device_state(self):
if not self._local_contact_metadata_poll_enabled:
return
if self.__base_url is None:
self.local_contact_device_state_loaded.emit({})
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/local_contact/device_state"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self._handle_local_contact_device_state_response(reply))
def _handle_local_contact_device_state_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
payload = json.loads(response_data) if response_data else {}
if not isinstance(payload, dict):
raise RuntimeError("Invalid local contact device state payload")
self.local_contact_device_state_loaded.emit(payload)
except Exception as e:
message = self._build_local_contact_error_message(
"Error transferring information from DAQ while loading Local Contact device state.",
reply,
e,
)
self._disable_local_contact_metadata_polling(message)
@Slot()
def load_local_contact_links(self):
if not self._local_contact_metadata_poll_enabled:
return
if self.__base_url is None:
self.local_contact_links_loaded.emit({})
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/local_contact/links"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self._handle_local_contact_links_response(reply))
def _handle_local_contact_links_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
payload = json.loads(response_data) if response_data else {}
if not isinstance(payload, dict):
raise RuntimeError("Invalid local contact links payload")
self.local_contact_links_loaded.emit(payload)
except Exception as e:
message = self._build_local_contact_error_message(
"Error transferring information from DAQ while loading Local Contact links.",
reply,
e,
)
self._disable_local_contact_metadata_polling(message)
@Slot(str, bool)
def set_local_contact_simulation(self, device: str, enabled: bool):
self.generic_post(f"local_contact/simulate/{device}?enabled={str(enabled).lower()}")
@Slot(str)
def restart_local_contact_device(self, device: str):
self.generic_post(f"local_contact/restart/{device}")
@Slot()
def bec_load_user_macros(self):
self.generic_post("bec/load_user_macros")
@Slot()
def bec_list_all_user_macros(self):
if self.__base_url is None:
logger.info("GET /bec/user_macros")
self.bec_user_macros_loaded.emit([])
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/bec/user_macros"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self._handle_bec_user_macros_response(reply))
def _handle_bec_user_macros_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
payload = json.loads(response_data) if response_data else []
if not isinstance(payload, list):
raise RuntimeError("Invalid BEC user macros payload")
self.bec_user_macros_loaded.emit([str(item) for item in payload])
except Exception as e:
logger.error(f"Failed to list BEC user macros: {e}")
self.http_error.emit(str(e))
@Slot()
def bec_list_all_devices(self):
if self.__base_url is None:
logger.info("GET /bec/devices")
self.bec_devices_loaded.emit([])
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/bec/devices"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self._handle_bec_devices_response(reply))
def _handle_bec_devices_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
payload = json.loads(response_data) if response_data else []
if not isinstance(payload, list):
raise RuntimeError("Invalid BEC devices payload")
self.bec_devices_loaded.emit([str(item) for item in payload])
except Exception as e:
logger.error(f"Failed to list BEC devices: {e}")
self.http_error.emit(str(e))
@Slot(str)
def bec_reinitialise_planner_and_position_devices(self, method: str = "auto"):
self.generic_post(f"bec/reinitialise_planner_and_position_devices?method={method}")
@Slot()
def initialise_smargon(self):
self.generic_post("smargon/initialize")
@Slot()
def initialise_aerotech(self):
logger.info("initisalisation does not initisalise aareSCAN but runs homing script")
self.generic_post("aerotech/initialize")
@Slot()
def detector_take_pedestal(self):
self.generic_post("detector/take_pedestal")
@Slot()
def initialise_detector(self):
self.generic_post("detector/initialize")
@Slot()
def unmount(self):
self.generic_post("sample/unmount")
@Slot(SampleShortInfo)
def mount(self, s: SampleShortInfo, reference: bool = False):
self.generic_post(f"sample/mount?dbid={s.db_id}&reference={reference}")
@Slot(SampleShortInfo)
def sample_manual(self, s: SampleShortInfo):
self.generic_post(f"sample/manual", s.model_dump_json())
@Slot()
def cancel(self):
self.generic_post("scan/cancel")
def handle_ml_box_response(self, reply):
try:
response_data = self.handle_response(reply)
if response_data != "":
parsed_response = RasterGridRequest.model_validate_json(response_data)
self.raster_generated_by_ml.emit(parsed_response)
except Exception as e:
logger.error(f"Error in ml box: {e}")
self.http_error.emit(str(e))
@Slot()
def ml_bounding_box(self):
"""
Request an ML-based bounding box for the current sample.
"""
if self.__base_url is None:
logger.info(f"POST /alc/ml_bounding_box")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/alc/ml_bounding_box"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.post(request, QByteArray(b""))
reply.finished.connect(lambda: self.handle_ml_box_response(reply))
def _handle_face_detection_response(self, reply: QNetworkReply):
try:
if reply.error() != QNetworkReply.NetworkError.NoError:
logger.error(f"Error in face detection: {reply.errorString()}")
raise RuntimeError(reply.errorString())
payload = reply.readAll().data().decode("utf-8") or "{}"
import json
data = json.loads(payload)
self.face_detection_result.emit(data)
except Exception as e:
logger.error(f"Error in face detection: {e}")
self.http_error.emit(str(e))
finally:
reply.deleteLater()
def _read_face_detection_stream(self, reply: QNetworkReply):
try:
chunk = reply.readAll().data().decode("utf-8")
for line in chunk.splitlines():
if line.startswith("data:"):
payload = line[5:].strip()
if payload:
data = json.loads(payload)
self.face_detection_result.emit(data)
except Exception as e:
logger.error(f"Face detection stream parse error: {e}")
def _restart_face_detection_stream(self):
reply = self._face_detection_stream_reply
self._face_detection_stream_reply = None
status = None
if reply is not None:
try:
status = reply.attribute(QNetworkRequest.Attribute.HttpStatusCodeAttribute)
except Exception:
status = None
if status == 403:
self._face_detection_stream_blocked_403 = True
logger.info("Face detection SSE access denied; waiting for access change before retrying.")
return
if self.__base_url is not None:
QTimer.singleShot(1000, self.start_face_detection_stream)
@staticmethod
def _parse_automation_progress(progress_payload: dict) -> AutomationProgress:
steps: list[StepState] = []
events: list[LogEvent] = []
step_aliases = {
"final": WorkflowStateKind.FINAL,
"Paused/Finished": WorkflowStateKind.FINAL,
}
for raw_step in progress_payload.get("steps", []):
raw_kind = raw_step.get("step")
raw_status = raw_step.get("status", StepStatus.PENDING.value)
raw_message = raw_step.get("message", "")
step_kind = step_aliases.get(raw_kind)
if step_kind is None:
step_kind = WorkflowStateKind(raw_kind)
step_status = StepStatus(raw_status)
steps.append(
StepState(
step=step_kind,
status=step_status,
message=str(raw_message or ""),
started_at=raw_step.get("started_at"),
completed_at=raw_step.get("completed_at"),
error_code=raw_step.get("error_code"),
)
)
for raw_event in progress_payload.get("events", []):
raw_ts = raw_event.get("ts")
if isinstance(raw_ts, str) and raw_ts:
ts = datetime.fromisoformat(raw_ts)
elif isinstance(raw_ts, datetime):
ts = raw_ts
else:
continue
raw_level = str(raw_event.get("level", "INFO")).upper()
if raw_level not in {"INFO", "WARNING", "ERROR"}:
raw_level = "INFO"
level = cast(Literal["INFO", "WARNING", "ERROR"], raw_level)
events.append(
LogEvent(
ts=ts,
level=level,
code=str(raw_event.get("code", "UNKNOWN")),
exception_class=raw_event.get("exception_class"),
message=str(raw_event.get("message", "")),
sample_id=raw_event.get("sample_id"),
context=dict(raw_event.get("context") or {}),
)
)
return AutomationProgress(
current_step=progress_payload.get("current_step"),
steps=steps,
events=events,
finished=bool(progress_payload.get("finished", False)),
success=progress_payload.get("success"),
samples_in_queue=int(progress_payload.get("samples_in_queue", 0) or 0),
avg_time_per_sample=float(progress_payload.get("avg_time_per_sample", 0.0) or 0.0),
current_sample_name=str(progress_payload.get("current_sample_name") or ""),
)
@staticmethod
def _log_level_for_event(level: str) -> int:
return {
"ERROR": logging.ERROR,
"WARNING": logging.WARNING,
"INFO": logging.INFO,
}.get(str(level).upper(), logging.INFO)
def _emit_automation_progress_events(self, progress: AutomationProgress) -> None:
for event in progress.events:
sample_id = event.sample_id
last_seen = self._last_seen_automation_event_ts.get(sample_id)
if last_seen is not None and event.ts <= last_seen:
continue
self._last_seen_automation_event_ts[sample_id] = event.ts
logger.log(
self._log_level_for_event(event.level),
event.message,
extra={
"code": event.code,
"exception_class": event.exception_class,
"sample_id": event.sample_id,
**dict(event.context or {}),
},
)
watcher_trip = self._observe_recurrence(event.exception_class)
if watcher_trip is not None:
self._emit_watcher_trip(watcher_trip)
def _handle_automation_progress_event(self, payload: str) -> None:
if not payload:
return
logger.info(f"[automation_progress raw] {payload}")
outer = json.loads(payload)
progress_payload = outer.get("progress")
if progress_payload is None:
logger.info("[automation_progress raw] no progress payload in SSE event")
return
progress = self._parse_automation_progress(progress_payload)
self.automation_progress.emit(progress)
self._emit_automation_progress_events(progress)
current = progress.current_step or "Idle"
if progress.finished:
if progress.success is True:
logger.info(f"Automation progress: {current} - finished successfully")
elif progress.success is False:
logger.info(f"Automation progress: {current} - finished with error")
else:
logger.info(f"Automation progress: {current} - finished")
else:
logger.info(f"Automation progress: {current}")
def _process_automation_progress_buffer(self) -> None:
while "\n\n" in self._automation_progress_buffer:
event_data, self._automation_progress_buffer = self._automation_progress_buffer.split("\n\n", 1)
data_lines: list[str] = []
for line in event_data.splitlines():
if line.startswith("data:"):
data_lines.append(line[5:].lstrip())
if not data_lines:
continue
payload = "\n".join(data_lines)
self._handle_automation_progress_event(payload)
def _read_automation_progress_stream(self, reply: QNetworkReply):
try:
chunk = reply.readAll().data().decode("utf-8")
if not chunk:
return
self._automation_progress_buffer += chunk
self._process_automation_progress_buffer()
except Exception as e:
logger.error(f"Automation progress stream parse error: {e}")
def _restart_automation_progress_stream(self):
reply = self._automation_progress_stream_reply
self._automation_progress_stream_reply = None
self._automation_progress_buffer = ""
status = None
if reply is not None:
try:
status = reply.attribute(QNetworkRequest.Attribute.HttpStatusCodeAttribute)
except Exception:
status = None
if status == 403:
self._automation_progress_stream_blocked_403 = True
logger.info("Automation progress SSE access denied; waiting for access change before retrying.")
return
if self.__base_url is not None:
QTimer.singleShot(1000, self.start_automation_progress_stream)
def start_automation_progress_stream(self):
if self.__base_url is None:
return
if self._automation_progress_stream_reply is not None:
return
if self._automation_progress_stream_blocked_403:
return
self._automation_progress_buffer = ""
request = QNetworkRequest(QUrl(f"{self.__base_url}/sse/automation_progress"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.readyRead.connect(lambda: self._read_automation_progress_stream(reply))
reply.finished.connect(self._restart_automation_progress_stream)
self._automation_progress_stream_reply = reply
def start_face_detection_stream(self):
if self.__base_url is None:
return
if self._face_detection_stream_reply is not None:
return
if self._face_detection_stream_blocked_403:
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/sse/face_detection"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.readyRead.connect(lambda: self._read_face_detection_stream(reply))
reply.finished.connect(self._restart_face_detection_stream)
self._face_detection_stream_reply = reply
@Slot()
def face_detection(self, steps: int, step_size: int):
"""
Trigger the face detection procedure.
Args:
steps: Number of rotation steps.
step_size: Degrees per step.
"""
if self.__base_url is None:
logger.info(f"POST /face_detection/run?steps={steps}&step_size={step_size}")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/face_detection/run?steps={steps}&step_size={step_size}"))
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_face_detection_response(reply))
@Slot()
def fluorimeter_start(self):
self.generic_post("fluorimeter/start")
@Slot()
def fluorimeter_start_erase(self):
self.generic_post("fluorimeter/start?erase=true")
@Slot()
def fluorimeter_stop(self):
self.generic_post("fluorimeter/stop")
def _handle_fluorimeter_data(self, reply: QNetworkReply, emit_status: bool = False):
try:
response_data = self.handle_response(reply)
import json
data = json.loads(response_data) if response_data else []
if emit_status:
# fetch status and bkg in parallel (simple sequential here)
status_req = QNetworkRequest(QUrl(f"{self.__base_url}/fluorimeter/status"))
status_req.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
status_reply = self.__net_manager.get(status_req)
status_reply.finished.connect(lambda: self._handle_fluorimeter_status_and_emit(data, status_reply))
else:
self.fluorimeter_update.emit(data, [], -1)
except Exception as e:
logger.error(f"Fluorimeter data error: {e}")
self.http_error.emit(str(e))
def _handle_fluorimeter_status_and_emit(self, data, status_reply: QNetworkReply):
try:
s_payload = self.handle_response(status_reply)
s = int(s_payload) if s_payload not in ("", "null") else -1
b_req = QNetworkRequest(QUrl(f"{self.__base_url}/fluorimeter/background"))
b_req.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
b_reply = self.__net_manager.get(b_req)
b_reply.finished.connect(lambda: self._emit_fluorimeter_with_bkg(data, s, b_reply))
except Exception as e:
logger.error(f"Fluorimeter status error: {e}")
self.http_error.emit(str(e))
def _emit_fluorimeter_with_bkg(self, data, s, b_reply: QNetworkReply):
try:
bkg_json = self.handle_response(b_reply)
import json
bkg = json.loads(bkg_json) if bkg_json else []
self.fluorimeter_update.emit(data, bkg, s)
except Exception as e:
logger.error(f"Fluorimeter background error: {e}")
self.http_error.emit(str(e))
@Slot()
def fluorimeter_spectrum(self, f: FluorescenceSpectrumParameterModel):
if self.__base_url is None:
return
req = QNetworkRequest(QUrl(f"{self.__base_url}/fluorimeter/spectrum"))
req.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
req.setRawHeader(b"Content-Type", b"application/json")
body = f.model_dump_json()
reply = self.__net_manager.post(req, QByteArray(body.encode("utf-8")))
reply.finished.connect(lambda: self._handle_fluorimeter_spectrum(reply))
def _handle_fluorimeter_spectrum(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
parsed_response = FluorescenceSpectrumOutputModel.model_validate_json(response_data)
self.fluorimeter_spectrum_update.emit(parsed_response)
except Exception as e:
logger.error(f"Exception from fluorimeter spectrum: {e}")
self.http_error.emit(str(e))
@Slot()
def fluorimeter_request_snapshot(self):
if self.__base_url is None:
return
req = QNetworkRequest(QUrl(f"{self.__base_url}/fluorimeter/data"))
req.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(req)
reply.finished.connect(lambda: self._handle_fluorimeter_data(reply, emit_status=True))
# Optional: SSE listener for live updates
def start_fluorimeter_stream(self):
if self.__base_url is None:
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/sse/fluorimeter"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.readyRead.connect(lambda: self._read_fluorimeter_stream(reply))
reply.finished.connect(lambda: reply.deleteLater())
def _read_fluorimeter_stream(self, reply: QNetworkReply):
try:
chunk = reply.readAll().data().decode("utf-8")
for line in chunk.splitlines():
if line.startswith("data:"):
payload = line[5:].strip()
if payload:
obj = json.loads(payload)
data = obj.get("data") or []
bkg = obj.get("background") or []
status = obj.get("status", -1)
self.fluorimeter_update.emit(data, bkg, status)
except Exception as e:
logger.error(f"SSE parse error: {e}")
@staticmethod
def _flatten_error_codes_payload(obj: dict) -> dict[str, str]:
"""
Accept either:
- flat: {"INVALID_TOKEN": "INVALID_TOKEN"}
- grouped: {"AuthErrorCode": {"INVALID_TOKEN": "INVALID_TOKEN"}, "AareErrorCode": {...}}
Output is always flat strings, using 'Group.KEY' for grouped input.
"""
out: dict[str, str] = {}
for k, v in (obj or {}).items():
if isinstance(v, dict):
group = str(k)
for kk, vv in v.items():
out[f"{group}.{str(kk)}"] = str(vv)
else:
out[str(k)] = str(v)
return out
@Slot()
def get_error_codes(self) -> None:
"""
Fetch server error codes registry for developer/help UI.
Emits error_codes_loaded(dict).
Server default is grouped. We flatten grouped payloads for existing UI.
"""
if self.__base_url is None:
self.error_codes_loaded.emit(export_error_codes())
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/meta/error-codes"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self._handle_error_codes_response(reply))
def _retry_error_codes_legacy(self) -> None:
request = QNetworkRequest(QUrl(f"{self.__base_url}/meta/error-codes/flat"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self._handle_error_codes_response(reply))
@Slot(QNetworkReply)
def _handle_error_codes_response(self, reply: QNetworkReply) -> None:
try:
status = reply.attribute(QNetworkRequest.Attribute.HttpStatusCodeAttribute)
try:
url = reply.request().url().toString()
except Exception:
url = ""
if int(status) == 404 and url.endswith("/meta/error-codes"):
reply.deleteLater()
logger.error(f"Error codes not found on server.")
return
payload = self.handle_response(reply)
obj = json.loads(payload) if payload else {}
if not isinstance(obj, dict):
raise RuntimeError("Invalid error-codes payload (expected JSON object)")
out = self._flatten_error_codes_payload(obj)
self.error_codes_loaded.emit(out)
except Exception as e:
logger.error(f"Failed to load error codes: {e}")
self.http_error.emit(str(e))
@Slot(str, str)
def send_screenshot_db(self, filename: str = "", message: str = ""):
if self.__base_url is None:
logger.info(f"POST /samcam/send_screenshot_db?filename={filename}&message={message}")
return
from urllib.parse import quote
query = []
filename = filename.strip()
message = message.strip()
if filename:
query.append(f"filename={quote(filename)}")
if message:
query.append(f"message={quote(message)}")
suffix = f"?{'&'.join(query)}" if query else ""
self.generic_post(f"samcam/send_screenshot_db{suffix}")
def start_baton_stream(self):
"""Start SSE stream for baton status updates."""
if self.__base_url is None:
return
if self._baton_stream_reply is not None:
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/sse/baton"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.readyRead.connect(lambda: self._read_baton_stream(reply))
reply.finished.connect(self._restart_baton_stream)
self._baton_stream_reply = reply
def _restart_baton_stream(self):
self._baton_stream_reply = None
if self.__base_url is not None:
QTimer.singleShot(1000, self.start_baton_stream)
def _read_baton_stream(self, reply: QNetworkReply):
try:
chunk = reply.readAll().data().decode("utf-8")
for line in chunk.splitlines():
if line.startswith("data:"):
payload = line[5:].strip()
if payload:
status = BatonStatus.model_validate_json(payload)
if status.you_have_pending_request:
if not self._baton_timeout_timer.isActive():
self._baton_timeout_timer.start()
else:
if self._baton_timeout_timer.isActive():
self._baton_timeout_timer.stop()
if (status.incoming_request and
(self._last_baton_status is None or
not self._last_baton_status.incoming_request)):
self.baton_incoming_request.emit({
"requester": status.pending_request.requester_username if status.pending_request else "Unknown",
"timeout": status.pending_request.timeout_seconds if status.pending_request else 30
})
self._last_baton_status = status
self.baton_status_changed.emit(status)
except Exception as e:
logger.error(f"Baton stream parse error: {e}")
@Slot()
def request_baton(self):
"""
Request the baton to gain write access to the beamline.
"""
if self.__base_url is None:
logger.info("POST /baton/request")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/baton/request"))
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_baton_request_response(reply))
def _handle_baton_request_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
result = json.loads(response_data) if response_data else {}
self.baton_request_result.emit(result)
if result.get("granted"):
self.status_message.emit("Baton acquired", False)
self.send_status_request()
self._restart_blocked_sse_streams_if_access_restored()
elif result.get("pending"):
self.status_message.emit(
f"Request sent - waiting for response ({result.get('timeout_seconds', 30)}s timeout)",
False
)
elif result.get("error"):
self.status_message.emit(result.get("message", "Request failed"), True)
except Exception as e:
logger.error(f"Baton request failed: {e}")
self.http_error.emit(str(e))
@Slot(bool)
def respond_to_baton_request(self, accept: bool):
"""
Respond to an incoming baton request from another user.
Args:
accept: True to grant the baton, False to refuse.
"""
if self.__base_url is None:
logger.info(f"POST /baton/respond?accept={accept}")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/baton/respond?accept={str(accept).lower()}"))
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_baton_response_result(reply))
def _handle_baton_response_result(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
result = json.loads(response_data) if response_data else {}
self.baton_response_result.emit(result)
if result.get("accepted") or result.get("refused"):
self.send_status_request()
self.check_baton_timeout()
self.start_baton_stream()
self._restart_blocked_sse_streams_if_access_restored()
except Exception as e:
logger.error(f"Baton response failed: {e}")
self.http_error.emit(str(e))
@Slot()
def release_baton(self):
"""Release the baton voluntarily."""
self.generic_post("baton/release")
@Slot()
def cancel_baton_request(self):
"""Cancel your pending baton request."""
self.generic_post("baton/cancel")
@Slot()
def check_baton_timeout(self):
"""Poll to check if timeout has been reached."""
if self.__base_url is None:
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/baton/check_timeout"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self._handle_baton_timeout_response(reply))
def _handle_baton_timeout_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
result = json.loads(response_data) if response_data else {}
self.baton_timeout_checked.emit(result)
if result.get("granted") or result.get("queued") or result.get("refused"):
self.send_status_request()
self.start_baton_stream()
self._restart_blocked_sse_streams_if_access_restored()
except Exception as e:
logger.error(f"Baton timeout check failed: {e}")
self.http_error.emit(str(e))
def release_baton_on_close(self):
"""Backward-compatible fallback."""
self.end_session_on_close()
def end_session_on_close(self):
"""
End the GUI session when the window is closed.
This is preferred over only releasing the baton because the backend can:
- remove the Redis GUI session entry immediately
- clear the active beamline session if this GUI owns it
"""
try:
if hasattr(self, "_baton_timeout_timer") and self._baton_timeout_timer is not None:
self._baton_timeout_timer.stop()
self.end_session()
from PySide6.QtCore import QEventLoop, QTimer
loop = QEventLoop()
QTimer.singleShot(500, loop.quit)
loop.exec()
logger.info("End session requested on GUI close")
except Exception as e:
logger.warning(f"Error ending session on close: {e}")
def _handle_gui_sessions_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
payload = json.loads(response_data) if response_data else []
sessions = [OpenGuiSessionInfo.model_validate(item) for item in payload]
self.gui_sessions_loaded.emit(sessions)
except Exception as e:
logger.error(f"Failed to load GUI sessions: {e}")
self.http_error.emit(str(e))
@Slot()
def load_gui_sessions(self):
if self.__base_url is None:
self.gui_sessions_loaded.emit([])
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/admin/gui_sessions"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.get(request)
reply.finished.connect(lambda: self._handle_gui_sessions_response(reply))
@Slot(int, int)
def request_gui_close(self, session_id: int, grace_seconds: int = 60):
if self.__base_url is None:
logger.info(f"POST /admin/gui_sessions/{session_id}/request_close?grace_seconds={grace_seconds}")
return
request = QNetworkRequest(
QUrl(f"{self.__base_url}/admin/gui_sessions/{session_id}/request_close?grace_seconds={grace_seconds}")
)
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_gui_session_mutation_response(reply))
@Slot(int)
def force_remove_gui_session(self, session_id: int):
if self.__base_url is None:
logger.info(f"DELETE /admin/gui_sessions/{session_id}")
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/admin/gui_sessions/{session_id}"))
request.setRawHeader(b"Authorization", f"Bearer {self.__token}".encode("utf-8"))
reply = self.__net_manager.deleteResource(request)
reply.finished.connect(lambda: self._handle_gui_session_mutation_response(reply))
def _handle_gui_session_mutation_response(self, reply: QNetworkReply):
try:
if reply.error() != QNetworkReply.NetworkError.NoError:
self.handle_req_response(reply)
self.load_gui_sessions()
return
_ = self.handle_response(reply)
self.load_gui_sessions()
except Exception as e:
logger.error(f"GUI session mutation failed: {e}")
self.http_error.emit(str(e))
self.load_gui_sessions()
@Slot(int)
def report_gui_interaction(self, session_id: int):
if self.__base_url is None:
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/admin/gui_sessions/{session_id}/interaction"))
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: reply.deleteLater())
def cleanup(self) -> None:
if getattr(self, "_cleanup_done", False):
return
self._cleanup_done = True
try:
if hasattr(self, "_baton_timeout_timer") and self._baton_timeout_timer is not None:
self._baton_timeout_timer.stop()
except Exception as e:
logger.warning(f"Failed to stop _baton_timeout_timer: {e}")
try:
if hasattr(self, "_DAQWorker__timer") and self.__timer is not None:
self.__timer.stop()
except Exception as e:
logger.warning(f"Failed to stop __timer: {e}")
for attr_name in (
"_baton_stream_reply",
"_face_detection_stream_reply",
"_automation_progress_stream_reply",
):
reply = getattr(self, attr_name, None)
if reply is None:
continue
try:
reply.abort()
except Exception as e:
logger.warning(f"Failed to abort {attr_name}: {e}")
try:
reply.deleteLater()
except Exception as e:
logger.warning(f"Failed to delete {attr_name}: {e}")
setattr(self, attr_name, None)
self._face_detection_stream_blocked_403 = False
self._automation_progress_stream_blocked_403 = False
self._automation_progress_buffer = ""