From 8b9dcd8cef09e88115b65f80a91508d851c48642 Mon Sep 17 00:00:00 2001 From: GotthardG <51994228+GotthardG@users.noreply.github.com> Date: Fri, 11 Jul 2025 21:23:01 +0200 Subject: [PATCH] Update dependencies and enhance WebSocket processing Updated `aareDB` and `aaredaqlib` dependencies to the latest versions. Enhanced WebSocket connection handling with improved message processing, error handling, and environment variable validation. Introduced additional logging for debugging and streamlined Redis spreadsheet updates. --- common/pyproject.toml | 2 +- daq/pyproject.toml | 4 +- daq/src/aaredaq/spreadsheetupdater.py | 74 +++++++++++++++++++-------- 3 files changed, 56 insertions(+), 24 deletions(-) diff --git a/common/pyproject.toml b/common/pyproject.toml index ac16cff9..bc35aeff 100644 --- a/common/pyproject.toml +++ b/common/pyproject.toml @@ -8,7 +8,7 @@ dependencies = [ "pydantic==2.11.4", "numpy==2.2.5", "jfjoch_client==1.0.0rc61", - "aareDB==0.1.1a11" + "aareDB==0.1.1a12" ] [lint] diff --git a/daq/pyproject.toml b/daq/pyproject.toml index 7ec7f374..ef713a07 100644 --- a/daq/pyproject.toml +++ b/daq/pyproject.toml @@ -13,12 +13,12 @@ dependencies = [ "fastapi==0.115.13", "uvicorn==0.34.2", "ultralytics==8.3.133", - "aaredb==0.1.1a11", + "aaredb==0.1.1a12", "opencv-python-headless==4.11.0.86", "python_multipart==0.0.20", "websocket-client==1.8.0", "sseclient-py==1.8.0", - "aaredaqlib==0.2.3" + "aaredaqlib==0.2.3", ] [lint] diff --git a/daq/src/aaredaq/spreadsheetupdater.py b/daq/src/aaredaq/spreadsheetupdater.py index 4260c61a..57324495 100644 --- a/daq/src/aaredaq/spreadsheetupdater.py +++ b/daq/src/aaredaq/spreadsheetupdater.py @@ -3,68 +3,83 @@ import json import websocket import time import signal -from aareDBclient.models import PuckWithTellPosition - +from aareDBclient.models import PuckWithTellPosition, DataCollectionParameters from aaredaqlib.models import SampleShortInfoList, SampleShortInfo, DewarAddress -from config import BeamlineConfig # Your BeamlineConfig implementation -from aaredaqlib.beamline import MXBeamline # For beamline id/enum +from config import BeamlineConfig +from aaredaqlib.beamline import MXBeamline -SLOT_IDENTIFIER = "X06DA" # or whatever is relevant -WS_URL = f"wss://mx-db-01.psi.ch/dispatcher/protected_router/tell_runner_router/ws/samples-spreadsheet/{SLOT_IDENTIFIER}" -WS_HEADERS = [f"X-Shared-Password: {os.getenv('AAREDB_SHARED_PASSWORD')}"] +SLOT_IDENTIFIER = "X06DA" +WS_URL = f"wss://mx-db-01.psi.ch/dispatcher/protected_router/tell_runner/ws/samples-spreadsheet/{SLOT_IDENTIFIER}" -beamline = MXBeamline.X06DA # Use your beamline enum/value +# Ensure the environment variable for the shared password is set +password = os.getenv("AAREDB_SHARED_PASSWORD") +if not password: + raise ValueError("The AAREDB_SHARED_PASSWORD environment variable is not set.") +WS_HEADERS = [f"X-Shared-Password: {password}"] + +beamline = MXBeamline.X06DA config = BeamlineConfig(bl=beamline) -current_spreadsheet = None # Keep last spreadsheet cached for change detection +# Cache the current spreadsheet for change detection +current_spreadsheet = None def set_spreadsheet_in_redis(spreadsheet): """ Store the spreadsheet in Redis using your BeamlineConfig instance. - Could be JSON-serialized or fields set atomically as per your usual convention. """ + print('[REDIS][DEBUG] Data to write:', json.dumps(spreadsheet, indent=4)) # Pretty-print the data print('[REDIS][INFO] Writing spreadsheet to Redis...') config.__client.set(f"{config._BeamlineConfig__bl}:spreadsheet", json.dumps(spreadsheet)) def on_message(ws, message): + """ + Handle incoming WebSocket messages. + """ try: data = json.loads(message) - print("[WS][RAW] Message received:", message) - print("[WS][DEBUG] Incoming data:", data) - # Most messages will be dicts with a "samples" list + + # Process the pucks data from the "samples" key in the message if isinstance(data, dict) and "samples" in data: pucks_data = data["samples"] else: - pucks_data = data # fallback: it might be a list directly + pucks_data = data # fallback: might be a list directly + # Convert dictionary data into PuckWithTellPosition Pydantic model instances pucks = [PuckWithTellPosition(**item) for item in pucks_data] sample_short_infos = [] for p in pucks: + # Iterate through samples in the puck; `s` is a Pydantic `Sample` object for s in p.samples or []: + # Handle dewar_address dewar_address = None if p.tell_position is not None and len(p.tell_position) == 2: dewar_address = DewarAddress( segment=p.tell_position[0], pos=p.tell_position[1] ) + aaredb_params = s.data_collection_parameters + + # Build SampleShortInfo using attribute access. sample_short_infos.append( SampleShortInfo( db_id=s.id, puck_name=p.puck_name, dewar_name="Unknown" if p.dewar_name is None else p.dewar_name, sample_name=s.sample_name, - user=p.pgroup, - pin=s.position, + user=s.pgroup, + pin=s.position, # NOTE: This might need conversion if `pin` is an int and `position` is a string. location=dewar_address, - priority=1.0 if getattr(s, "priority", None) is None else s.priority, - comment=getattr(s, "comments", None), - mount_count=getattr(s, "mount_count", 0) or 0, - aaredb_params=getattr(s, "data_collection_parameters", None) + priority=s.priority if hasattr(s, "priority") and s.priority is not None else 1.0, + comment=s.comments if hasattr(s, "comments") else None, + mount_count=s.mount_count or 0, + aaredb_params=aaredb_params ) ) + + # Update Redis with the new spreadsheet ret = SampleShortInfoList(s=sample_short_infos) config.spreadsheet = ret print("DEBUG: Writing to redis key:", config._BeamlineConfig__bl + ":sample_spreadsheet") @@ -73,17 +88,34 @@ def on_message(ws, message): except Exception as exc: print("[WS][ERROR] Failed to parse or convert message:", exc) + def on_error(ws, error): + """ + Handle WebSocket errors. + """ print("[WS][ERROR]", error) + def on_close(ws, close_status_code, close_msg): + """ + Handle WebSocket closure. + """ print("[WS][CLOSE]", close_status_code, close_msg) + def on_open(ws): + """ + Handle WebSocket connection opening. + """ print("[WS][OPEN] WebSocket opened.") + def main(): - # Let's support clean exit on SIGINT/SIGTERM + """ + Main function to initiate WebSocket connection. + """ + + # Support clean exit on SIGINT/SIGTERM def exit_gracefully(*args): print('\n[EXIT] Caught signal, exiting.') exit(0)