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

636 lines
26 KiB
Python

import copy
import random
import time
import json
from PySide6.QtCore import Signal, QUrl, Slot, QTimer, QObject, QByteArray
from PySide6.QtNetwork import QNetworkAccessManager, QNetworkRequest, QNetworkReply
from jfjoch_client import ScanResult, ScanResultImagesInner
from aaredaqlib.coordinate import SmargonCoordinate, Coordinate
from aaredaqlib.models import DAQStatusModel, SampleShortInfoList, SampleShortInfo, SampleCameraSettings, \
AutofocusSettings, SimpleScanParameters, FluorescenceSpectrumParameterModel, FluorescenceSpectrumOutputModel
from aaredaqlib.raster_grid import RasterGridRequest, CompletedRasterGrid
from aaredaqlib.rotation_scan import RotationScanRequest, CompletedRotationScan
from aaredaqlib.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)
auth_error = Signal()
automated_scan_done = Signal(int, bool) # 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)
def __init__(self, base_url: str | None, token: str, parent=None):
super().__init__(parent)
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._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 = 10.0
self._last_status_error = None
self._last_smart_params_json: str | None = None
self._smart_params_debounce = QTimer(self)
self._smart_params_debounce.setInterval(500) # ms
self._smart_params_debounce.setSingleShot(True)
self._smart_params_pending: str | None = None
self._smart_params_debounce.timeout.connect(self._flush_smart_params)
@Slot()
def regular_update(self):
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):
if self.__base_url is None:
return
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):
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())
@Slot(QNetworkReply)
def handle_status_response(self, reply: QNetworkReply):
try:
response_data = self.handle_response(reply)
parsed_response = DAQStatusModel.model_validate_json(response_data)
self.update.emit(parsed_response)
except Exception as e:
now = time.monotonic()
if str(e) == str(self._last_status_error):
if now - self._last_auth_error_log_ts > self._auth_error_min_interval:
self._last_auth_error_log_ts = now
logger.error(f"Exception from status response: {e}")
else:
self._last_status_error = e
self._last_auth_error_log_ts = now
logger.error(f"Exception from status response: {e}")
self.http_error.emit(str(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_str = reply.errorString()
if status== 401:
now = time.monotonic()
if now - self._last_auth_error_log_ts > self._auth_error_min_interval:
logger.error(f"{err_str}: baton taken by another user")
self._last_auth_error_log_ts = now
self.auth_error.emit()
else:
logger.error(f"{err_str}")
self.http_error.emit(err_str)
reply.deleteLater()
def generic_post(self, url: str, body: str = ""):
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 = ""):
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):
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 zoom(self, f: float):
self.generic_put(f"beamline/zoom?val={f:.3f}")
@Slot(int)
def light(self, v: int):
self.generic_put(f"beamline/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 alc_background(self):
self.generic_post("alc/background")
@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(str)
def set_pgroup(self, val: str):
if val == "":
self.generic_delete("access/pgroup")
else:
self.generic_put(f"access/pgroup?val={val}")
@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:
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))
@Slot(RotationScanRequest)
def standard_scan(self, r: RotationScanRequest):
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:
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))
@Slot(RasterGridRequest)
def raster_scan(self, r: RasterGridRequest):
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),
mos = random.uniform(0, 0.1),
b= random.uniform(15.0, 80.0)
))
logger.debug("check that this works - raster scan - complete raster grid")
reply = CompletedRasterGrid(request = new_copy,
result = ScanResult(file_prefix=r.file_prefix, images=images))
logger.debug(f"It appears to work {reply}")
self.raster_scan_completed.emit(reply)
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/scan/raster?auto=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):
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),
mos=random.uniform(0, 0.1),
b=random.uniform(15.0, 80.0)
))
logger.debug("check that this works - raster scan auto - complete raster grid")
reply = CompletedRasterGrid(request=new_copy,
result=ScanResult(file_prefix=r.file_prefix, images=images))
logger.debug(f"It appears to work {reply}")
self.raster_scan_completed.emit(reply)
return
request = QNetworkRequest(QUrl(f"{self.__base_url}/scan/raster?auto=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 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:
if reply.attribute(QNetworkRequest.Attribute.HttpStatusCodeAttribute) == 401:
logger.error(f"Error in auto scan: {reply.errorString()}")
self.auth_error.emit()
else:
logger.error(f"Error in auto scan: {reply.errorString()}")
self.http_error.emit(reply.errorString())
self.automated_scan_done.emit(sample_id, False)
reply.deleteLater()
@Slot(SampleShortInfo)
def automated_scan(self, s: SampleShortInfo):
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))
def _flush_smart_params(self):
body = self._smart_params_pending
if body is None:
return
if body == self._last_smart_params_json:
return # no-op
self._last_smart_params_json = body
self.generic_post("scan/smart_params", body)
self._smart_params_pending = None
@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
body = p.model_dump_json()
# dedupe and debounce
self._smart_params_pending = body
self._smart_params_debounce.start()
@Slot(Coordinate)
def abr_tweak(self, c: Coordinate):
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 unmount(self):
self.generic_post("sample/unmount")
@Slot(SampleShortInfo)
def mount(self, s: SampleShortInfo, reference: bool = False):
self.generic_post(f"sample/mount?dbid={s.db_id}&reference={reference}")
@Slot(SampleShortInfo)
def sample_manual(self, s: SampleShortInfo):
self.generic_post(f"sample/manual", s.model_dump_json())
@Slot()
def cancel(self):
self.generic_post("scan/cancel")
def handle_ml_box_response(self, reply):
try:
response_data = self.handle_response(reply)
if response_data != "":
parsed_response = RasterGridRequest.model_validate_json(response_data)
self.raster_generated_by_ml.emit(parsed_response)
except Exception as e:
logger.error(f"Error in ml box: {e}")
self.http_error.emit(str(e))
@Slot()
def ml_bounding_box(self):
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()
@Slot()
def face_detection(self, steps:int, step_size:int):
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}")