Files
AareDAQ/src/aare/gui/threads/daq_worker.py
T

1936 lines
77 KiB
Python

import copy
import random
import time
import json
from collections import deque
from PySide6.QtCore import Signal, QUrl, Slot, QTimer, QObject, QByteArray
from PySide6.QtNetwork import QNetworkAccessManager, QNetworkRequest, QNetworkReply
from jfjoch_client import ScanResult, ScanResultImagesInner
from aare.common.coordinate import SmargonCoordinate, Coordinate, AerotechCoordinate
from aare.common.auth_models import BatonStatus
from aare.common.coordinate import SmargonCoordinate, Coordinate
from aare.common.error_codes import export_error_codes, DAQErrorCode
from aare.common.exception_handler import JFJochCommunicationError
from aare.common.models import (
DAQStatusModel,
SampleShortInfoList,
SampleShortInfo,
SampleCameraSettings,
AutofocusSettings,
SimpleScanParameters,
FluorescenceSpectrumParameterModel,
FluorescenceSpectrumOutputModel,
OpenGuiSessionInfo,
)
from aare.common.automation_models import (
AutomationProgress,
StepState,
StepStatus,
WorkflowStateKind,
)
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
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)
#dedicated signals for polled device errors and request-time errors
polled_devices_status = Signal(str, bool) # (message, is_error)
detector_error = Signal(str, bool) # (message, is_error)
auth_error = Signal()
sample_missing = Signal(str)
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)
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.__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._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
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
# Deduplication tracking for polled device messages
self._last_polled_msg: str | None = None
self._last_polled_is_error: bool | None = None
# Deduplication tracking for detector/request-time messages
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
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}")
@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))
@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)
# 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:
pass
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 = reply.attribute(QNetworkRequest.Attribute.HttpStatusCodeAttribute)
err_details = reply.errorString()
raw_body = ""
body_json = None
try:
raw_body = reply.readAll().data().decode("utf-8")
if raw_body:
body_json = json.loads(raw_body)
if isinstance(body_json, dict):
if "detail" in body_json:
err_details = body_json["detail"]
elif "message" in body_json:
err_details = body_json["message"]
else:
err_details = raw_body
else:
err_details = raw_body
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 status == 401:
now = time.monotonic()
if now - self._last_auth_error_log_ts > self._auth_error_min_interval:
logger.error(f"{err_details}: 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(err_details)
else:
logger.error(f"{err_details}")
self.http_error.emit(err_details)
reply.deleteLater()
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_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 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 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.send_status_request()
@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 == "JFJOCH_UNAVAILABLE" 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 == "JFJOCH_UNAVAILABLE" 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 _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_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 == "JFJOCH_UNAVAILABLE" 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
@staticmethod
def _is_critical_automation_failure(
status: int | None,
body_json: dict | None = None,
err_msg: str | None = None,
) -> bool:
"""Detect server-side critical automation errors that require operator recovery."""
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")
if code == DAQErrorCode.AUTOMATION_CRITICAL.value:
return True
if status_int is not None and 500 <= status_int < 600 and code is None:
return True
text = str(err_msg or "").lower()
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.automated_scan_done.emit(sample_id, True, "")
else:
status, err_str, body_json = self._extract_reply_error_details(reply)
if self._is_detector_state_failure_message(err_str):
self._emit_detector_message(f"JFJoch: {err_str}", is_error=True)
if status == 401:
logger.error(f"Error in auto scan: {err_str}")
self.automated_scan_done.emit(sample_id, False, "Authentication Error")
self.auth_error.emit()
elif status == 404:
self.sample_missing.emit(err_str)
self.automated_scan_done.emit(sample_id, False, "Missing")
elif status == 410:
self.sample_missing.emit(err_str)
self.automated_scan_done.emit(sample_id, False, "Warning")
elif status == 417:
self.sample_missing.emit(err_str)
self.automated_scan_done.emit(sample_id, False, "Critical")
elif self._is_critical_automation_failure(status, body_json, err_str):
logger.critical(f"Critical automation failure: {err_str}")
self.http_error.emit(err_str)
self.automated_scan_done.emit(sample_id, False, "Critical")
self.automation_critical_failure.emit(err_str)
else:
logger.error(f"Error in auto scan: {err_str}")
self.http_error.emit(err_str)
self.automated_scan_done.emit(sample_id, False, err_str)
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 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()
def park_and_dry(self):
self.generic_post("sample/park_and_dry")
@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):
self._face_detection_stream_reply = None
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] = []
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"),
)
)
return AutomationProgress(
current_step=progress_payload.get("current_step"),
steps=steps,
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 ""),
)
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)
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):
self._automation_progress_stream_reply = None
self._automation_progress_buffer = ""
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
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
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"}, "DAQErrorCode": {...}}
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()
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()
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()
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._automation_progress_buffer = ""