Files
AareDAQ/src/aare/gui/threads/daq_worker.py
T
duan_jandClaude Fable 5 8fed60e73b
CI / lint (push) Skipped
CI / test (3.11) (push) Skipped
CI / test (3.12) (push) Skipped
CI / test (3.13) (push) Skipped
CI / test-with-beamline-plugins (pxi_bec) (push) Skipped
CI / test-with-beamline-plugins (pxii_bec) (push) Skipped
CI / test-with-beamline-plugins (pxiii_bec) (push) Skipped
CI / test (3.11) (pull_request) Successful in 1m6s
CI / test (3.13) (pull_request) Successful in 1m11s
CI / test-with-beamline-plugins (pxii_bec) (pull_request) Successful in 1m12s
CI / test (3.12) (pull_request) Successful in 1m23s
CI / test-with-beamline-plugins (pxi_bec) (pull_request) Successful in 1m18s
CI / test-with-beamline-plugins (pxiii_bec) (pull_request) Successful in 1m35s
CI / test-with-coverage (pull_request) Successful in 1m43s
CI / coverage-analysis (pull_request) Failing after 3s
CI / lint (pull_request) Failing after 2m42s
docs: note the spreadsheet log-spam TODO in daq_worker
load_spreadsheet/load_reference_tools log a GET they never send every
~12.5s poll while base_url is None; fix later by demoting to debug or
logging once on the None->set edge.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-10 09:45:07 +02:00

2570 lines
102 KiB
Python

import copy
import json
import logging
import os
import random
import re
import time
from collections import deque
from dataclasses import dataclass
from datetime import datetime
from typing import ClassVar, Literal, cast
from aarecommon.config.logger import setup_logger
from aarecommon.errors.codes import AareErrorCode, AuthErrorCode, export_error_codes
from aarecommon.math.coordinate import AerotechCoordinate, SmargonCoordinate
from aarecommon.models.auth import BatonStatus
from aarecommon.models.automation import (
AutomationProgress,
LogEvent,
StepState,
StepStatus,
WorkflowStateKind,
)
from aarecommon.models.models import (
AutofocusSettings,
DAQStatusModel,
FluorescenceSpectrumOutputModel,
FluorescenceSpectrumParameterModel,
OpenGuiSessionInfo,
SampleCameraSettings,
SampleShortInfo,
SampleShortInfoList,
SimpleScanParameters,
)
from aarecommon.models.raster_grid import (
CompletedRasterGrid,
CompletedRasterGridElem,
RasterGridRequest,
)
from aarecommon.models.rotation_scan import CompletedRotationScan, RotationScanRequest
from aarecommon.recurrence_watcher import (
DEFAULT_WATCHERS,
WatcherTrip,
create_default_watchers,
load_watcher_threshold_overrides,
redis_key_to_env_var,
resolve_exception_class,
)
from jfjoch_client import ScanResult, ScanResultImagesInner
from PySide6.QtCore import QByteArray, QObject, QTimer, QUrl, Signal, Slot
from PySide6.QtNetwork import QNetworkAccessManager, QNetworkReply, QNetworkRequest, QSslError
from aare.gui.constants import LOGGER_NAME
logger = setup_logger(LOGGER_NAME)
SPREADHSEET_FREQUENCY = 25 # Every 5 seconds
# Operation errors that just mean "you can't do that right now" (e.g. the user
# triggered an action while the beamline/robot was busy). These are blocking but
# benign, so they are surfaced quietly (log + status bar) rather than as a modal
# pop-up, regardless of the exception's critical flag.
QUIET_OPERATION_ERROR_CODES = frozenset(
{
AareErrorCode.BEAMLINE_BUSY_EXCEPTION.value,
AareErrorCode.TELL_COMMAND_WHILE_BUSY_EXCEPTION.value,
}
)
@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_config_loaded = Signal(dict)
local_contact_config_saved = 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())
reply = self._net_manager.get(request)
reply.finished.connect(lambda: self.handle_status_response(reply))
def server_about_info(self) -> tuple[str, str]:
version_request = QNetworkRequest(QUrl(f"{self._base_url}/about/running_version"))
version_reply = self._net_manager.get(version_request)
file_request = QNetworkRequest(QUrl(f"{self._base_url}/about/server_file_path"))
file_reply = self._net_manager.get(file_request)
version = (
str(version_reply.readAll()) if version_reply.waitForReadyRead(500) else "Not connected"
)
file = str(file_reply.readAll()) if file_reply.waitForReadyRead(500) else "Not connected"
return version, file
_SELF_SIGNED_ERRORS: ClassVar[set[QSslError.SslError]] = {
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.exception("Exception reading status response")
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.exception("Exception from status response")
@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.exception("Exception from spreadsheet response")
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.exception("Exception from reference tools response")
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:
logger.debug("Could not read the raw reply body", exc_info=True)
try:
url = reply.request().url().toString()
except Exception:
logger.debug("Could not read the request URL", exc_info=True)
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)
elif error_info.code in QUIET_OPERATION_ERROR_CODES:
# Blocking-but-benign (e.g. "Beamline is busy"): inform quietly,
# no modal pop-up and no automation pause.
logger.info(f"Action unavailable: {error_info.message}")
self.status_message.emit(error_info.message, True)
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.exception("Sample resync failed")
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.exception("Hardware metadata resync failed")
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.exception("Recovery action failed")
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())
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())
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())
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("beamline/shutter?val=false")
@Slot()
def open_shutter(self):
self.generic_post("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())
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())
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())
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())
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 TypeError("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.exception("Exception from all_pgroups response")
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())
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.exception("Exception from rotation scan response")
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())
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.exception("Exception from raster scan response")
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())
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())
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:
# TODO: log spam. This (and load_reference_tools) fires every
# SPREADHSEET_FREQUENCY cycle (~12.5s) while base_url is None,
# logging a GET it never actually sends -> two INFO lines every
# poll. Fix by demoting to logger.debug, or log once on the
# None->set edge rather than on every poll.
logger.info("GET /sample/spreadsheet")
return
request = QNetworkRequest(QUrl(f"{self._base_url}/sample/spreadsheet"))
request.setRawHeader(b"Authorization", f"Bearer {self._token}".encode())
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("GET /sample/reference_tools")
return
request = QNetworkRequest(QUrl(f"{self._base_url}/sample/reference_tools"))
request.setRawHeader(b"Authorization", f"Bearer {self._token}".encode())
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:
logger.debug("Could not parse the error body as JSON", exc_info=True)
err_msg = raw_body
except Exception:
logger.debug("Could not extract error details from the reply body", exc_info=True)
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:
logger.debug("Could not classify the reply as an auth error", exc_info=True)
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:
logger.debug("Could not read the reply status code", exc_info=True)
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
return bool(
"daq state error" in text
or "must be idle to start measurement" in text
or "must be idle" in text
)
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())
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("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())
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())
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:
logger.debug("Could not read the request URL", exc_info=True)
url = "unknown-url"
return f"{context}\n\nURL: {url}\nHTTP status: {status}\nError: {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())
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 TypeError("Invalid local contact simulation state payload")
self.local_contact_simulation_state_loaded.emit(payload)
except Exception as e:
logger.warning("Local Contact simulation state request failed", exc_info=True)
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())
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 TypeError("Invalid local contact device state payload")
self.local_contact_device_state_loaded.emit(payload)
except Exception as e:
logger.warning("Local Contact device state request failed", exc_info=True)
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())
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 TypeError("Invalid local contact links payload")
self.local_contact_links_loaded.emit(payload)
except Exception as e:
logger.warning("Local Contact links request failed", exc_info=True)
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()
def load_local_contact_config(self):
if self._base_url is None:
self.local_contact_config_loaded.emit(
{
"mount_to_center_sleep_s": 0.0,
"line_scan_loop_face_y_padding_fraction_each_side": 0.5,
"line_scan_loop_all_y_padding_fraction_each_side": 0.5,
}
)
return
request = QNetworkRequest(QUrl(f"{self._base_url}/local_contact/config"))
request.setRawHeader(b"Authorization", f"Bearer {self._token}".encode())
reply = self._net_manager.get(request)
reply.finished.connect(lambda: self._handle_local_contact_config_response(reply))
def _handle_local_contact_config_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 TypeError("Invalid local contact config payload")
self.local_contact_config_loaded.emit(payload)
except Exception as e:
message = (
"Error transferring information from DAQ while loading Local Contact config.\n\n"
f"{e}"
)
logger.exception(message)
self.local_contact_transfer_error.emit(message)
@Slot(dict)
def set_local_contact_config(self, payload: dict):
if self._base_url is None:
self.local_contact_config_saved.emit(payload)
return
request = QNetworkRequest(QUrl(f"{self._base_url}/local_contact/config"))
request.setRawHeader(b"Authorization", f"Bearer {self._token}".encode())
request.setRawHeader(b"Content-Type", b"application/json")
body = QByteArray(json.dumps(payload).encode("utf-8"))
reply = self._net_manager.put(request, body)
reply.finished.connect(lambda: self._handle_set_local_contact_config_response(reply))
def _handle_set_local_contact_config_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 TypeError("Invalid local contact config response")
self.local_contact_config_saved.emit(payload)
self.local_contact_config_loaded.emit(payload)
self.status_message.emit("Local Contact config saved.", False)
except Exception as e:
message = (
f"Error transferring information from DAQ while saving Local Contact config.\n\n{e}"
)
logger.exception(message)
self.local_contact_transfer_error.emit(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())
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 TypeError("Invalid BEC user macros payload")
self.bec_user_macros_loaded.emit([str(item) for item in payload])
except Exception as e:
logger.exception("Failed to list BEC user macros")
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())
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 TypeError("Invalid BEC devices payload")
self.bec_devices_loaded.emit([str(item) for item in payload])
except Exception as e:
logger.exception("Failed to list BEC devices")
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 bec_save_current_bs_pos(self):
self.generic_post("bec/save_current_bs_pos")
@Slot()
def bec_save_current_collimator_pos(self):
self.generic_post("bec/save_current_collimator_pos")
@Slot()
def bec_save_current_aerotech_position(self):
self.generic_post("bec/save_current_aerotech_position")
@Slot()
def initialise_smargon(self):
self.generic_post("smargon/initialize")
@Slot()
def save_beam_location_camera_setting(self):
self.generic_post("beamline/save_beam_location_camera_setting")
@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("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.exception("Error in ml box")
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("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())
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.exception("Error in face detection")
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:
logger.exception("Face detection stream parse error")
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:
logger.debug("Could not read the face detection stream status", exc_info=True)
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:
logger.exception("Automation progress stream parse error")
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:
logger.debug("Could not read the automation progress stream status", exc_info=True)
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())
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())
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())
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())
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.exception("Fluorimeter data error")
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())
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.exception("Fluorimeter status error")
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.exception("Fluorimeter background error")
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())
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.exception("Exception from fluorimeter spectrum")
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())
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())
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:
logger.exception("SSE parse error")
@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}.{kk!s}"] = 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())
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())
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:
logger.debug("Could not read the request URL", exc_info=True)
url = ""
if int(status) == 404 and url.endswith("/meta/error-codes"):
reply.deleteLater()
logger.error("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 TypeError("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.exception("Failed to load error codes")
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())
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:
logger.exception("Baton stream parse error")
@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())
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.exception("Baton request failed")
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())
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.exception("Baton response failed")
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())
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.exception("Baton timeout check failed")
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}", exc_info=True)
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.exception("Failed to load GUI sessions")
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())
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())
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())
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.exception("GUI session mutation failed")
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())
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}", exc_info=True)
try:
if hasattr(self, "_timer") and self._timer is not None:
self._timer.stop()
except Exception as e:
logger.warning(f"Failed to stop __timer: {e}", exc_info=True)
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}", exc_info=True)
try:
reply.deleteLater()
except Exception as e:
logger.warning(f"Failed to delete {attr_name}: {e}", exc_info=True)
setattr(self, attr_name, None)
self._face_detection_stream_blocked_403 = False
self._automation_progress_stream_blocked_403 = False
self._automation_progress_buffer = ""