From c6741ed9ce370bfe28099c231a06f7b88548b5d7 Mon Sep 17 00:00:00 2001 From: David Schote Date: Wed, 3 Sep 2025 18:46:25 +0200 Subject: [PATCH 01/57] Revision of orchestration engine --- .../src/scanhub_libraries/resources.py | 2 +- .../orchestrator/assets/__init__.py | 0 .../orchestrator/assets/acquisition.py | 34 ++++++++++ .../orchestrator/io_managers/__init__.py | 0 .../io_managers/idata_io_manager.py | 41 ++++++++++++ .../orchestrator/jobs/mrpro_reco_asset_job.py | 26 ++++++++ .../orchestrator/repository copy.py | 16 +++++ .../orchestrator/repository.py | 30 +++++++-- .../orchestrator/ressources/__init__.py | 0 .../orchestrator/ressources/datalake.py | 64 +++++++++++++++++++ .../app/api/dagster_queries.py | 11 ++-- .../app/api/manager_endpoints.py | 37 ++++++----- 12 files changed, 234 insertions(+), 27 deletions(-) create mode 100644 services/orchestration-engine/orchestrator/assets/__init__.py create mode 100644 services/orchestration-engine/orchestrator/assets/acquisition.py create mode 100644 services/orchestration-engine/orchestrator/io_managers/__init__.py create mode 100644 services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py create mode 100644 services/orchestration-engine/orchestrator/jobs/mrpro_reco_asset_job.py create mode 100644 services/orchestration-engine/orchestrator/repository copy.py create mode 100644 services/orchestration-engine/orchestrator/ressources/__init__.py create mode 100644 services/orchestration-engine/orchestrator/ressources/datalake.py diff --git a/services/base/shared_libs/src/scanhub_libraries/resources.py b/services/base/shared_libs/src/scanhub_libraries/resources.py index 38be40c9..0f4521ec 100644 --- a/services/base/shared_libs/src/scanhub_libraries/resources.py +++ b/services/base/shared_libs/src/scanhub_libraries/resources.py @@ -10,8 +10,8 @@ class JobConfigResource(ConfigurableResource): callback_url: str | None = None user_access_token: str | None = None + update_device_parameter_base_url: str | None = None input_files: list[str] output_dir: str task_id: str # serves as series id exam_id: str # serves as study id - update_device_parameter_base_url: str | None = None diff --git a/services/orchestration-engine/orchestrator/assets/__init__.py b/services/orchestration-engine/orchestrator/assets/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/services/orchestration-engine/orchestrator/assets/acquisition.py b/services/orchestration-engine/orchestrator/assets/acquisition.py new file mode 100644 index 00000000..ee8946ce --- /dev/null +++ b/services/orchestration-engine/orchestrator/assets/acquisition.py @@ -0,0 +1,34 @@ +"""Definition of acquisiton data assets.""" +from typing import Generator + +from dagster import DynamicOutput, MetadataValue, OpExecutionContext, asset +from scanhub_libraries.models import ResultOut + + +@asset( + required_resource_keys={"datalake"}, + config_schema={"acquisitions": list}, +) +def acquisition_results(context: OpExecutionContext) -> Generator[DynamicOutput]: + """Define acquisition result asset.""" + for acq in context.op_config["acquisitions"]: + result = ResultOut(**acq) + context.log.info("Processing acquisition result: %s", result) + mrd_file = context.resources.datalake.get_mrd_path(result.directory, result.files) + context.log.info("MRD file path: %s", str(mrd_file)) + device_parameter = context.resources.datalake.get_device_parameter(result.directory, result.files) + context.log.info("Device parameters: %s", str(device_parameter)) + + yield DynamicOutput( + value={ + "mrd_file": mrd_file, + "device_parameter": device_parameter, + "result": result.model_dump(), + }, + mapping_key=str(result.id), + metadata={ + "mrd_file": MetadataValue.path(str(mrd_file)), + "result": MetadataValue.json(result.model_dump()), + "device_parameter": MetadataValue.json(device_parameter), + }, + ) diff --git a/services/orchestration-engine/orchestrator/io_managers/__init__.py b/services/orchestration-engine/orchestrator/io_managers/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py b/services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py new file mode 100644 index 00000000..ebaa551f --- /dev/null +++ b/services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py @@ -0,0 +1,41 @@ +"""Dagster IO Manager for mrpro IData object.""" +from pathlib import Path + +from dagster import IOManager, OutputContext, InputContext, io_manager +from mrpro.data import IData + + +class IDataIOManager(IOManager): + """IO manager for mrpro IData object.""" + + def __init__(self, base_dir: str) -> None: + """Init.""" + self.base_dir = Path(base_dir) + + def handle_output(self, context: OutputContext, idata: IData) -> None: + """Write idata to dicom folder.""" + # Decide where to write based on asset key + asset_key = "_".join(context.asset_key.path) + output_dir = self.base_dir / asset_key + output_dir.mkdir(parents=True, exist_ok=False) + + idata.to_dicom_folder(output_dir) + + context.add_output_metadata({ + "stored_files": [f.name for f in output_dir.iterdir() if f.is_file()], + "output_dir": str(output_dir), + }) + + def load_input(self, context: InputContext) -> IData: + """Read idata from dicom folder.""" + # Reload previously written files if needed + upstream_metadata = context.upstream_output.metadata + dcm_folder = Path(upstream_metadata.get("output_dir")) + if not dcm_folder.exists(): + context.log.error("Dicom folder does not exist: %s", dcm_folder) + return IData.from_dicom_folder(dcm_folder) + + +@io_manager(config_schema={"base_dir": str}) +def idata_io_manager(init_context) -> IDataIOManager: + return IDataIOManager(base_dir=init_context.resource_config["base_dir"]) diff --git a/services/orchestration-engine/orchestrator/jobs/mrpro_reco_asset_job.py b/services/orchestration-engine/orchestrator/jobs/mrpro_reco_asset_job.py new file mode 100644 index 00000000..e8e818be --- /dev/null +++ b/services/orchestration-engine/orchestrator/jobs/mrpro_reco_asset_job.py @@ -0,0 +1,26 @@ +import mrpro +from dagster import asset, define_asset_job, AssetSelection + + +@asset(required_resource_keys={"datalake"}, io_manager_key="idata_io_manager") +def mrpro_direct_reconstructed_dcm(context, acquisition_results: list[dict]) -> mrpro.data.IData: + """Reconstruct image from a list acquisition results. + + 1. Loads acquisition results using the DataLakeRessource providing mrd path and meta data. + 2. Loads mrpro KData object from mrd path + 3. Performs image reconstruction using the direct reconstruction method from mrpro. + 4. Return the reconstructed IData object -> Return is passed to the idata_io_manager. + """ + mrd_input = acquisition_results[0]["mrd_file"] + context.log.info("Reading MRD input: %s", str(mrd_input)) + trajectory_calculator = mrpro.data.traj_calculators.KTrajectoryCartesian() + kdata = mrpro.data.KData.from_file(mrd_input, trajectory_calculator) + context.log.info("Loaded data: %s", kdata.shape) + reconstruction = mrpro.algorithms.reconstruction.DirectReconstruction(kdata) + context.log.info("Performing direct reconstruction using mrpro...") + return reconstruction(kdata) + +mrpro_direct_reconstruction_job = define_asset_job( + name="mrpro_reconstruction_job", + selection=AssetSelection.assets(mrpro_direct_reconstructed_dcm), +) diff --git a/services/orchestration-engine/orchestrator/repository copy.py b/services/orchestration-engine/orchestrator/repository copy.py new file mode 100644 index 00000000..616e2357 --- /dev/null +++ b/services/orchestration-engine/orchestrator/repository copy.py @@ -0,0 +1,16 @@ +"""Define dagster repository.""" +from dagster import Definitions, in_process_executor + +from orchestrator.jobs.frequency_calibration import frequency_calibration_job +from orchestrator.jobs.mrpro_image_reconstruction import mrpro_reconstruction_job + +defs = Definitions( + jobs=[ + mrpro_reconstruction_job, + frequency_calibration_job, + ], + executor=in_process_executor, # default executor for all jobs +) + + + diff --git a/services/orchestration-engine/orchestrator/repository.py b/services/orchestration-engine/orchestrator/repository.py index 8a49f734..e9ff5bbd 100644 --- a/services/orchestration-engine/orchestrator/repository.py +++ b/services/orchestration-engine/orchestrator/repository.py @@ -1,13 +1,33 @@ """Define dagster repository.""" +import os + from dagster import Definitions, in_process_executor -from orchestrator.jobs.frequency_calibration import frequency_calibration_job -from orchestrator.jobs.mrpro_image_reconstruction import mrpro_reconstruction_job +from orchestrator.assets.acquisition import acquisition_results +from orchestrator.io_managers.idata_io_manager import idata_io_manager +# from orchestrator.jobs.frequency_calibration import frequency_calibration_job +# from orchestrator.jobs.mrpro_image_reconstruction import mrpro_reconstruction_job +from orchestrator.jobs.mrpro_reco_asset_job import mrpro_direct_reconstructed_dcm, mrpro_direct_reconstruction_job +from orchestrator.ressources.datalake import DataLakeResource + +DATA_LAKE_DIR = os.getenv("DATA_LAKE_DIRECTORY", "data") +# if DATA_LAKE_DIR is None: # ensure that DATA_LAKE_DIR is set +# raise OSError("Missing `DATA_LAKE_DIRECTORY` environment variable.") -defs = Definitions( + +definitions = Definitions( + assets=[ + acquisition_results, + mrpro_direct_reconstructed_dcm, + ], + resources={ + "datalake": DataLakeResource(), + "idata_io_manager": idata_io_manager.configured({"base_dir": DATA_LAKE_DIR}), + }, jobs=[ - mrpro_reconstruction_job, - frequency_calibration_job, + mrpro_direct_reconstruction_job, + # mrpro_reconstruction_job, + # frequency_calibration_job, ], executor=in_process_executor, # default executor for all jobs ) diff --git a/services/orchestration-engine/orchestrator/ressources/__init__.py b/services/orchestration-engine/orchestrator/ressources/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/services/orchestration-engine/orchestrator/ressources/datalake.py b/services/orchestration-engine/orchestrator/ressources/datalake.py new file mode 100644 index 00000000..8bf5ca4d --- /dev/null +++ b/services/orchestration-engine/orchestrator/ressources/datalake.py @@ -0,0 +1,64 @@ +"""Definition of dagster data lake ressource for acquisition data.""" +import json +from pathlib import Path + +from dagster import ConfigurableResource + + +class DataLakeResource(ConfigurableResource): + """Dagster data lake ressource.""" + + def get_mrd_path(self, directory: str, filenames: list[str]) -> Path: + """Construct and validate the MRD path inside the data lake. + + Parameters + ---------- + directory : str + Absolute path to result data (contains data lake path). + filenames: list[str] + List of filenames which can be found in directory. + + Returns + ------- + path + Path to acquisition ISMRMRD file. + + """ + filename = next((f for f in filenames if f.lower().endswith(".mrd")), None) + if filename is None: + raise FileNotFoundError(f"Acquisition result does not specify mrd filename.") + mrd_path = Path(directory) / filename + if not mrd_path.is_file(): + raise FileNotFoundError(f"MRD file does not exist: {mrd_path}") + return mrd_path + + def get_device_parameter(self, directory: str, filenames: list[str]) -> dict: + """Return the path to the device parameter JSON file if it exists. + + Parameters + ---------- + directory : str + Absolute path to result data (contains data lake path). + filenames: list[str] + List of filenames which can be found in directory. + + Returns + ------- + dict + Dictionary containing device parameters + + """ + filename = next((f for f in filenames if f.lower().endswith(".json")), None) + if filename is None: + raise FileNotFoundError(f"Acquisition result does not specify device parameter file.") + json_path = Path(directory) / filename + # Check if parameter file exists + if not json_path.exists(): + raise FileExistsError(f"Device parameter file does not exist: {json_path}") + # Load parameter file + with json_path.open("r") as fh: + data = json.load(fh) + # Check if parameter file contains device id and parameter + if "device_id" not in data or "parameter" not in data: + raise AttributeError(f"Invalid paraeter file: {json_path}") + return data diff --git a/services/workflow-manager/app/api/dagster_queries.py b/services/workflow-manager/app/api/dagster_queries.py index 71cada9c..b9d94ad7 100644 --- a/services/workflow-manager/app/api/dagster_queries.py +++ b/services/workflow-manager/app/api/dagster_queries.py @@ -12,12 +12,8 @@ def list_dagster_jobs() -> list[dict]: ... on RepositoryConnection { nodes { name - location { - name - } - pipelines { - name - } + location { name } + jobs { name } } } } @@ -26,12 +22,15 @@ def list_dagster_jobs() -> list[dict]: response = requests.post(DAGSTER_URL, json={"query": query}, timeout=3) response.raise_for_status() data = response.json() + print("Received dagster jobs: ", data) jobs = [] for repo in data["data"]["repositoriesOrError"]["nodes"]: repo_name = repo["name"] location_name = repo["location"]["name"] for pipeline in repo["pipelines"]: job_name = pipeline["name"] + if job_name == "__ASSET_JOB": + continue job_id = f"{location_name}::{repo_name}::{job_name}" jobs.append({ "job_id": job_id, diff --git a/services/workflow-manager/app/api/manager_endpoints.py b/services/workflow-manager/app/api/manager_endpoints.py index 96910d5f..6720fbda 100644 --- a/services/workflow-manager/app/api/manager_endpoints.py +++ b/services/workflow-manager/app/api/manager_endpoints.py @@ -201,10 +201,11 @@ def handle_dag_task_trigger( detail = f"No results found for input task {input_task.id}." raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=detail) latest_result = sorted(input_task.results, key=lambda r: r.datetime_created, reverse=True)[0] - for _file in latest_result.files: - file_path = Path(latest_result.directory) / _file - if file_path.exists(): - job_inputs.append(str(file_path)) + job_inputs.append(latest_result.model_dump()) + # for _file in latest_result.files: + # file_path = Path(latest_result.directory) / _file + # if file_path.exists(): + # job_inputs.append(str(file_path)) # Update the task status to IN_PROGRESS task.status = ItemStatus.INPROGRESS @@ -216,7 +217,7 @@ def handle_dag_task_trigger( # Use internal url and http (not https and port 8443) because callback endpoint is requested from another docker container callback_endpoint = f"{WORKFLOW_MANAGER_URI}/result_ready/{task.id}/{new_result_out.id}" device_parameter_update_endpoint = f"{DEVICE_MANAGER_URI}/parameter/" - result_directory = f"/data/{str(task.workflow_id)}/{str(task.id)}/{str(new_result_out.id)}/" + result_directory = f"/{DATA_LAKE_DIR}/{str(task.workflow_id)}/{str(task.id)}/{str(new_result_out.id)}/" # Trigger dagster job job_name, repository, location = parse_job_id(task.dag_id) @@ -225,18 +226,24 @@ def handle_dag_task_trigger( job_name=job_name, repository_location_name=location, repository_name=repository, - run_config=RunConfig(resources={ - SCANHUB_RESOURCE_KEY: JobConfigResource( - callback_url=callback_endpoint, - user_access_token=access_token, - input_files=job_inputs, - output_dir=result_directory, - task_id=task_id, - exam_id=exam_id, - update_device_parameter_base_url=device_parameter_update_endpoint, - ), + # run_config=RunConfig(resources={ + # SCANHUB_RESOURCE_KEY: JobConfigResource( + # callback_url=callback_endpoint, + # user_access_token=access_token, + # input_files=job_inputs, + # output_dir=result_directory, + # task_id=task_id, + # exam_id=exam_id, + # update_device_parameter_base_url=device_parameter_update_endpoint, + # ), + # }), + run_config=RunConfig(ops={ + "acquisition_results": { + "config": { "acquisitions": job_inputs }, + }, }), ) + if run_id: # Update result new_result_out = set_result( From 7f21a83865b0481e4fd4674541bd6404747d1dda Mon Sep 17 00:00:00 2001 From: David Schote Date: Thu, 4 Sep 2025 13:53:28 +0200 Subject: [PATCH 02/57] Fixed asset based dagster job definition --- .../src/scanhub_libraries/resources.py | 14 +- .../orchestrator/assets/acquisition.py | 64 ++++---- .../mrpro_direct_reconstruction.py} | 27 ++-- .../orchestrator/hooks/__init__.py | 0 .../orchestrator/hooks/scanhub.py | 29 ++++ .../io_managers/idata_io_manager.py | 58 ++++--- .../orchestrator/repository.py | 48 +++--- .../orchestrator/ressources/datalake.py | 16 +- .../orchestrator/ressources/notifier.py | 40 +++++ .../app/api/dagster_queries.py | 3 +- .../app/api/manager_endpoints.py | 149 ++++++++++-------- 11 files changed, 285 insertions(+), 163 deletions(-) rename services/orchestration-engine/orchestrator/{jobs/mrpro_reco_asset_job.py => assets/mrpro_direct_reconstruction.py} (52%) create mode 100644 services/orchestration-engine/orchestrator/hooks/__init__.py create mode 100644 services/orchestration-engine/orchestrator/hooks/scanhub.py create mode 100644 services/orchestration-engine/orchestrator/ressources/notifier.py diff --git a/services/base/shared_libs/src/scanhub_libraries/resources.py b/services/base/shared_libs/src/scanhub_libraries/resources.py index 0f4521ec..354b8976 100644 --- a/services/base/shared_libs/src/scanhub_libraries/resources.py +++ b/services/base/shared_libs/src/scanhub_libraries/resources.py @@ -2,7 +2,10 @@ from dagster import ConfigurableResource -SCANHUB_RESOURCE_KEY = "job_config" +DAG_CONFIG_KEY = "dag_config" +DATA_LAKE_KEY = "data_lake" +IDATA_IO_KEY = "idata_io_manager" +NOTIFIER_KEY = "scanhub_notifier" class JobConfigResource(ConfigurableResource): @@ -15,3 +18,12 @@ class JobConfigResource(ConfigurableResource): output_dir: str task_id: str # serves as series id exam_id: str # serves as study id + + +class DAGConfiguration(ConfigurableResource): + """Run-scoped parameters accessible from assets + IO managers.""" + + output_directory: str = "" + input_files: list[str] = [] + user_access_token: str = "" + output_result_id: str = "" diff --git a/services/orchestration-engine/orchestrator/assets/acquisition.py b/services/orchestration-engine/orchestrator/assets/acquisition.py index ee8946ce..47dfb116 100644 --- a/services/orchestration-engine/orchestrator/assets/acquisition.py +++ b/services/orchestration-engine/orchestrator/assets/acquisition.py @@ -1,34 +1,44 @@ """Definition of acquisiton data assets.""" -from typing import Generator +from dataclasses import dataclass +from pathlib import Path -from dagster import DynamicOutput, MetadataValue, OpExecutionContext, asset -from scanhub_libraries.models import ResultOut +from dagster import AssetExecutionContext, MetadataValue, asset +from scanhub_libraries.resources import DAGConfiguration + +from orchestrator.ressources.datalake import DataLakeResource + + +@dataclass +class AcquisitionData: + """Acquisition data output of read acquisition data asset.""" + + mrd_path: Path + device_parameter: dict @asset( - required_resource_keys={"datalake"}, - config_schema={"acquisitions": list}, + group_name="io", + description="Provides acquired ISMRMRD result", ) -def acquisition_results(context: OpExecutionContext) -> Generator[DynamicOutput]: +def read_acquisition_data( + context: AssetExecutionContext, + dag_config: DAGConfiguration, + data_lake: DataLakeResource, +) -> AcquisitionData: """Define acquisition result asset.""" - for acq in context.op_config["acquisitions"]: - result = ResultOut(**acq) - context.log.info("Processing acquisition result: %s", result) - mrd_file = context.resources.datalake.get_mrd_path(result.directory, result.files) - context.log.info("MRD file path: %s", str(mrd_file)) - device_parameter = context.resources.datalake.get_device_parameter(result.directory, result.files) - context.log.info("Device parameters: %s", str(device_parameter)) - - yield DynamicOutput( - value={ - "mrd_file": mrd_file, - "device_parameter": device_parameter, - "result": result.model_dump(), - }, - mapping_key=str(result.id), - metadata={ - "mrd_file": MetadataValue.path(str(mrd_file)), - "result": MetadataValue.json(result.model_dump()), - "device_parameter": MetadataValue.json(device_parameter), - }, - ) + mrd_file = data_lake.get_mrd_path(dag_config.input_files) + context.log.info("MRD file path: %s", str(mrd_file)) + device_parameter = data_lake.get_device_parameter(dag_config.input_files) + context.log.info("Device parameters: %s", str(device_parameter)) + + # Optional: surface helpful metadata in the Dagster UI + context.add_output_metadata({ + "mrd_path": MetadataValue.path(str(mrd_file)), + "device_parameter": MetadataValue.json(device_parameter), + "num_input_files": len(dag_config.input_files), + "output_directory": MetadataValue.path(dag_config.output_directory), + "output_result_id": dag_config.output_result_id, + "access_token": dag_config.user_access_token, + }) + + return AcquisitionData(mrd_path=mrd_file, device_parameter=device_parameter) diff --git a/services/orchestration-engine/orchestrator/jobs/mrpro_reco_asset_job.py b/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py similarity index 52% rename from services/orchestration-engine/orchestrator/jobs/mrpro_reco_asset_job.py rename to services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py index e8e818be..055547a0 100644 --- a/services/orchestration-engine/orchestrator/jobs/mrpro_reco_asset_job.py +++ b/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py @@ -1,9 +1,20 @@ import mrpro -from dagster import asset, define_asset_job, AssetSelection +from dagster import AssetIn, AssetKey, asset +from scanhub_libraries.resources import DAGConfiguration +from orchestrator.assets.acquisition import AcquisitionData +from orchestrator.hooks.scanhub import notify_dag_success +from orchestrator.io_managers.idata_io_manager import IDataContext -@asset(required_resource_keys={"datalake"}, io_manager_key="idata_io_manager") -def mrpro_direct_reconstructed_dcm(context, acquisition_results: list[dict]) -> mrpro.data.IData: + +@asset( + group_name="reconstruction", + description="MRpro direct reconstruction.", + ins={"data": AssetIn(key=AssetKey("read_acquisition_data"))}, + io_manager_key="idata_io_manager", + hooks={notify_dag_success}, +) +def mrpro_direct_reconstruction(context, data: AcquisitionData, dag_config: DAGConfiguration) -> IDataContext: """Reconstruct image from a list acquisition results. 1. Loads acquisition results using the DataLakeRessource providing mrd path and meta data. @@ -11,16 +22,12 @@ def mrpro_direct_reconstructed_dcm(context, acquisition_results: list[dict]) -> 3. Performs image reconstruction using the direct reconstruction method from mrpro. 4. Return the reconstructed IData object -> Return is passed to the idata_io_manager. """ - mrd_input = acquisition_results[0]["mrd_file"] + mrd_input = data.mrd_path context.log.info("Reading MRD input: %s", str(mrd_input)) trajectory_calculator = mrpro.data.traj_calculators.KTrajectoryCartesian() kdata = mrpro.data.KData.from_file(mrd_input, trajectory_calculator) context.log.info("Loaded data: %s", kdata.shape) reconstruction = mrpro.algorithms.reconstruction.DirectReconstruction(kdata) context.log.info("Performing direct reconstruction using mrpro...") - return reconstruction(kdata) - -mrpro_direct_reconstruction_job = define_asset_job( - name="mrpro_reconstruction_job", - selection=AssetSelection.assets(mrpro_direct_reconstructed_dcm), -) + idata = reconstruction(kdata) + return IDataContext(data=idata, dag_config=dag_config) diff --git a/services/orchestration-engine/orchestrator/hooks/__init__.py b/services/orchestration-engine/orchestrator/hooks/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/services/orchestration-engine/orchestrator/hooks/scanhub.py b/services/orchestration-engine/orchestrator/hooks/scanhub.py new file mode 100644 index 00000000..8757e4b9 --- /dev/null +++ b/services/orchestration-engine/orchestrator/hooks/scanhub.py @@ -0,0 +1,29 @@ +# orchestrator/hooks/recon_hooks.py +from pathlib import Path + +from dagster import HookContext, success_hook +from scanhub_libraries.resources import NOTIFIER_KEY + + +@success_hook(required_resource_keys={NOTIFIER_KEY}) +def notify_dag_success(context: HookContext) -> None: + """Check if data has been written and notify backend.""" + # Expect the asset to return an IDataContext + result = next(iter(context.op_output_values.values()), None) + if result is None: + context.log.warning("NOTIFY-HOOK: No return value; skipping.") + return + dag_cfg = getattr(result, "dag_config", None) + out_dir = getattr(dag_cfg, "output_directory", None) + if not out_dir: + context.log.warning("NOTIFY-HOOK: No output_directory; skipping.") + return + + p = Path(out_dir) + files = [str(f) for f in p.glob("**/*") if f.is_file()] + if not files: + context.log.warning(f"NOTIFY-HOOK: No files in {p}; skipping.") + return + + notifier = getattr(context.resources, NOTIFIER_KEY) + notifier.send_dag_success(success=True) diff --git a/services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py b/services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py index ebaa551f..b9ac53f1 100644 --- a/services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py +++ b/services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py @@ -1,41 +1,51 @@ """Dagster IO Manager for mrpro IData object.""" +from dataclasses import dataclass from pathlib import Path -from dagster import IOManager, OutputContext, InputContext, io_manager +from dagster import ConfigurableIOManager, InputContext, OutputContext from mrpro.data import IData +from scanhub_libraries.resources import DAGConfiguration -class IDataIOManager(IOManager): - """IO manager for mrpro IData object.""" +@dataclass +class IDataContext: + """Context definition for IData IO manager.""" + + data: IData + dag_config: DAGConfiguration + - def __init__(self, base_dir: str) -> None: - """Init.""" - self.base_dir = Path(base_dir) +class IDataIOManager(ConfigurableIOManager): + """IO manager for mrpro IData object.""" - def handle_output(self, context: OutputContext, idata: IData) -> None: + def handle_output(self, context: OutputContext, obj: IDataContext) -> None: """Write idata to dicom folder.""" # Decide where to write based on asset key - asset_key = "_".join(context.asset_key.path) - output_dir = self.base_dir / asset_key - output_dir.mkdir(parents=True, exist_ok=False) + if not obj.dag_config.output_directory: + context.log.error("Output directory not defined") + raise AttributeError + directory_path = Path(obj.dag_config.output_directory) - idata.to_dicom_folder(output_dir) + obj.data.to_dicom_folder(directory_path) + # Surface paths in the UI and for hooks/sensors context.add_output_metadata({ - "stored_files": [f.name for f in output_dir.iterdir() if f.is_file()], - "output_dir": str(output_dir), + "output_directory": str(directory_path.resolve()), + "stored_files": [p.name for p in directory_path.iterdir() if p.is_file()], }) def load_input(self, context: InputContext) -> IData: """Read idata from dicom folder.""" - # Reload previously written files if needed - upstream_metadata = context.upstream_output.metadata - dcm_folder = Path(upstream_metadata.get("output_dir")) - if not dcm_folder.exists(): - context.log.error("Dicom folder does not exist: %s", dcm_folder) - return IData.from_dicom_folder(dcm_folder) - - -@io_manager(config_schema={"base_dir": str}) -def idata_io_manager(init_context) -> IDataIOManager: - return IDataIOManager(base_dir=init_context.resource_config["base_dir"]) + if (meta := context.metadata) is None: + context.log.error("No metadata, cannot save result") + raise AttributeError + dcm_folder = meta.get("output_directory") + if not dcm_folder: + context.log.error("Missing directory to load dicom files") + raise AttributeError + dcm_path = Path(dcm_folder) + if not dcm_path.exists(): + msg = f"DICOM folder does not exist: {dcm_path}" + context.log.error(msg) + raise FileNotFoundError(msg) + return IData.from_dicom_folder(str(dcm_path)) diff --git a/services/orchestration-engine/orchestrator/repository.py b/services/orchestration-engine/orchestrator/repository.py index e9ff5bbd..40cd7f27 100644 --- a/services/orchestration-engine/orchestrator/repository.py +++ b/services/orchestration-engine/orchestrator/repository.py @@ -1,33 +1,41 @@ """Define dagster repository.""" import os -from dagster import Definitions, in_process_executor +from dagster import AssetSelection, Definitions, define_asset_job, in_process_executor -from orchestrator.assets.acquisition import acquisition_results -from orchestrator.io_managers.idata_io_manager import idata_io_manager # from orchestrator.jobs.frequency_calibration import frequency_calibration_job # from orchestrator.jobs.mrpro_image_reconstruction import mrpro_reconstruction_job -from orchestrator.jobs.mrpro_reco_asset_job import mrpro_direct_reconstructed_dcm, mrpro_direct_reconstruction_job +from orchestrator.assets.acquisition import read_acquisition_data +from orchestrator.assets.mrpro_direct_reconstruction import mrpro_direct_reconstruction +from orchestrator.io_managers.idata_io_manager import IDataIOManager from orchestrator.ressources.datalake import DataLakeResource +from orchestrator.ressources.notifier import BackendNotifier +from scanhub_libraries.resources import DAGConfiguration, DAG_CONFIG_KEY, IDATA_IO_KEY, DATA_LAKE_KEY, NOTIFIER_KEY DATA_LAKE_DIR = os.getenv("DATA_LAKE_DIRECTORY", "data") # if DATA_LAKE_DIR is None: # ensure that DATA_LAKE_DIR is set # raise OSError("Missing `DATA_LAKE_DIRECTORY` environment variable.") -definitions = Definitions( - assets=[ - acquisition_results, - mrpro_direct_reconstructed_dcm, - ], - resources={ - "datalake": DataLakeResource(), - "idata_io_manager": idata_io_manager.configured({"base_dir": DATA_LAKE_DIR}), - }, - jobs=[ - mrpro_direct_reconstruction_job, - # mrpro_reconstruction_job, - # frequency_calibration_job, - ], - executor=in_process_executor, # default executor for all jobs -) +assets = [ + read_acquisition_data, + mrpro_direct_reconstruction, +] + +ressources = { + DAG_CONFIG_KEY: DAGConfiguration(), + DATA_LAKE_KEY: DataLakeResource(), + IDATA_IO_KEY: IDataIOManager(), + NOTIFIER_KEY: BackendNotifier(), +} + +jobs = [ + define_asset_job( + name="mrpro_reconstruction_job", + selection=AssetSelection.keys("read_acquisition_data", "mrpro_direct_reconstruction"), + # selection=[dg.AssetSelection.keys("reconstruct_numpy")], + # tags={"workflow_id": "numpy"}, + ), +] + +defs = Definitions(assets=assets, resources=ressources, jobs=jobs, executor=in_process_executor) diff --git a/services/orchestration-engine/orchestrator/ressources/datalake.py b/services/orchestration-engine/orchestrator/ressources/datalake.py index 8bf5ca4d..5fbe538e 100644 --- a/services/orchestration-engine/orchestrator/ressources/datalake.py +++ b/services/orchestration-engine/orchestrator/ressources/datalake.py @@ -8,7 +8,7 @@ class DataLakeResource(ConfigurableResource): """Dagster data lake ressource.""" - def get_mrd_path(self, directory: str, filenames: list[str]) -> Path: + def get_mrd_path(self, files: list[str]) -> Path: """Construct and validate the MRD path inside the data lake. Parameters @@ -24,15 +24,14 @@ def get_mrd_path(self, directory: str, filenames: list[str]) -> Path: Path to acquisition ISMRMRD file. """ - filename = next((f for f in filenames if f.lower().endswith(".mrd")), None) + filename = next((f for f in files if f.lower().endswith(".mrd")), None) if filename is None: raise FileNotFoundError(f"Acquisition result does not specify mrd filename.") - mrd_path = Path(directory) / filename - if not mrd_path.is_file(): + if not (mrd_path := Path(filename)).is_file(): raise FileNotFoundError(f"MRD file does not exist: {mrd_path}") return mrd_path - def get_device_parameter(self, directory: str, filenames: list[str]) -> dict: + def get_device_parameter(self, files: list[str]) -> dict: """Return the path to the device parameter JSON file if it exists. Parameters @@ -48,12 +47,11 @@ def get_device_parameter(self, directory: str, filenames: list[str]) -> dict: Dictionary containing device parameters """ - filename = next((f for f in filenames if f.lower().endswith(".json")), None) - if filename is None: + json_file = next((f for f in files if f.lower().endswith(".json")), None) + if json_file is None: raise FileNotFoundError(f"Acquisition result does not specify device parameter file.") - json_path = Path(directory) / filename # Check if parameter file exists - if not json_path.exists(): + if not (json_path := Path(json_file)).exists(): raise FileExistsError(f"Device parameter file does not exist: {json_path}") # Load parameter file with json_path.open("r") as fh: diff --git a/services/orchestration-engine/orchestrator/ressources/notifier.py b/services/orchestration-engine/orchestrator/ressources/notifier.py new file mode 100644 index 00000000..82fc76c1 --- /dev/null +++ b/services/orchestration-engine/orchestrator/ressources/notifier.py @@ -0,0 +1,40 @@ +# orchestrator/resources/notifier.py +import json + +import httpx +from dagster import ConfigurableResource +from fastapi.encoders import jsonable_encoder + + +class BackendNotifier(ConfigurableResource): + """Backend notifier.""" + + success_callback_url: str | None = None + devicemanager_url: str | None = None + access_token: str | None = None + timeout: float = 5.0 + + def send_dag_success(self, success: bool) -> None: + """Notify backend about successful execution of dagster job/dag.""" + if self.success_callback_url is None: + raise AttributeError + if self.access_token is None: + raise AttributeError + headers = {"Authorization": "Bearer " + self.access_token} + payload = {"success": success} + with httpx.Client(timeout=self.timeout) as client: + response = client.post(self.success_callback_url, json=payload, headers=headers) + response.raise_for_status() + + def send_device_parameter_update(self, device_id: str, parameter: dict) -> None: + """Notify backend about device parameter update and send parameters.""" + if self.devicemanager_url is None: + raise AttributeError + if self.access_token is None: + raise AttributeError + headers = {"Authorization": "Bearer " + self.access_token} + url = self.devicemanager_url.rstrip("/") + f"/{device_id}" + payload = json.dumps(parameter, default=jsonable_encoder) + with httpx.Client(timeout=self.timeout) as client: + response = client.put(url, json=payload, headers=headers) + response.raise_for_status() diff --git a/services/workflow-manager/app/api/dagster_queries.py b/services/workflow-manager/app/api/dagster_queries.py index b9d94ad7..ae3488af 100644 --- a/services/workflow-manager/app/api/dagster_queries.py +++ b/services/workflow-manager/app/api/dagster_queries.py @@ -22,12 +22,11 @@ def list_dagster_jobs() -> list[dict]: response = requests.post(DAGSTER_URL, json={"query": query}, timeout=3) response.raise_for_status() data = response.json() - print("Received dagster jobs: ", data) jobs = [] for repo in data["data"]["repositoriesOrError"]["nodes"]: repo_name = repo["name"] location_name = repo["location"]["name"] - for pipeline in repo["pipelines"]: + for pipeline in repo["jobs"]: job_name = pipeline["name"] if job_name == "__ASSET_JOB": continue diff --git a/services/workflow-manager/app/api/manager_endpoints.py b/services/workflow-manager/app/api/manager_endpoints.py index 6720fbda..b3901edd 100644 --- a/services/workflow-manager/app/api/manager_endpoints.py +++ b/services/workflow-manager/app/api/manager_endpoints.py @@ -19,8 +19,8 @@ import requests from dagster import RunConfig -from dagster_graphql import DagsterGraphQLClient -from fastapi import APIRouter, Depends, HTTPException, status +from dagster_graphql import DagsterGraphQLClient, DagsterGraphQLClientError +from fastapi import APIRouter, Depends, HTTPException, status, Body from fastapi.encoders import jsonable_encoder from fastapi.security import OAuth2PasswordBearer from scanhub_libraries.models import ( @@ -36,7 +36,7 @@ TaskType, WorkflowOut, ) -from scanhub_libraries.resources import SCANHUB_RESOURCE_KEY, JobConfigResource +from scanhub_libraries.resources import NOTIFIER_KEY, DAG_CONFIG_KEY, DAGConfiguration from scanhub_libraries.security import get_current_user from scanhub_libraries.utils import calc_age_from_date @@ -201,11 +201,10 @@ def handle_dag_task_trigger( detail = f"No results found for input task {input_task.id}." raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail=detail) latest_result = sorted(input_task.results, key=lambda r: r.datetime_created, reverse=True)[0] - job_inputs.append(latest_result.model_dump()) - # for _file in latest_result.files: - # file_path = Path(latest_result.directory) / _file - # if file_path.exists(): - # job_inputs.append(str(file_path)) + for _file in latest_result.files: + file_path = Path(latest_result.directory) / _file + if file_path.exists(): + job_inputs.append(str(file_path)) # Update the task status to IN_PROGRESS task.status = ItemStatus.INPROGRESS @@ -217,32 +216,46 @@ def handle_dag_task_trigger( # Use internal url and http (not https and port 8443) because callback endpoint is requested from another docker container callback_endpoint = f"{WORKFLOW_MANAGER_URI}/result_ready/{task.id}/{new_result_out.id}" device_parameter_update_endpoint = f"{DEVICE_MANAGER_URI}/parameter/" - result_directory = f"/{DATA_LAKE_DIR}/{str(task.workflow_id)}/{str(task.id)}/{str(new_result_out.id)}/" + result_directory = f"{DATA_LAKE_DIR}/{str(task.workflow_id)}/{str(task.id)}/{str(new_result_out.id)}/" # Trigger dagster job job_name, repository, location = parse_job_id(task.dag_id) print(f"Triggering job: {job_name} in repository: {repository} at location: {location}") - run_id = dg_client.submit_job_execution( - job_name=job_name, - repository_location_name=location, - repository_name=repository, - # run_config=RunConfig(resources={ - # SCANHUB_RESOURCE_KEY: JobConfigResource( - # callback_url=callback_endpoint, - # user_access_token=access_token, - # input_files=job_inputs, - # output_dir=result_directory, - # task_id=task_id, - # exam_id=exam_id, - # update_device_parameter_base_url=device_parameter_update_endpoint, - # ), - # }), - run_config=RunConfig(ops={ - "acquisition_results": { - "config": { "acquisitions": job_inputs }, - }, - }), - ) + + resource_cfg = { + DAG_CONFIG_KEY: DAGConfiguration( + output_directory=result_directory, + input_files=job_inputs, + user_access_token=access_token, + output_result_id=str(new_result_out.id), + ), + NOTIFIER_KEY: { "config": { + "success_callback_url": f"{WORKFLOW_MANAGER_URI}/result_ready/{new_result_out.id}", + "devicemanager_url": f"{DEVICE_MANAGER_URI}/parameter/", + "access_token": access_token, + }} + } + + try: + run_id = dg_client.submit_job_execution( + job_name=job_name, + repository_location_name=location, + repository_name=repository, + run_config=RunConfig(resources=resource_cfg), + # run_config=RunConfig(resources={ + # SCANHUB_RESOURCE_KEY: JobConfigResource( + # callback_url=callback_endpoint, + # user_access_token=access_token, + # input_files=job_inputs, + # output_dir=result_directory, + # task_id=task_id, + # exam_id=exam_id, + # update_device_parameter_base_url=device_parameter_update_endpoint, + # ), + # }), + ) + except DagsterGraphQLClientError as exc: + raise HTTPException(status_code=400, detail=f"Dagster submission failed: {exc}") from exc if run_id: # Update result @@ -263,15 +276,15 @@ def handle_dag_task_trigger( _ = set_task(task.id, updated_task, access_token) return {"message": "Failed to start DAG, no run_id..."} except Exception as exc: - logging.error(f"Failed to trigger DAG: {exc}") + print(f"Failed to trigger DAG: {exc}") raise HTTPException(status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=str(exc)) -@router.post("/result_ready/{task_id}/{result_id}", tags=["WorkflowManager"]) +@router.post("/result_ready/{result_id}", tags=["WorkflowManager"]) async def callback_results_ready( - task_id: UUID | str, result_id: UUID | str, - access_token: Annotated[str, Depends(oauth2_scheme)] + access_token: Annotated[str, Depends(oauth2_scheme)], + success: bool = Body(..., embed=True), ) -> dict[str, Any]: """ Notify that results are ready via callback endpoint. @@ -284,49 +297,45 @@ async def callback_results_ready( ------- dict: A dictionary containing a success message. """ - if not isinstance(task := get_task(task_id, access_token), DAGTaskOut): + if not isinstance(result := get_result(result_id, access_token), ResultOut): raise HTTPException( status_code=status.HTTP_400_BAD_REQUEST, - detail=f"Result ready callback received invalid task_id: {task.id}", + detail=f"Result ready callback received invalid result_id: {result_id.id}", ) - if not isinstance(result := get_result(result_id, access_token), ResultOut): + if not isinstance(task := get_task(result.task_id, access_token), DAGTaskOut): raise HTTPException( status_code=status.HTTP_400_BAD_REQUEST, - detail=f"Result ready callback received invalid task_id: {task.id}", + detail=f"Result ready callback received invalid task_id: {result.task_id}", ) - - print("\n>>>>>\nCallback received: ", result.model_dump_json()) - - # Update task status - task.status = ItemStatus.FINISHED - task.progress = 100 - _ = set_task(task.id, task, access_token) - - # Get a list of dicom files from the result directory - result_dir = Path(result.directory) - dicom_files = sorted(result_dir.rglob("*.dcm")) - - # Get dicom location in shared data lake - workflow_folder = result_dir.parts[-3] - task_folder = result_dir.parts[-2] - result_folder = result_dir.parts[-1] - - # Add file names to the result - result.files = [str(_file.name) for _file in dicom_files] - # Add meta information - meta_update = { - "instance_count": len(dicom_files), - "instances": [ - f"{DICOM_BASE_URI}{workflow_folder}/{task_folder}/{result_folder}/{_file.name}" for _file in dicom_files - ], - } - if result.meta is not None: - result.meta.update(meta_update) + # Update task/result depending on success + if not success: + task.status = ItemStatus.ERROR + task.progress = 0 + _ = set_task(task.id, task, access_token) else: - result.meta = meta_update - print(f"Updated result: {result.model_dump_json()}") - _ = set_result(result_id=result.id, payload=result, user_access_token=access_token) - + task.status = ItemStatus.FINISHED + task.progress = 100 + _ = set_task(task.id, task, access_token) + + # Get a list of dicom files from the result directory + result_dir = Path(result.directory) + dcm_files = sorted(result_dir.rglob("*.dcm")) + result.files = [str(_file.name) for _file in dcm_files] + + # Get dicom location in data lake and update meta + workflow_folder, task_folder, result_folder = result_dir.parts[-3:] + meta_update = { + "instance_count": len(dcm_files), + "instances": [ + f"{DICOM_BASE_URI}{workflow_folder}/{task_folder}/{result_folder}/{_file.name}" for _file in dcm_files + ], + } + if result.meta is not None: + result.meta.update(meta_update) + else: + result.meta = meta_update + print(f"Updated result: {result.model_dump_json()}") + _ = set_result(result_id=result.id, payload=result, user_access_token=access_token) return {"message": "Results ready notification sent successfully."} From bc4f2dbbbe5dc5440784cd6763ee11327c76615e Mon Sep 17 00:00:00 2001 From: David Schote Date: Thu, 4 Sep 2025 22:38:02 +0200 Subject: [PATCH 03/57] Added dagster resource definitions to base image (can be used in manager) --- .../src/scanhub_libraries/resources.py | 29 --------- .../scanhub_libraries/resources/__init__.py | 6 ++ .../scanhub_libraries/resources/dag_config.py | 11 ++++ .../scanhub_libraries/resources/data_lake.py | 63 +++++++++++++++++++ .../scanhub_libraries/resources/notifier.py | 36 +++++++++++ 5 files changed, 116 insertions(+), 29 deletions(-) delete mode 100644 services/base/shared_libs/src/scanhub_libraries/resources.py create mode 100644 services/base/shared_libs/src/scanhub_libraries/resources/__init__.py create mode 100644 services/base/shared_libs/src/scanhub_libraries/resources/dag_config.py create mode 100644 services/base/shared_libs/src/scanhub_libraries/resources/data_lake.py create mode 100644 services/base/shared_libs/src/scanhub_libraries/resources/notifier.py diff --git a/services/base/shared_libs/src/scanhub_libraries/resources.py b/services/base/shared_libs/src/scanhub_libraries/resources.py deleted file mode 100644 index 354b8976..00000000 --- a/services/base/shared_libs/src/scanhub_libraries/resources.py +++ /dev/null @@ -1,29 +0,0 @@ -"""Define the resource used by Dagster jobs.""" -from dagster import ConfigurableResource - - -DAG_CONFIG_KEY = "dag_config" -DATA_LAKE_KEY = "data_lake" -IDATA_IO_KEY = "idata_io_manager" -NOTIFIER_KEY = "scanhub_notifier" - - -class JobConfigResource(ConfigurableResource): - """Resource used by the job.""" - - callback_url: str | None = None - user_access_token: str | None = None - update_device_parameter_base_url: str | None = None - input_files: list[str] - output_dir: str - task_id: str # serves as series id - exam_id: str # serves as study id - - -class DAGConfiguration(ConfigurableResource): - """Run-scoped parameters accessible from assets + IO managers.""" - - output_directory: str = "" - input_files: list[str] = [] - user_access_token: str = "" - output_result_id: str = "" diff --git a/services/base/shared_libs/src/scanhub_libraries/resources/__init__.py b/services/base/shared_libs/src/scanhub_libraries/resources/__init__.py new file mode 100644 index 00000000..ce9945e5 --- /dev/null +++ b/services/base/shared_libs/src/scanhub_libraries/resources/__init__.py @@ -0,0 +1,6 @@ +"""Ressource init.""" +# Define keys for ressources +DAG_CONFIG_KEY = "dag_config" +DATA_LAKE_KEY = "data_lake" +IDATA_IO_KEY = "idata_io_manager" +NOTIFIER_KEY = "scanhub_notifier" diff --git a/services/base/shared_libs/src/scanhub_libraries/resources/dag_config.py b/services/base/shared_libs/src/scanhub_libraries/resources/dag_config.py new file mode 100644 index 00000000..585786ed --- /dev/null +++ b/services/base/shared_libs/src/scanhub_libraries/resources/dag_config.py @@ -0,0 +1,11 @@ +"""Define the resource used by Dagster jobs.""" +from dagster import ConfigurableResource + + +class DAGConfiguration(ConfigurableResource): + """Run-scoped parameters accessible from assets + IO managers.""" + + output_directory: str = "" + input_files: list[str] = [] + user_access_token: str = "" + output_result_id: str = "" diff --git a/services/base/shared_libs/src/scanhub_libraries/resources/data_lake.py b/services/base/shared_libs/src/scanhub_libraries/resources/data_lake.py new file mode 100644 index 00000000..e36ce208 --- /dev/null +++ b/services/base/shared_libs/src/scanhub_libraries/resources/data_lake.py @@ -0,0 +1,63 @@ +"""Definition of dagster data lake ressource for acquisition data.""" +import json +from pathlib import Path + +from dagster import ConfigurableResource + + +class DataLakeResource(ConfigurableResource): + """Dagster data lake ressource.""" + + def get_mrd_path(self, files: list[str]) -> Path: + """Construct and validate the MRD path inside the data lake. + + Parameters + ---------- + directory : str + Absolute path to result data (contains data lake path). + filenames: list[str] + List of filenames which can be found in directory. + + Returns + ------- + path + Path to acquisition ISMRMRD file. + + """ + filename = next((f for f in files if f.lower().endswith(".mrd")), None) + if filename is None: + raise FileNotFoundError(f"Acquisition result does not specify mrd filename.") + if not (mrd_path := Path(filename)).is_file(): + raise FileNotFoundError(f"MRD file does not exist: {mrd_path}") + return mrd_path + + def get_device_parameter(self, files: list[str]) -> tuple[str, dict]: + """Return the path to the device parameter JSON file if it exists. + + Parameters + ---------- + directory : str + Absolute path to result data (contains data lake path). + filenames: list[str] + List of filenames which can be found in directory. + + Returns + ------- + dict + Dictionary containing device parameters + + """ + json_file = next((f for f in files if f.lower().endswith(".json")), None) + if json_file is None: + raise FileNotFoundError(f"Acquisition result does not specify device parameter file.") + # Check if parameter file exists + if not (json_path := Path(json_file)).exists(): + raise FileExistsError(f"Device parameter file does not exist: {json_path}") + # Load parameter file + with json_path.open("r") as fh: + data = json.load(fh) + # Check if parameter file contains device id and parameter + if "device_id" not in data or "parameter" not in data: + raise AttributeError(f"Invalid paraeter file: {json_path}") + + return (str(data["device_id"]), data["parameter"]) diff --git a/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py b/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py new file mode 100644 index 00000000..3126e3ee --- /dev/null +++ b/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py @@ -0,0 +1,36 @@ +# orchestrator/resources/notifier.py +import httpx +from dagster import ConfigurableResource + + +class BackendNotifier(ConfigurableResource): + """Backend notifier.""" + + success_callback_url: str | None = None + devicemanager_url: str | None = None + access_token: str | None = None + timeout: float = 5.0 + + def send_dag_success(self, success: bool = True) -> None: + """Notify backend about successful execution of dagster job/dag.""" + if self.success_callback_url is None: + raise AttributeError + if self.access_token is None: + raise AttributeError + headers = {"Authorization": "Bearer " + self.access_token} + payload = {"success": success} + with httpx.Client(timeout=self.timeout) as client: + response = client.post(self.success_callback_url, json=payload, headers=headers) + response.raise_for_status() + + def send_device_parameter_update(self, device_id: str, parameter: dict) -> None: + """Notify backend about device parameter update and send parameters.""" + if self.devicemanager_url is None: + raise AttributeError + if self.access_token is None: + raise AttributeError + headers = {"Authorization": "Bearer " + self.access_token} + url = self.devicemanager_url.rstrip("/") + f"/{device_id}" + with httpx.Client(timeout=self.timeout) as client: + response = client.put(url, json=parameter, headers=headers) + response.raise_for_status() From 9095bfd68ac974cbed751ffcb424bd3f4bdd2c02 Mon Sep 17 00:00:00 2001 From: David Schote Date: Thu, 4 Sep 2025 22:38:32 +0200 Subject: [PATCH 04/57] Fixed dagster asset and job definitions for reco and calibration --- services/orchestration-engine/README.md | 12 ++- .../orchestrator/assets/acquisition.py | 44 --------- .../assets/mrpro_direct_reconstruction.py | 13 +-- .../{hooks/scanhub.py => hooks.py} | 24 +++++ .../orchestrator/{hooks => io}/__init__.py | 0 .../orchestrator/io/acquisition_data.py | 61 ++++++++++++ .../{io_managers => io}/idata_io_manager.py | 3 +- .../orchestrator/io_managers/__init__.py | 0 .../jobs/frequency_calibration.py | 94 ------------------- .../orchestrator/jobs/mri_calibration.py | 47 ++++++++++ .../jobs/mrpro_image_reconstruction.py | 93 ------------------ .../orchestrator/ops/__init__.py | 0 .../orchestrator/ops/scanhub.py | 49 ---------- .../orchestrator/repository copy.py | 16 ---- .../orchestrator/repository.py | 34 ++++--- .../orchestrator/ressources/__init__.py | 0 .../orchestrator/ressources/datalake.py | 62 ------------ .../orchestrator/ressources/notifier.py | 40 -------- .../app/api/manager_endpoints.py | 27 ++---- 19 files changed, 173 insertions(+), 446 deletions(-) delete mode 100644 services/orchestration-engine/orchestrator/assets/acquisition.py rename services/orchestration-engine/orchestrator/{hooks/scanhub.py => hooks.py} (53%) rename services/orchestration-engine/orchestrator/{hooks => io}/__init__.py (100%) create mode 100644 services/orchestration-engine/orchestrator/io/acquisition_data.py rename services/orchestration-engine/orchestrator/{io_managers => io}/idata_io_manager.py (94%) delete mode 100644 services/orchestration-engine/orchestrator/io_managers/__init__.py delete mode 100644 services/orchestration-engine/orchestrator/jobs/frequency_calibration.py create mode 100644 services/orchestration-engine/orchestrator/jobs/mri_calibration.py delete mode 100644 services/orchestration-engine/orchestrator/jobs/mrpro_image_reconstruction.py delete mode 100644 services/orchestration-engine/orchestrator/ops/__init__.py delete mode 100644 services/orchestration-engine/orchestrator/ops/scanhub.py delete mode 100644 services/orchestration-engine/orchestrator/repository copy.py delete mode 100644 services/orchestration-engine/orchestrator/ressources/__init__.py delete mode 100644 services/orchestration-engine/orchestrator/ressources/datalake.py delete mode 100644 services/orchestration-engine/orchestrator/ressources/notifier.py diff --git a/services/orchestration-engine/README.md b/services/orchestration-engine/README.md index d85841a0..150e1451 100644 --- a/services/orchestration-engine/README.md +++ b/services/orchestration-engine/README.md @@ -13,8 +13,9 @@ orchestration-engine ├── LICENSE ├── orchestrator/ │ ├── jobs/ -│ ├── ops/ -│ ├── ressources/ +│ ├── assets/ +│ ├── io/ +│ ├── hooks.py │ └── repository.py ├── poetry.lock └── pyproject.toml @@ -24,9 +25,10 @@ orchestration-engine - Dockerfile — Container build definition for this service - orchestrator/ — Python package with orchestration logic - jobs/ — Dagster jobs (pipelines/graphs) such as frequency calibration and image reconstruction - - ops/ — Reusable Dagster ops (atomic computation steps) - - repository.py — Dagster repository definition that registers jobs, ops, and resources - - ressources/ — Dagster resources (interfaces to external systems like device manager, storage, etc.) + - assets/ — Reusable Dagster ops (atomic computation steps) + - io/ — Assets, ops and io managers related to input and output + - hooks.py - Hooks are used to notify the backend upon success of job execution + - repository.py - Configuration/definition of the dagster repo components ## 🚀 Getting Started diff --git a/services/orchestration-engine/orchestrator/assets/acquisition.py b/services/orchestration-engine/orchestrator/assets/acquisition.py deleted file mode 100644 index 47dfb116..00000000 --- a/services/orchestration-engine/orchestrator/assets/acquisition.py +++ /dev/null @@ -1,44 +0,0 @@ -"""Definition of acquisiton data assets.""" -from dataclasses import dataclass -from pathlib import Path - -from dagster import AssetExecutionContext, MetadataValue, asset -from scanhub_libraries.resources import DAGConfiguration - -from orchestrator.ressources.datalake import DataLakeResource - - -@dataclass -class AcquisitionData: - """Acquisition data output of read acquisition data asset.""" - - mrd_path: Path - device_parameter: dict - - -@asset( - group_name="io", - description="Provides acquired ISMRMRD result", -) -def read_acquisition_data( - context: AssetExecutionContext, - dag_config: DAGConfiguration, - data_lake: DataLakeResource, -) -> AcquisitionData: - """Define acquisition result asset.""" - mrd_file = data_lake.get_mrd_path(dag_config.input_files) - context.log.info("MRD file path: %s", str(mrd_file)) - device_parameter = data_lake.get_device_parameter(dag_config.input_files) - context.log.info("Device parameters: %s", str(device_parameter)) - - # Optional: surface helpful metadata in the Dagster UI - context.add_output_metadata({ - "mrd_path": MetadataValue.path(str(mrd_file)), - "device_parameter": MetadataValue.json(device_parameter), - "num_input_files": len(dag_config.input_files), - "output_directory": MetadataValue.path(dag_config.output_directory), - "output_result_id": dag_config.output_result_id, - "access_token": dag_config.user_access_token, - }) - - return AcquisitionData(mrd_path=mrd_file, device_parameter=device_parameter) diff --git a/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py b/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py index 055547a0..f9c73005 100644 --- a/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py +++ b/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py @@ -1,17 +1,18 @@ import mrpro from dagster import AssetIn, AssetKey, asset -from scanhub_libraries.resources import DAGConfiguration +from scanhub_libraries.resources import IDATA_IO_KEY +from scanhub_libraries.resources.dag_config import DAGConfiguration -from orchestrator.assets.acquisition import AcquisitionData -from orchestrator.hooks.scanhub import notify_dag_success -from orchestrator.io_managers.idata_io_manager import IDataContext +from orchestrator.hooks import notify_dag_success +from orchestrator.io.acquisition_data import AcquisitionData, acquisition_data_asset +from orchestrator.io.idata_io_manager import IDataContext @asset( group_name="reconstruction", description="MRpro direct reconstruction.", - ins={"data": AssetIn(key=AssetKey("read_acquisition_data"))}, - io_manager_key="idata_io_manager", + ins={"data": AssetIn(key=acquisition_data_asset.key)}, + io_manager_key=IDATA_IO_KEY, hooks={notify_dag_success}, ) def mrpro_direct_reconstruction(context, data: AcquisitionData, dag_config: DAGConfiguration) -> IDataContext: diff --git a/services/orchestration-engine/orchestrator/hooks/scanhub.py b/services/orchestration-engine/orchestrator/hooks.py similarity index 53% rename from services/orchestration-engine/orchestrator/hooks/scanhub.py rename to services/orchestration-engine/orchestrator/hooks.py index 8757e4b9..5bbd4c99 100644 --- a/services/orchestration-engine/orchestrator/hooks/scanhub.py +++ b/services/orchestration-engine/orchestrator/hooks.py @@ -27,3 +27,27 @@ def notify_dag_success(context: HookContext) -> None: notifier = getattr(context.resources, NOTIFIER_KEY) notifier.send_dag_success(success=True) + + +@success_hook(required_resource_keys={NOTIFIER_KEY}) +def notify_device_parameter_update(context: HookContext) -> None: + """Check if data has been written and notify backend.""" + # Expect the asset to return an IDataContext + result = next(iter(context.op_output_values.values()), None) + if result is None: + context.log.warning("NOTIFY-HOOK: No return value; skipping.") + return + dag_cfg = getattr(result, "dag_config", None) + out_dir = getattr(dag_cfg, "output_directory", None) + if not out_dir: + context.log.warning("NOTIFY-HOOK: No output_directory; skipping.") + return + + p = Path(out_dir) + files = [str(f) for f in p.glob("**/*") if f.is_file()] + if not files: + context.log.warning(f"NOTIFY-HOOK: No files in {p}; skipping.") + return + + notifier = getattr(context.resources, NOTIFIER_KEY) + notifier.send_dag_success(success=True) diff --git a/services/orchestration-engine/orchestrator/hooks/__init__.py b/services/orchestration-engine/orchestrator/io/__init__.py similarity index 100% rename from services/orchestration-engine/orchestrator/hooks/__init__.py rename to services/orchestration-engine/orchestrator/io/__init__.py diff --git a/services/orchestration-engine/orchestrator/io/acquisition_data.py b/services/orchestration-engine/orchestrator/io/acquisition_data.py new file mode 100644 index 00000000..5415578a --- /dev/null +++ b/services/orchestration-engine/orchestrator/io/acquisition_data.py @@ -0,0 +1,61 @@ +"""Definition of acquisiton data assets.""" +from dataclasses import dataclass +from pathlib import Path + +from dagster import AssetExecutionContext, MetadataValue, OpExecutionContext, asset, op +from scanhub_libraries.resources.dag_config import DAGConfiguration +from scanhub_libraries.resources.data_lake import DataLakeResource + + +@dataclass +class AcquisitionData: + """Acquisition data output of read acquisition data asset.""" + + mrd_path: Path + device_id: str + device_parameter: dict + + +def _load_acquisition_data(dag_config: DAGConfiguration, data_lake: DataLakeResource) -> AcquisitionData: + """Load acquisition data.""" + mrd_file = data_lake.get_mrd_path(dag_config.input_files) + device_id, device_parameter = data_lake.get_device_parameter(dag_config.input_files) + return AcquisitionData(mrd_path=mrd_file, device_id=device_id, device_parameter=device_parameter) + + +@asset( + group_name="io", + description="Provides acquired ISMRMRD result", +) +def acquisition_data_asset( + context: AssetExecutionContext, + dag_config: DAGConfiguration, + data_lake: DataLakeResource, +) -> AcquisitionData: + """Define acquisition result asset.""" + data = _load_acquisition_data(dag_config, data_lake) + context.log.info("Parameters for device id %s: %s", data.device_id, data.device_parameter) + + # Optional: Add meta data for dagster UI + context.add_output_metadata({ + "mrd_path": MetadataValue.path(str(data.mrd_path)), + "device_parameter": MetadataValue.json(data.device_parameter), + "num_input_files": len(dag_config.input_files), + "output_directory": MetadataValue.path(dag_config.output_directory), + "output_result_id": dag_config.output_result_id, + "access_token": dag_config.user_access_token, + }) + return data + + +@op +def acquisition_data_op( + context: OpExecutionContext, + dag_config: DAGConfiguration, + data_lake: DataLakeResource, +) -> AcquisitionData: + """Op version of the loader so the job can run end-to-end without assets.""" + data = _load_acquisition_data(dag_config, data_lake) + context.log.info("MRD file path: %s", str(data.mrd_path)) + context.log.info("Parameters for device id %s: %s", data.device_id, data.device_parameter) + return data diff --git a/services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py b/services/orchestration-engine/orchestrator/io/idata_io_manager.py similarity index 94% rename from services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py rename to services/orchestration-engine/orchestrator/io/idata_io_manager.py index b9ac53f1..ea1a6a64 100644 --- a/services/orchestration-engine/orchestrator/io_managers/idata_io_manager.py +++ b/services/orchestration-engine/orchestrator/io/idata_io_manager.py @@ -4,7 +4,7 @@ from dagster import ConfigurableIOManager, InputContext, OutputContext from mrpro.data import IData -from scanhub_libraries.resources import DAGConfiguration +from scanhub_libraries.resources.dag_config import DAGConfiguration @dataclass @@ -26,6 +26,7 @@ def handle_output(self, context: OutputContext, obj: IDataContext) -> None: raise AttributeError directory_path = Path(obj.dag_config.output_directory) + # Save data as dicom to result folder obj.data.to_dicom_folder(directory_path) # Surface paths in the UI and for hooks/sensors diff --git a/services/orchestration-engine/orchestrator/io_managers/__init__.py b/services/orchestration-engine/orchestrator/io_managers/__init__.py deleted file mode 100644 index e69de29b..00000000 diff --git a/services/orchestration-engine/orchestrator/jobs/frequency_calibration.py b/services/orchestration-engine/orchestrator/jobs/frequency_calibration.py deleted file mode 100644 index 399ff591..00000000 --- a/services/orchestration-engine/orchestrator/jobs/frequency_calibration.py +++ /dev/null @@ -1,94 +0,0 @@ -# %% -import json -import tempfile -from pathlib import Path - -from dagster import RunConfig, graph, op -from scanhub_libraries.resources import SCANHUB_RESOURCE_KEY, JobConfigResource - -from orchestrator.ops.scanhub import notify_op, send_device_parameter_update - - -@op(required_resource_keys={SCANHUB_RESOURCE_KEY}) -def load_parameter(context) -> dict | None: - """Load device parameter file.""" - job_config: JobConfigResource = context.resources.job_config - # Get correct input: - parameter_file = None - for _file in job_config.input_files: - if _file.endswith(".json"): - parameter_file = Path(_file) - if not parameter_file.exists(): - context.log.error("Parameter file does not exist") - - if parameter_file is not None: - with parameter_file.open("r") as fh: - data = json.load(fh) - - if "device_id" not in data or "parameter" not in data: - context.log.error("Invalid parameter file") - - context.log.info("Loaded parameter file: %s", data) - return data - - context.log.error("No parameter file found in job input files.") - return None - - -@op -def update_parameters(context, data: dict) -> dict: - """Modify device parameter.""" - context.log.info("Device ID: %s", data["device_id"]) - context.log.info("Parameter: %s", data["parameter"]) - data["parameter"]["larmor_frequency"] += 10 - context.log.info("Updated parameter: %s", data["parameter"]) - return data - - -@graph -def frequency_calibration_graph() -> None: - """Define frequency calibration graph.""" - data = load_parameter() - if data is not None: - data = update_parameters(data) - done = send_device_parameter_update(data) - notify_op(done) - -# Define the job with the resource definition (uses defaults here) -frequency_calibration_job = frequency_calibration_graph.to_job( - name="frequency_calibration_job", - resource_defs={SCANHUB_RESOURCE_KEY: JobConfigResource.configure_at_launch()}, -) - -# %% -if __name__ == "__main__": - print("Testing frequency calibration job...") - data = {'device_id': 'f977740a-5ddf-4364-9a40-941cd49d5050', 'parameter': {'larmor_frequency': 2025000.0}} - with tempfile.TemporaryDirectory() as tmpdir: - file_path = Path(tmpdir) / "device_config.json" - with file_path.open("w") as fh: - json.dump(data, fh) - - result = frequency_calibration_job.execute_in_process( - run_config=RunConfig( - resources={ - SCANHUB_RESOURCE_KEY: JobConfigResource( - callback_url=None, - user_access_token=None, - input_files=[str(file_path)], - output_dir=tmpdir + "/output", - update_device_parameter_base_url=None, - task_id="", - exam_id="", - ) - }, - loggers={"console": {"config": {"log_level": "INFO"}}}, - ) - ) - - # Report logs - for event in result.all_events: - print(f"{event.event_type_value} > {event.message}") - - -# %% diff --git a/services/orchestration-engine/orchestrator/jobs/mri_calibration.py b/services/orchestration-engine/orchestrator/jobs/mri_calibration.py new file mode 100644 index 00000000..0f985ac0 --- /dev/null +++ b/services/orchestration-engine/orchestrator/jobs/mri_calibration.py @@ -0,0 +1,47 @@ +"""MRI calibration assets.""" +import ismrmrd +from dagster import In, Nothing, OpExecutionContext, job, op, In, AssetKey +from scanhub_libraries.resources.notifier import BackendNotifier + +from orchestrator.io.acquisition_data import AcquisitionData, acquisition_data_op + + +@op +def frequency_calibration_op( + context: OpExecutionContext, + data: AcquisitionData, + scanhub_notifier: BackendNotifier, +) -> None: + """Perform frequency calibration.""" + context.log.info("AcquisitionData: %s", data) + parameter = data.device_parameter + context.log.info(f"Received device parameters (device id: {data.device_id}): {parameter}") + + with ismrmrd.File(data.mrd_path, "r") as f: + ds = f[list(f.keys())[0]] + header = ds.header + acquisitions = ds.acquisitions[:] + + f0 = header.experimentalConditions.H1resonanceFrequency_Hz + context.log.info(f"Larmor frequency from ismrmrd file: {f0}") + context.log.info(f"Number of readouts: {len(acquisitions)}") + + parameter["larmor_frequency"] = f0 + scanhub_notifier.send_device_parameter_update(device_id=data.device_id, parameter=parameter) + + +@op(ins={"start": In(Nothing), }) +def success_notify_op(context: OpExecutionContext, scanhub_notifier: BackendNotifier) -> None: + """Single, end-of-run success signal.""" + scanhub_notifier.send_dag_success() + context.log.info("Backend notification send.") + +@job( + name="frequency_calibration", + description="Read ismrmrd data and adjust Larmor frequency.", +) +def frequency_calibration_job() -> None: + """Calibrate Larmor frequency.""" + data = acquisition_data_op() + done = frequency_calibration_op(data) + success_notify_op(start=done) diff --git a/services/orchestration-engine/orchestrator/jobs/mrpro_image_reconstruction.py b/services/orchestration-engine/orchestrator/jobs/mrpro_image_reconstruction.py deleted file mode 100644 index 1e708d8c..00000000 --- a/services/orchestration-engine/orchestrator/jobs/mrpro_image_reconstruction.py +++ /dev/null @@ -1,93 +0,0 @@ -# %% -import tempfile -from pathlib import Path - -import mrpro -from dagster import Nothing, OpExecutionContext, Out, RunConfig, graph, op -from scanhub_libraries.resources import SCANHUB_RESOURCE_KEY, JobConfigResource - -from orchestrator.ops.scanhub import notify_op - - -@op(required_resource_keys={SCANHUB_RESOURCE_KEY}) -def load_kdata(context: OpExecutionContext) -> mrpro.data.KData: - """Load k-space data from MRD file.""" - job_config: JobConfigResource = context.resources.job_config - path = Path(job_config.input_files[0]) # expect a single input file - context.log.info("Loading data from %s", path) - - if not path.exists(): - detail = f"Input file {path} does not exist" - context.log.error(detail) - raise FileNotFoundError(detail) - - # trajectory_calculator = mrpro.data.traj_calculators.KTrajectoryIsmrmrd() - trajectory_calculator = mrpro.data.traj_calculators.KTrajectoryCartesian() - kdata = mrpro.data.KData.from_file(path, trajectory_calculator) - context.log.info("Loaded data: %s", kdata.shape) - - return kdata - - -@op -def direct_reconstruction(context: OpExecutionContext, kdata: mrpro.data.KData) -> mrpro.data.IData: - """Perform direct reconstruction on k-space data.""" - context.log.info("Performing direct reconstruction using mrpro...") - - reconstruction = mrpro.algorithms.reconstruction.DirectReconstruction(kdata) - - return reconstruction(kdata) - - -@op(required_resource_keys={SCANHUB_RESOURCE_KEY}, out=Out(Nothing)) -def save_as_dicom(context: OpExecutionContext, img: mrpro.data.IData) -> None: - """Save reconstructed image as DICOM files.""" - job_config: JobConfigResource = context.resources.job_config - path = Path(job_config.output_dir) - if path.exists(): - context.log.error("Output directory %s already exists.", path) - context.log.info("Saving dicom files to %s", path) - img.to_dicom_folder(path) - context.log.info("Data saved successfully.") - return - - -@graph -def mrpro_reconstruct_graph() -> None: - """Graph for MRPro image reconstruction.""" - ksp = load_kdata() - img = direct_reconstruction(ksp) - done = save_as_dicom(img) - notify_op(done) - - -# Define the job with the resource definition (uses defaults here) -mrpro_reconstruction_job = mrpro_reconstruct_graph.to_job( - name="mrpro_direct_reconstruct_job", - resource_defs={SCANHUB_RESOURCE_KEY: JobConfigResource.configure_at_launch()}, -) - -# %% -if __name__ == "__main__": - with tempfile.TemporaryDirectory() as tmpdir: - input_path = "../../../../tools/examples/data.mrd" - - result = mrpro_reconstruction_job.execute_in_process( - run_config=RunConfig(resources={ - SCANHUB_RESOURCE_KEY: JobConfigResource( - callback_url=None, - user_access_token=None, - input_files=[str(input_path)], - output_dir=tmpdir + "/output", # Ensure output directory is specified - task_id="", - exam_id="", - update_device_parameter_base_url=None, - ) - }) - ) - - # Report logs - for event in result.all_events: - print(f"{event.event_type_value} > {event.message}") - -# %% diff --git a/services/orchestration-engine/orchestrator/ops/__init__.py b/services/orchestration-engine/orchestrator/ops/__init__.py deleted file mode 100644 index e69de29b..00000000 diff --git a/services/orchestration-engine/orchestrator/ops/scanhub.py b/services/orchestration-engine/orchestrator/ops/scanhub.py deleted file mode 100644 index 998b9c86..00000000 --- a/services/orchestration-engine/orchestrator/ops/scanhub.py +++ /dev/null @@ -1,49 +0,0 @@ -"""Definition of shared dagster operations.""" -import json - -import requests -from dagster import In, Nothing, OpExecutionContext, Out, op -from fastapi.encoders import jsonable_encoder -from scanhub_libraries.resources import SCANHUB_RESOURCE_KEY, JobConfigResource - - -@op(required_resource_keys={SCANHUB_RESOURCE_KEY}, ins={"after_save": In(Nothing)}) -def notify_op(context: OpExecutionContext) -> None: - """Notify backend about finished job using the callback url from the job configuration. - - This should always be the last step within the graph/job definition. - Input needs to be None, if operation should be executed sequentially. - """ - job_config: JobConfigResource = context.resources.job_config - if job_config.callback_url and job_config.user_access_token: - try: - headers = {"Authorization": "Bearer " + job_config.user_access_token} - with requests.post(job_config.callback_url, headers=headers, timeout=3) as callback_response: - if callback_response.status_code != 200: - context.log.error(f"Callback failed with status code {callback_response.status_code}") - else: - context.log.info("Callback successful") - except Exception as e: - context.log.error(f"Callback failed: {e}") - - -@op(required_resource_keys={SCANHUB_RESOURCE_KEY}, out=Out(Nothing)) -def send_device_parameter_update(context: OpExecutionContext, data: dict) -> None: - """Send an updated version of device parameters to device manager.""" - job_config: JobConfigResource = context.resources.job_config - if (url := job_config.update_device_parameter_base_url) and job_config.user_access_token: - try: - if not "device_id" in data or not "parameter" in data: - context.log.error("Invalid dictionary") - device_id = data["device_id"] - payload = json.dumps(data["parameter"], default=jsonable_encoder) - headers = {"Authorization": "Bearer " + job_config.user_access_token} - endpoint = url+device_id - context.log.info(f"Device parameter update endpoint: {endpoint}") - with requests.put(endpoint, data=payload, headers=headers, timeout=3) as callback_response: - if callback_response.status_code != 200: - context.log.error(f"Callback failed with status code {callback_response.status_code}") - else: - context.log.info(f"Callback successful. Updated device:\n{callback_response.json()}") - except Exception as e: - context.log.error(f"Callback failed: {e}") diff --git a/services/orchestration-engine/orchestrator/repository copy.py b/services/orchestration-engine/orchestrator/repository copy.py deleted file mode 100644 index 616e2357..00000000 --- a/services/orchestration-engine/orchestrator/repository copy.py +++ /dev/null @@ -1,16 +0,0 @@ -"""Define dagster repository.""" -from dagster import Definitions, in_process_executor - -from orchestrator.jobs.frequency_calibration import frequency_calibration_job -from orchestrator.jobs.mrpro_image_reconstruction import mrpro_reconstruction_job - -defs = Definitions( - jobs=[ - mrpro_reconstruction_job, - frequency_calibration_job, - ], - executor=in_process_executor, # default executor for all jobs -) - - - diff --git a/services/orchestration-engine/orchestrator/repository.py b/services/orchestration-engine/orchestrator/repository.py index 40cd7f27..fdec6387 100644 --- a/services/orchestration-engine/orchestrator/repository.py +++ b/services/orchestration-engine/orchestrator/repository.py @@ -2,40 +2,38 @@ import os from dagster import AssetSelection, Definitions, define_asset_job, in_process_executor +from scanhub_libraries.resources import DAG_CONFIG_KEY, DATA_LAKE_KEY, IDATA_IO_KEY, NOTIFIER_KEY +from scanhub_libraries.resources.dag_config import DAGConfiguration +from scanhub_libraries.resources.data_lake import DataLakeResource +from scanhub_libraries.resources.notifier import BackendNotifier -# from orchestrator.jobs.frequency_calibration import frequency_calibration_job -# from orchestrator.jobs.mrpro_image_reconstruction import mrpro_reconstruction_job -from orchestrator.assets.acquisition import read_acquisition_data from orchestrator.assets.mrpro_direct_reconstruction import mrpro_direct_reconstruction -from orchestrator.io_managers.idata_io_manager import IDataIOManager -from orchestrator.ressources.datalake import DataLakeResource -from orchestrator.ressources.notifier import BackendNotifier -from scanhub_libraries.resources import DAGConfiguration, DAG_CONFIG_KEY, IDATA_IO_KEY, DATA_LAKE_KEY, NOTIFIER_KEY +from orchestrator.io.acquisition_data import acquisition_data_asset +from orchestrator.io.idata_io_manager import IDataIOManager +from orchestrator.jobs.mri_calibration import frequency_calibration_job DATA_LAKE_DIR = os.getenv("DATA_LAKE_DIRECTORY", "data") -# if DATA_LAKE_DIR is None: # ensure that DATA_LAKE_DIR is set -# raise OSError("Missing `DATA_LAKE_DIRECTORY` environment variable.") - assets = [ - read_acquisition_data, + acquisition_data_asset, mrpro_direct_reconstruction, ] ressources = { - DAG_CONFIG_KEY: DAGConfiguration(), - DATA_LAKE_KEY: DataLakeResource(), - IDATA_IO_KEY: IDataIOManager(), - NOTIFIER_KEY: BackendNotifier(), + DAG_CONFIG_KEY: DAGConfiguration.configure_at_launch(), + DATA_LAKE_KEY: DataLakeResource.configure_at_launch(), + IDATA_IO_KEY: IDataIOManager.configure_at_launch(), + NOTIFIER_KEY: BackendNotifier.configure_at_launch(), } -jobs = [ +jobs = ( define_asset_job( name="mrpro_reconstruction_job", - selection=AssetSelection.keys("read_acquisition_data", "mrpro_direct_reconstruction"), + selection=AssetSelection.keys(acquisition_data_asset.key, mrpro_direct_reconstruction.key), # selection=[dg.AssetSelection.keys("reconstruct_numpy")], # tags={"workflow_id": "numpy"}, ), -] + frequency_calibration_job, +) defs = Definitions(assets=assets, resources=ressources, jobs=jobs, executor=in_process_executor) diff --git a/services/orchestration-engine/orchestrator/ressources/__init__.py b/services/orchestration-engine/orchestrator/ressources/__init__.py deleted file mode 100644 index e69de29b..00000000 diff --git a/services/orchestration-engine/orchestrator/ressources/datalake.py b/services/orchestration-engine/orchestrator/ressources/datalake.py deleted file mode 100644 index 5fbe538e..00000000 --- a/services/orchestration-engine/orchestrator/ressources/datalake.py +++ /dev/null @@ -1,62 +0,0 @@ -"""Definition of dagster data lake ressource for acquisition data.""" -import json -from pathlib import Path - -from dagster import ConfigurableResource - - -class DataLakeResource(ConfigurableResource): - """Dagster data lake ressource.""" - - def get_mrd_path(self, files: list[str]) -> Path: - """Construct and validate the MRD path inside the data lake. - - Parameters - ---------- - directory : str - Absolute path to result data (contains data lake path). - filenames: list[str] - List of filenames which can be found in directory. - - Returns - ------- - path - Path to acquisition ISMRMRD file. - - """ - filename = next((f for f in files if f.lower().endswith(".mrd")), None) - if filename is None: - raise FileNotFoundError(f"Acquisition result does not specify mrd filename.") - if not (mrd_path := Path(filename)).is_file(): - raise FileNotFoundError(f"MRD file does not exist: {mrd_path}") - return mrd_path - - def get_device_parameter(self, files: list[str]) -> dict: - """Return the path to the device parameter JSON file if it exists. - - Parameters - ---------- - directory : str - Absolute path to result data (contains data lake path). - filenames: list[str] - List of filenames which can be found in directory. - - Returns - ------- - dict - Dictionary containing device parameters - - """ - json_file = next((f for f in files if f.lower().endswith(".json")), None) - if json_file is None: - raise FileNotFoundError(f"Acquisition result does not specify device parameter file.") - # Check if parameter file exists - if not (json_path := Path(json_file)).exists(): - raise FileExistsError(f"Device parameter file does not exist: {json_path}") - # Load parameter file - with json_path.open("r") as fh: - data = json.load(fh) - # Check if parameter file contains device id and parameter - if "device_id" not in data or "parameter" not in data: - raise AttributeError(f"Invalid paraeter file: {json_path}") - return data diff --git a/services/orchestration-engine/orchestrator/ressources/notifier.py b/services/orchestration-engine/orchestrator/ressources/notifier.py deleted file mode 100644 index 82fc76c1..00000000 --- a/services/orchestration-engine/orchestrator/ressources/notifier.py +++ /dev/null @@ -1,40 +0,0 @@ -# orchestrator/resources/notifier.py -import json - -import httpx -from dagster import ConfigurableResource -from fastapi.encoders import jsonable_encoder - - -class BackendNotifier(ConfigurableResource): - """Backend notifier.""" - - success_callback_url: str | None = None - devicemanager_url: str | None = None - access_token: str | None = None - timeout: float = 5.0 - - def send_dag_success(self, success: bool) -> None: - """Notify backend about successful execution of dagster job/dag.""" - if self.success_callback_url is None: - raise AttributeError - if self.access_token is None: - raise AttributeError - headers = {"Authorization": "Bearer " + self.access_token} - payload = {"success": success} - with httpx.Client(timeout=self.timeout) as client: - response = client.post(self.success_callback_url, json=payload, headers=headers) - response.raise_for_status() - - def send_device_parameter_update(self, device_id: str, parameter: dict) -> None: - """Notify backend about device parameter update and send parameters.""" - if self.devicemanager_url is None: - raise AttributeError - if self.access_token is None: - raise AttributeError - headers = {"Authorization": "Bearer " + self.access_token} - url = self.devicemanager_url.rstrip("/") + f"/{device_id}" - payload = json.dumps(parameter, default=jsonable_encoder) - with httpx.Client(timeout=self.timeout) as client: - response = client.put(url, json=payload, headers=headers) - response.raise_for_status() diff --git a/services/workflow-manager/app/api/manager_endpoints.py b/services/workflow-manager/app/api/manager_endpoints.py index b3901edd..9bc04802 100644 --- a/services/workflow-manager/app/api/manager_endpoints.py +++ b/services/workflow-manager/app/api/manager_endpoints.py @@ -20,7 +20,7 @@ import requests from dagster import RunConfig from dagster_graphql import DagsterGraphQLClient, DagsterGraphQLClientError -from fastapi import APIRouter, Depends, HTTPException, status, Body +from fastapi import APIRouter, Body, Depends, HTTPException, status from fastapi.encoders import jsonable_encoder from fastapi.security import OAuth2PasswordBearer from scanhub_libraries.models import ( @@ -36,7 +36,9 @@ TaskType, WorkflowOut, ) -from scanhub_libraries.resources import NOTIFIER_KEY, DAG_CONFIG_KEY, DAGConfiguration +from scanhub_libraries.resources import DAG_CONFIG_KEY, NOTIFIER_KEY +from scanhub_libraries.resources.dag_config import DAGConfiguration +from scanhub_libraries.resources.notifier import BackendNotifier from scanhub_libraries.security import get_current_user from scanhub_libraries.utils import calc_age_from_date @@ -229,11 +231,11 @@ def handle_dag_task_trigger( user_access_token=access_token, output_result_id=str(new_result_out.id), ), - NOTIFIER_KEY: { "config": { - "success_callback_url": f"{WORKFLOW_MANAGER_URI}/result_ready/{new_result_out.id}", - "devicemanager_url": f"{DEVICE_MANAGER_URI}/parameter/", - "access_token": access_token, - }} + NOTIFIER_KEY: BackendNotifier( + success_callback_url=f"{WORKFLOW_MANAGER_URI}/result_ready/{new_result_out.id}", + devicemanager_url=f"{DEVICE_MANAGER_URI}/parameter/", + access_token=access_token, + ) } try: @@ -242,17 +244,6 @@ def handle_dag_task_trigger( repository_location_name=location, repository_name=repository, run_config=RunConfig(resources=resource_cfg), - # run_config=RunConfig(resources={ - # SCANHUB_RESOURCE_KEY: JobConfigResource( - # callback_url=callback_endpoint, - # user_access_token=access_token, - # input_files=job_inputs, - # output_dir=result_directory, - # task_id=task_id, - # exam_id=exam_id, - # update_device_parameter_base_url=device_parameter_update_endpoint, - # ), - # }), ) except DagsterGraphQLClientError as exc: raise HTTPException(status_code=400, detail=f"Dagster submission failed: {exc}") from exc From d6ffedb17ebd2a6103ce7e058eba86c2d291443a Mon Sep 17 00:00:00 2001 From: David Schote Date: Fri, 5 Sep 2025 01:36:24 +0200 Subject: [PATCH 05/57] Added sensors to orchestration engine --- .../scanhub_libraries/resources/__init__.py | 3 +- .../scanhub_libraries/resources/notifier.py | 36 +++--- .../assets/mrpro_direct_reconstruction.py | 5 +- .../orchestrator/hooks.py | 106 +++++++++--------- .../orchestrator/io/acquisition_data.py | 1 + .../orchestrator/jobs/mri_calibration.py | 47 -------- .../jobs/mri_frequency_calibration.py | 81 +++++++++++++ .../orchestrator/repository.py | 22 ++-- .../orchestrator/sensors.py | 65 +++++++++++ .../orchestrator/utils/__init__.py | 0 .../orchestrator/utils/snr.py | 35 ++++++ .../app/api/manager_endpoints.py | 16 ++- 12 files changed, 278 insertions(+), 139 deletions(-) delete mode 100644 services/orchestration-engine/orchestrator/jobs/mri_calibration.py create mode 100644 services/orchestration-engine/orchestrator/jobs/mri_frequency_calibration.py create mode 100644 services/orchestration-engine/orchestrator/sensors.py create mode 100644 services/orchestration-engine/orchestrator/utils/__init__.py create mode 100644 services/orchestration-engine/orchestrator/utils/snr.py diff --git a/services/base/shared_libs/src/scanhub_libraries/resources/__init__.py b/services/base/shared_libs/src/scanhub_libraries/resources/__init__.py index ce9945e5..e4288a75 100644 --- a/services/base/shared_libs/src/scanhub_libraries/resources/__init__.py +++ b/services/base/shared_libs/src/scanhub_libraries/resources/__init__.py @@ -3,4 +3,5 @@ DAG_CONFIG_KEY = "dag_config" DATA_LAKE_KEY = "data_lake" IDATA_IO_KEY = "idata_io_manager" -NOTIFIER_KEY = "scanhub_notifier" +NOTIFIER_WM_KEY = "notifier_workflow_manager" +NOTIFIER_DM_KEY = "notifier_device_manager" diff --git a/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py b/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py index 3126e3ee..416030a3 100644 --- a/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py +++ b/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py @@ -3,34 +3,32 @@ from dagster import ConfigurableResource -class BackendNotifier(ConfigurableResource): - """Backend notifier.""" +class WorkflowManagerNotifier(ConfigurableResource): + """Notifies device manager.""" - success_callback_url: str | None = None - devicemanager_url: str | None = None - access_token: str | None = None + base_url: str timeout: float = 5.0 - def send_dag_success(self, success: bool = True) -> None: + def send_dag_success(self, result_id: str, access_token: str, success: bool = True) -> None: """Notify backend about successful execution of dagster job/dag.""" - if self.success_callback_url is None: - raise AttributeError - if self.access_token is None: - raise AttributeError - headers = {"Authorization": "Bearer " + self.access_token} + headers = {"Authorization": "Bearer " + access_token} payload = {"success": success} + url = self.base_url.rstrip("/") + f"/result_ready/{result_id}" with httpx.Client(timeout=self.timeout) as client: - response = client.post(self.success_callback_url, json=payload, headers=headers) + response = client.post(url, json=payload, headers=headers) response.raise_for_status() - def send_device_parameter_update(self, device_id: str, parameter: dict) -> None: + +class DeviceManagerNotifier(ConfigurableResource): + """Notifies device manager.""" + + base_url: str + timeout: float = 5.0 + + def send_device_parameter_update(self, device_id: str, access_token: str, parameter: dict) -> None: """Notify backend about device parameter update and send parameters.""" - if self.devicemanager_url is None: - raise AttributeError - if self.access_token is None: - raise AttributeError - headers = {"Authorization": "Bearer " + self.access_token} - url = self.devicemanager_url.rstrip("/") + f"/{device_id}" + headers = {"Authorization": "Bearer " + access_token} + url = self.base_url.rstrip("/") + f"/{device_id}" with httpx.Client(timeout=self.timeout) as client: response = client.put(url, json=parameter, headers=headers) response.raise_for_status() diff --git a/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py b/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py index f9c73005..ae01916d 100644 --- a/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py +++ b/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py @@ -1,9 +1,9 @@ import mrpro -from dagster import AssetIn, AssetKey, asset +from dagster import AssetIn, asset from scanhub_libraries.resources import IDATA_IO_KEY from scanhub_libraries.resources.dag_config import DAGConfiguration -from orchestrator.hooks import notify_dag_success +# from orchestrator.hooks import notify_dag_success from orchestrator.io.acquisition_data import AcquisitionData, acquisition_data_asset from orchestrator.io.idata_io_manager import IDataContext @@ -13,7 +13,6 @@ description="MRpro direct reconstruction.", ins={"data": AssetIn(key=acquisition_data_asset.key)}, io_manager_key=IDATA_IO_KEY, - hooks={notify_dag_success}, ) def mrpro_direct_reconstruction(context, data: AcquisitionData, dag_config: DAGConfiguration) -> IDataContext: """Reconstruct image from a list acquisition results. diff --git a/services/orchestration-engine/orchestrator/hooks.py b/services/orchestration-engine/orchestrator/hooks.py index 5bbd4c99..1062dadb 100644 --- a/services/orchestration-engine/orchestrator/hooks.py +++ b/services/orchestration-engine/orchestrator/hooks.py @@ -1,53 +1,53 @@ -# orchestrator/hooks/recon_hooks.py -from pathlib import Path - -from dagster import HookContext, success_hook -from scanhub_libraries.resources import NOTIFIER_KEY - - -@success_hook(required_resource_keys={NOTIFIER_KEY}) -def notify_dag_success(context: HookContext) -> None: - """Check if data has been written and notify backend.""" - # Expect the asset to return an IDataContext - result = next(iter(context.op_output_values.values()), None) - if result is None: - context.log.warning("NOTIFY-HOOK: No return value; skipping.") - return - dag_cfg = getattr(result, "dag_config", None) - out_dir = getattr(dag_cfg, "output_directory", None) - if not out_dir: - context.log.warning("NOTIFY-HOOK: No output_directory; skipping.") - return - - p = Path(out_dir) - files = [str(f) for f in p.glob("**/*") if f.is_file()] - if not files: - context.log.warning(f"NOTIFY-HOOK: No files in {p}; skipping.") - return - - notifier = getattr(context.resources, NOTIFIER_KEY) - notifier.send_dag_success(success=True) - - -@success_hook(required_resource_keys={NOTIFIER_KEY}) -def notify_device_parameter_update(context: HookContext) -> None: - """Check if data has been written and notify backend.""" - # Expect the asset to return an IDataContext - result = next(iter(context.op_output_values.values()), None) - if result is None: - context.log.warning("NOTIFY-HOOK: No return value; skipping.") - return - dag_cfg = getattr(result, "dag_config", None) - out_dir = getattr(dag_cfg, "output_directory", None) - if not out_dir: - context.log.warning("NOTIFY-HOOK: No output_directory; skipping.") - return - - p = Path(out_dir) - files = [str(f) for f in p.glob("**/*") if f.is_file()] - if not files: - context.log.warning(f"NOTIFY-HOOK: No files in {p}; skipping.") - return - - notifier = getattr(context.resources, NOTIFIER_KEY) - notifier.send_dag_success(success=True) +# # orchestrator/hooks/recon_hooks.py +# from pathlib import Path + +# from dagster import HookContext, success_hook +# from scanhub_libraries.resources import NOTIFIER_KEY + + +# @success_hook(required_resource_keys={NOTIFIER_KEY}) +# def notify_dag_success(context: HookContext) -> None: +# """Check if data has been written and notify backend.""" +# # Expect the asset to return an IDataContext +# result = next(iter(context.op_output_values.values()), None) +# if result is None: +# context.log.warning("NOTIFY-HOOK: No return value; skipping.") +# return +# dag_cfg = getattr(result, "dag_config", None) +# out_dir = getattr(dag_cfg, "output_directory", None) +# if not out_dir: +# context.log.warning("NOTIFY-HOOK: No output_directory; skipping.") +# return + +# p = Path(out_dir) +# files = [str(f) for f in p.glob("**/*") if f.is_file()] +# if not files: +# context.log.warning(f"NOTIFY-HOOK: No files in {p}; skipping.") +# return + +# notifier = getattr(context.resources, NOTIFIER_KEY) +# notifier.send_dag_success(success=True) + + +# @success_hook(required_resource_keys={NOTIFIER_KEY}) +# def notify_device_parameter_update(context: HookContext) -> None: +# """Check if data has been written and notify backend.""" +# # Expect the asset to return an IDataContext +# result = next(iter(context.op_output_values.values()), None) +# if result is None: +# context.log.warning("NOTIFY-HOOK: No return value; skipping.") +# return +# dag_cfg = getattr(result, "dag_config", None) +# out_dir = getattr(dag_cfg, "output_directory", None) +# if not out_dir: +# context.log.warning("NOTIFY-HOOK: No output_directory; skipping.") +# return + +# p = Path(out_dir) +# files = [str(f) for f in p.glob("**/*") if f.is_file()] +# if not files: +# context.log.warning(f"NOTIFY-HOOK: No files in {p}; skipping.") +# return + +# notifier = getattr(context.resources, NOTIFIER_KEY) +# notifier.send_dag_success(success=True) diff --git a/services/orchestration-engine/orchestrator/io/acquisition_data.py b/services/orchestration-engine/orchestrator/io/acquisition_data.py index 5415578a..d9c8a37e 100644 --- a/services/orchestration-engine/orchestrator/io/acquisition_data.py +++ b/services/orchestration-engine/orchestrator/io/acquisition_data.py @@ -39,6 +39,7 @@ def acquisition_data_asset( # Optional: Add meta data for dagster UI context.add_output_metadata({ "mrd_path": MetadataValue.path(str(data.mrd_path)), + "device_id": data.device_id, "device_parameter": MetadataValue.json(data.device_parameter), "num_input_files": len(dag_config.input_files), "output_directory": MetadataValue.path(dag_config.output_directory), diff --git a/services/orchestration-engine/orchestrator/jobs/mri_calibration.py b/services/orchestration-engine/orchestrator/jobs/mri_calibration.py deleted file mode 100644 index 0f985ac0..00000000 --- a/services/orchestration-engine/orchestrator/jobs/mri_calibration.py +++ /dev/null @@ -1,47 +0,0 @@ -"""MRI calibration assets.""" -import ismrmrd -from dagster import In, Nothing, OpExecutionContext, job, op, In, AssetKey -from scanhub_libraries.resources.notifier import BackendNotifier - -from orchestrator.io.acquisition_data import AcquisitionData, acquisition_data_op - - -@op -def frequency_calibration_op( - context: OpExecutionContext, - data: AcquisitionData, - scanhub_notifier: BackendNotifier, -) -> None: - """Perform frequency calibration.""" - context.log.info("AcquisitionData: %s", data) - parameter = data.device_parameter - context.log.info(f"Received device parameters (device id: {data.device_id}): {parameter}") - - with ismrmrd.File(data.mrd_path, "r") as f: - ds = f[list(f.keys())[0]] - header = ds.header - acquisitions = ds.acquisitions[:] - - f0 = header.experimentalConditions.H1resonanceFrequency_Hz - context.log.info(f"Larmor frequency from ismrmrd file: {f0}") - context.log.info(f"Number of readouts: {len(acquisitions)}") - - parameter["larmor_frequency"] = f0 - scanhub_notifier.send_device_parameter_update(device_id=data.device_id, parameter=parameter) - - -@op(ins={"start": In(Nothing), }) -def success_notify_op(context: OpExecutionContext, scanhub_notifier: BackendNotifier) -> None: - """Single, end-of-run success signal.""" - scanhub_notifier.send_dag_success() - context.log.info("Backend notification send.") - -@job( - name="frequency_calibration", - description="Read ismrmrd data and adjust Larmor frequency.", -) -def frequency_calibration_job() -> None: - """Calibrate Larmor frequency.""" - data = acquisition_data_op() - done = frequency_calibration_op(data) - success_notify_op(start=done) diff --git a/services/orchestration-engine/orchestrator/jobs/mri_frequency_calibration.py b/services/orchestration-engine/orchestrator/jobs/mri_frequency_calibration.py new file mode 100644 index 00000000..1376ada6 --- /dev/null +++ b/services/orchestration-engine/orchestrator/jobs/mri_frequency_calibration.py @@ -0,0 +1,81 @@ +"""MRI calibration assets.""" +import ismrmrd +import numpy as np +from dagster import OpExecutionContext, job, op +from scanhub_libraries.resources.dag_config import DAGConfiguration +from scanhub_libraries.resources.notifier import DeviceManagerNotifier + +from orchestrator.io.acquisition_data import AcquisitionData, acquisition_data_op +from orchestrator.utils.snr import signal_to_noise_ratio + + +@op +def frequency_calibration_op( + context: OpExecutionContext, + data: AcquisitionData, + dag_config: DAGConfiguration, + notigier_device_manager: DeviceManagerNotifier, +) -> None: + """Perform frequency calibration.""" + context.log.info("AcquisitionData: %s", data) + parameter = data.device_parameter + context.log.info(f"Received device parameters (device id: {data.device_id}): {parameter}") + + with ismrmrd.File(data.mrd_path, "r") as f: + ds = f[list(f.keys())[0]] + header = ds.header + acquisitions = ds.acquisitions + + num_acquisitions = len(acquisitions) + if num_acquisitions > 1: + context.log.warning( + "Data is not 1D, got %s acquisition, processing first acquisition.", + num_acquisitions, + ) + + # Only consider first acquisition and data from coil at index 0 + acq = acquisitions[0] + acq_data = acq.data[0, ...] + dwell_time = acq.sample_time_us * 1e-6 + data_fft = np.fft.fftshift(np.fft.fft(np.fft.fftshift(acq_data))) + fft_freq = np.fft.fftshift(np.fft.fftfreq(acq_data.size, dwell_time)) + + # Calculate SNR and FWHM + bw_per_pixel = 1 / dwell_time / acq_data.size + snr, fwhm_px = signal_to_noise_ratio(data_fft) + fwhm_hz = round(fwhm_px*bw_per_pixel) + context.log.info("SNR: %s dB, FWHM: %s Hz", str(round(snr, 2)), str(fwhm_hz)) + + freq_offset = fft_freq[np.argmax(np.abs(data_fft))] + + # Get experiment larmor frequency + f0 = header.experimentalConditions.H1resonanceFrequency_Hz + context.log.info( + "Frequency from acquisition: %s Hz, Larmor frequency offset: %s Hz", + str(f0), + str(freq_offset), + ) + + # Correct device parameter + parameter["larmor_frequency"] = f0 + freq_offset + notigier_device_manager.send_device_parameter_update( + device_id=data.device_id, access_token=dag_config.user_access_token, parameter=parameter, + ) + + +# @op(ins={"start": In(Nothing), }) +# def success_notify_op(context: OpExecutionContext, scanhub_notifier: BackendNotifier) -> None: +# """Single, end-of-run success signal.""" +# scanhub_notifier.send_dag_success() +# context.log.info("Backend notification send.") + +@job( + name="frequency_calibration", + description="Read ismrmrd data and adjust Larmor frequency.", +) +def frequency_calibration_job() -> None: + """Calibrate Larmor frequency.""" + data = acquisition_data_op() + frequency_calibration_op(data) + # done = frequency_calibration_op(data) + # success_notify_op(start=done) diff --git a/services/orchestration-engine/orchestrator/repository.py b/services/orchestration-engine/orchestrator/repository.py index fdec6387..e8b45ad3 100644 --- a/services/orchestration-engine/orchestrator/repository.py +++ b/services/orchestration-engine/orchestrator/repository.py @@ -2,38 +2,46 @@ import os from dagster import AssetSelection, Definitions, define_asset_job, in_process_executor -from scanhub_libraries.resources import DAG_CONFIG_KEY, DATA_LAKE_KEY, IDATA_IO_KEY, NOTIFIER_KEY +from scanhub_libraries.resources import DAG_CONFIG_KEY, DATA_LAKE_KEY, IDATA_IO_KEY, NOTIFIER_DM_KEY, NOTIFIER_WM_KEY from scanhub_libraries.resources.dag_config import DAGConfiguration from scanhub_libraries.resources.data_lake import DataLakeResource -from scanhub_libraries.resources.notifier import BackendNotifier +from scanhub_libraries.resources.notifier import WorkflowManagerNotifier, DeviceManagerNotifier from orchestrator.assets.mrpro_direct_reconstruction import mrpro_direct_reconstruction from orchestrator.io.acquisition_data import acquisition_data_asset from orchestrator.io.idata_io_manager import IDataIOManager -from orchestrator.jobs.mri_calibration import frequency_calibration_job +from orchestrator.jobs.mri_frequency_calibration import frequency_calibration_job +from orchestrator.sensors import on_run_canceled, on_run_failure, on_run_success DATA_LAKE_DIR = os.getenv("DATA_LAKE_DIRECTORY", "data") +DEVICE_MANAGER_URI = "http://device-manager:8000/api/v1/device" +WORKFLOW_MANAGER_URI = "http://workflow-manager:8000/api/v1/workflowmanager" assets = [ acquisition_data_asset, mrpro_direct_reconstruction, ] +sensors = [ + on_run_success, + on_run_failure, + on_run_canceled, +] + ressources = { DAG_CONFIG_KEY: DAGConfiguration.configure_at_launch(), DATA_LAKE_KEY: DataLakeResource.configure_at_launch(), IDATA_IO_KEY: IDataIOManager.configure_at_launch(), - NOTIFIER_KEY: BackendNotifier.configure_at_launch(), + NOTIFIER_WM_KEY: WorkflowManagerNotifier(base_url=WORKFLOW_MANAGER_URI), + NOTIFIER_DM_KEY: DeviceManagerNotifier(base_url=DEVICE_MANAGER_URI), } jobs = ( define_asset_job( name="mrpro_reconstruction_job", selection=AssetSelection.keys(acquisition_data_asset.key, mrpro_direct_reconstruction.key), - # selection=[dg.AssetSelection.keys("reconstruct_numpy")], - # tags={"workflow_id": "numpy"}, ), frequency_calibration_job, ) -defs = Definitions(assets=assets, resources=ressources, jobs=jobs, executor=in_process_executor) +defs = Definitions(assets=assets, resources=ressources, jobs=jobs, sensors=sensors, executor=in_process_executor) diff --git a/services/orchestration-engine/orchestrator/sensors.py b/services/orchestration-engine/orchestrator/sensors.py new file mode 100644 index 00000000..02bec799 --- /dev/null +++ b/services/orchestration-engine/orchestrator/sensors.py @@ -0,0 +1,65 @@ +# sensors.py +from dagster import DagsterRunStatus, DefaultSensorStatus, RunStatusSensorContext, run_status_sensor +from scanhub_libraries.resources.notifier import WorkflowManagerNotifier +from scanhub_libraries.resources import DAG_CONFIG_KEY + + +def _get_dag_config_from_run(context: RunStatusSensorContext) -> dict: + return context.dagster_run.run_config.get("resources", {}).get(DAG_CONFIG_KEY, {}).get("config", {}) + + +@run_status_sensor( + run_status=DagsterRunStatus.SUCCESS, + default_status=DefaultSensorStatus.RUNNING, + monitor_all_code_locations=True, + minimum_interval_seconds=5, +) +def on_run_success(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier): + dag_config = _get_dag_config_from_run(context) + access_token = dag_config.get("user_access_token", "") + result_id = dag_config.get("output_result_id", "") + if result_id and access_token: + notifier_workflow_manager.send_dag_success(result_id=result_id, access_token=access_token, success=True) + context.log.info( + "%s succeeded (run_id=%s).", context.dagster_run.job_name, context.dagster_run.run_id, + ) + else: + context.log.info("Run succeeded, but can not report DAG status, missing access_token and/or result_id.") + + +@run_status_sensor( + run_status=DagsterRunStatus.SUCCESS, + default_status=DefaultSensorStatus.RUNNING, + monitor_all_code_locations=True, + minimum_interval_seconds=5, +) +def on_run_failure(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier): + dag_config = _get_dag_config_from_run(context) + access_token = dag_config.get("user_access_token", "") + result_id = dag_config.get("output_result_id", "") + if result_id and access_token: + notifier_workflow_manager.send_dag_success(result_id=result_id, access_token=access_token, success=False) + context.log.info( + "%s failed (run_id=%s).", context.dagster_run.job_name, context.dagster_run.run_id, + ) + else: + context.log.info("Run failed, but can not report DAG status, missing access_token and/or result_id.") + +@run_status_sensor( + run_status=DagsterRunStatus.SUCCESS, + default_status=DefaultSensorStatus.RUNNING, + monitor_all_code_locations=True, + minimum_interval_seconds=5, +) +def on_run_canceled(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier): + dag_config = _get_dag_config_from_run(context) + access_token = dag_config.get("user_access_token", "") + result_id = dag_config.get("output_result_id", "") + if result_id and access_token: + notifier_workflow_manager.send_dag_success(result_id=result_id, access_token=access_token, success=False) + context.log.info( + "%s canceled (run_id=%s).", context.dagster_run.job_name, context.dagster_run.run_id, + ) + else: + context.log.info("Run canceled, but can not report DAG status, missing access_token and/or result_id.") + diff --git a/services/orchestration-engine/orchestrator/utils/__init__.py b/services/orchestration-engine/orchestrator/utils/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/services/orchestration-engine/orchestrator/utils/snr.py b/services/orchestration-engine/orchestrator/utils/snr.py new file mode 100644 index 00000000..b5dd13fd --- /dev/null +++ b/services/orchestration-engine/orchestrator/utils/snr.py @@ -0,0 +1,35 @@ +"""Signal-to-noise ratio (SNR) calculation.""" +import numpy as np + + +def signal_to_noise_ratio(data: np.ndarray) -> tuple[float, int]: + """Calculate the signal to noise ratio in dB. + + Parameters + ---------- + data + Centered 1D spectrum of the acquired signal + + Returns + ------- + Tuple containing snr and fwhm in points + + """ + if data.ndim > 1: + raise AttributeError("Provide 1D data to calculate snr.") + data_abs = np.abs(data) + peak_idx = int(np.argmax(data_abs)) + peak = data_abs[peak_idx] + fwhm_left_idx = np.argmin(data_abs[:peak_idx] - peak / 2) + fwhm_right_idx = np.argmin(data_abs[peak_idx:] - peak / 2) + fwhm_px = int(fwhm_right_idx - fwhm_left_idx) + if fwhm_px <= 0: + raise ValueError("Invalid FWHM (result <= 0)") + + noise_left = data[:peak_idx-fwhm_px] + noise_right = data[peak_idx+fwhm_px:] + noise = np.concat((noise_left, noise_right)).std() + + snr_db = float(20 * np.log10(peak / noise)) + + return (snr_db, fwhm_px) diff --git a/services/workflow-manager/app/api/manager_endpoints.py b/services/workflow-manager/app/api/manager_endpoints.py index 9bc04802..1e5ea429 100644 --- a/services/workflow-manager/app/api/manager_endpoints.py +++ b/services/workflow-manager/app/api/manager_endpoints.py @@ -36,9 +36,9 @@ TaskType, WorkflowOut, ) -from scanhub_libraries.resources import DAG_CONFIG_KEY, NOTIFIER_KEY +from scanhub_libraries.resources import DAG_CONFIG_KEY from scanhub_libraries.resources.dag_config import DAGConfiguration -from scanhub_libraries.resources.notifier import BackendNotifier +# from scanhub_libraries.resources.notifier import BackendNotifier from scanhub_libraries.security import get_current_user from scanhub_libraries.utils import calc_age_from_date @@ -216,8 +216,6 @@ def handle_dag_task_trigger( new_result_out = create_blank_result(task_id, access_token) # Use internal url and http (not https and port 8443) because callback endpoint is requested from another docker container - callback_endpoint = f"{WORKFLOW_MANAGER_URI}/result_ready/{task.id}/{new_result_out.id}" - device_parameter_update_endpoint = f"{DEVICE_MANAGER_URI}/parameter/" result_directory = f"{DATA_LAKE_DIR}/{str(task.workflow_id)}/{str(task.id)}/{str(new_result_out.id)}/" # Trigger dagster job @@ -231,11 +229,11 @@ def handle_dag_task_trigger( user_access_token=access_token, output_result_id=str(new_result_out.id), ), - NOTIFIER_KEY: BackendNotifier( - success_callback_url=f"{WORKFLOW_MANAGER_URI}/result_ready/{new_result_out.id}", - devicemanager_url=f"{DEVICE_MANAGER_URI}/parameter/", - access_token=access_token, - ) + # NOTIFIER_KEY: BackendNotifier( + # success_callback_url=f"{WORKFLOW_MANAGER_URI}/result_ready/{new_result_out.id}", + # devicemanager_url=f"{DEVICE_MANAGER_URI}/parameter/", + # access_token=access_token, + # ), } try: From 7805cc072139e572a3b38b4b514d2a2ac478d2db Mon Sep 17 00:00:00 2001 From: David Schote Date: Fri, 5 Sep 2025 08:22:00 +0200 Subject: [PATCH 06/57] Added comments to sensors --- .../orchestrator/sensors.py | 19 +++++++++++++------ 1 file changed, 13 insertions(+), 6 deletions(-) diff --git a/services/orchestration-engine/orchestrator/sensors.py b/services/orchestration-engine/orchestrator/sensors.py index 02bec799..cfa7cce2 100644 --- a/services/orchestration-engine/orchestrator/sensors.py +++ b/services/orchestration-engine/orchestrator/sensors.py @@ -1,11 +1,15 @@ -# sensors.py +"""Sensors to notify workflow manager dependent on run status.""" from dagster import DagsterRunStatus, DefaultSensorStatus, RunStatusSensorContext, run_status_sensor -from scanhub_libraries.resources.notifier import WorkflowManagerNotifier from scanhub_libraries.resources import DAG_CONFIG_KEY +from scanhub_libraries.resources.notifier import WorkflowManagerNotifier def _get_dag_config_from_run(context: RunStatusSensorContext) -> dict: - return context.dagster_run.run_config.get("resources", {}).get(DAG_CONFIG_KEY, {}).get("config", {}) + """Get DAG configuration from run status sensor conrext.""" + run_config = getattr(context.dagster_run, "run_config", None) + if not isinstance(run_config, dict): + return {} + return run_config.get("resources", {}).get(DAG_CONFIG_KEY, {}).get("config", {}) @run_status_sensor( @@ -14,7 +18,8 @@ def _get_dag_config_from_run(context: RunStatusSensorContext) -> dict: monitor_all_code_locations=True, minimum_interval_seconds=5, ) -def on_run_success(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier): +def on_run_success(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier) -> None: + """Notify workflow manager about successful run using a sensor.""" dag_config = _get_dag_config_from_run(context) access_token = dag_config.get("user_access_token", "") result_id = dag_config.get("output_result_id", "") @@ -33,7 +38,8 @@ def on_run_success(context: RunStatusSensorContext, notifier_workflow_manager: W monitor_all_code_locations=True, minimum_interval_seconds=5, ) -def on_run_failure(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier): +def on_run_failure(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier) -> None: + """Notify workflow manager about failed run using a sensor.""" dag_config = _get_dag_config_from_run(context) access_token = dag_config.get("user_access_token", "") result_id = dag_config.get("output_result_id", "") @@ -51,7 +57,8 @@ def on_run_failure(context: RunStatusSensorContext, notifier_workflow_manager: W monitor_all_code_locations=True, minimum_interval_seconds=5, ) -def on_run_canceled(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier): +def on_run_canceled(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier) -> None: + """Notify workflow manager about cancelled run using a sensor.""" dag_config = _get_dag_config_from_run(context) access_token = dag_config.get("user_access_token", "") result_id = dag_config.get("output_result_id", "") From 3b0191b2d6fa53b040217fc8fcf520015ec2b778 Mon Sep 17 00:00:00 2001 From: David Schote Date: Fri, 5 Sep 2025 08:52:36 +0200 Subject: [PATCH 07/57] Documentation of orchestration engine --- .../assets/mrpro_direct_reconstruction.py | 25 +++++++-- .../orchestrator/hooks.py | 53 ------------------ .../orchestrator/io/acquisition_data.py | 48 +++++++++++++++- .../orchestrator/io/idata_io_manager.py | 35 +++++++++++- .../jobs/mri_frequency_calibration.py | 45 ++++++++++----- .../orchestrator/repository.py | 6 +- .../orchestrator/sensors.py | 55 +++++++++++++++++-- 7 files changed, 185 insertions(+), 82 deletions(-) delete mode 100644 services/orchestration-engine/orchestrator/hooks.py diff --git a/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py b/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py index ae01916d..bdd5651f 100644 --- a/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py +++ b/services/orchestration-engine/orchestrator/assets/mrpro_direct_reconstruction.py @@ -1,9 +1,12 @@ +# Copyright (C) 2023, BRAIN-LINK UG (haftungsbeschränkt). All Rights Reserved. +# SPDX-License-Identifier: GPL-3.0-only OR LicenseRef-ScanHub-Commercial + +"""MRpro direct image reconstruction using.""" import mrpro from dagster import AssetIn, asset from scanhub_libraries.resources import IDATA_IO_KEY from scanhub_libraries.resources.dag_config import DAGConfiguration -# from orchestrator.hooks import notify_dag_success from orchestrator.io.acquisition_data import AcquisitionData, acquisition_data_asset from orchestrator.io.idata_io_manager import IDataContext @@ -15,12 +18,22 @@ io_manager_key=IDATA_IO_KEY, ) def mrpro_direct_reconstruction(context, data: AcquisitionData, dag_config: DAGConfiguration) -> IDataContext: - """Reconstruct image from a list acquisition results. + """Reconstruct an image from acquisition results using the direct reconstruction method from mrpro. + + Args: + context: The execution context, typically used for logging and runtime information. + data (AcquisitionData): The acquisition data containing the MRD file path and associated metadata. + dag_config (DAGConfiguration): The configuration for the directed acyclic graph (DAG) execution. + + Returns: + IDataContext: An object containing the reconstructed image data and the DAG configuration. + + Workflow: + 1. Loads acquisition results using the DataLakeResource, providing the MRD path and metadata. + 2. Loads an mrpro KData object from the MRD path. + 3. Performs image reconstruction using mrpro's direct reconstruction algorithm. + 4. Returns the reconstructed IData object, which is passed to the idata_io_manager. - 1. Loads acquisition results using the DataLakeRessource providing mrd path and meta data. - 2. Loads mrpro KData object from mrd path - 3. Performs image reconstruction using the direct reconstruction method from mrpro. - 4. Return the reconstructed IData object -> Return is passed to the idata_io_manager. """ mrd_input = data.mrd_path context.log.info("Reading MRD input: %s", str(mrd_input)) diff --git a/services/orchestration-engine/orchestrator/hooks.py b/services/orchestration-engine/orchestrator/hooks.py deleted file mode 100644 index 1062dadb..00000000 --- a/services/orchestration-engine/orchestrator/hooks.py +++ /dev/null @@ -1,53 +0,0 @@ -# # orchestrator/hooks/recon_hooks.py -# from pathlib import Path - -# from dagster import HookContext, success_hook -# from scanhub_libraries.resources import NOTIFIER_KEY - - -# @success_hook(required_resource_keys={NOTIFIER_KEY}) -# def notify_dag_success(context: HookContext) -> None: -# """Check if data has been written and notify backend.""" -# # Expect the asset to return an IDataContext -# result = next(iter(context.op_output_values.values()), None) -# if result is None: -# context.log.warning("NOTIFY-HOOK: No return value; skipping.") -# return -# dag_cfg = getattr(result, "dag_config", None) -# out_dir = getattr(dag_cfg, "output_directory", None) -# if not out_dir: -# context.log.warning("NOTIFY-HOOK: No output_directory; skipping.") -# return - -# p = Path(out_dir) -# files = [str(f) for f in p.glob("**/*") if f.is_file()] -# if not files: -# context.log.warning(f"NOTIFY-HOOK: No files in {p}; skipping.") -# return - -# notifier = getattr(context.resources, NOTIFIER_KEY) -# notifier.send_dag_success(success=True) - - -# @success_hook(required_resource_keys={NOTIFIER_KEY}) -# def notify_device_parameter_update(context: HookContext) -> None: -# """Check if data has been written and notify backend.""" -# # Expect the asset to return an IDataContext -# result = next(iter(context.op_output_values.values()), None) -# if result is None: -# context.log.warning("NOTIFY-HOOK: No return value; skipping.") -# return -# dag_cfg = getattr(result, "dag_config", None) -# out_dir = getattr(dag_cfg, "output_directory", None) -# if not out_dir: -# context.log.warning("NOTIFY-HOOK: No output_directory; skipping.") -# return - -# p = Path(out_dir) -# files = [str(f) for f in p.glob("**/*") if f.is_file()] -# if not files: -# context.log.warning(f"NOTIFY-HOOK: No files in {p}; skipping.") -# return - -# notifier = getattr(context.resources, NOTIFIER_KEY) -# notifier.send_dag_success(success=True) diff --git a/services/orchestration-engine/orchestrator/io/acquisition_data.py b/services/orchestration-engine/orchestrator/io/acquisition_data.py index d9c8a37e..b6214dc0 100644 --- a/services/orchestration-engine/orchestrator/io/acquisition_data.py +++ b/services/orchestration-engine/orchestrator/io/acquisition_data.py @@ -1,3 +1,6 @@ +# Copyright (C) 2023, BRAIN-LINK UG (haftungsbeschränkt). All Rights Reserved. +# SPDX-License-Identifier: GPL-3.0-only OR LicenseRef-ScanHub-Commercial + """Definition of acquisiton data assets.""" from dataclasses import dataclass from pathlib import Path @@ -17,7 +20,16 @@ class AcquisitionData: def _load_acquisition_data(dag_config: DAGConfiguration, data_lake: DataLakeResource) -> AcquisitionData: - """Load acquisition data.""" + """Load acquisition data required for processing. + + Args: + dag_config (DAGConfiguration): The configuration object containing DAG and input file information. + data_lake (DataLakeResource): The data lake resource used to retrieve file paths and device parameters. + + Returns: + AcquisitionData: An object containing the MRD file path, device ID, and device parameters. + + """ mrd_file = data_lake.get_mrd_path(dag_config.input_files) device_id, device_parameter = data_lake.get_device_parameter(dag_config.input_files) return AcquisitionData(mrd_path=mrd_file, device_id=device_id, device_parameter=device_parameter) @@ -32,7 +44,24 @@ def acquisition_data_asset( dag_config: DAGConfiguration, data_lake: DataLakeResource, ) -> AcquisitionData: - """Define acquisition result asset.""" + """Define an acquisition result asset for the orchestration engine. + + This function loads acquisition data using the provided DAG configuration and data lake resource, + logs device parameters, and attaches relevant metadata for visualization in the Dagster UI. + + Args: + context (AssetExecutionContext): The execution context for the asset, used for logging and metadata. + dag_config (DAGConfiguration): The configuration object for the current DAG execution. + data_lake (DataLakeResource): The data lake resource used to load acquisition data. + + Returns: + AcquisitionData: The loaded acquisition data object containing device information and parameters. + + Side Effects: + - Logs device parameters to the context logger. + - Adds output metadata to the context for Dagster UI visualization. + + """ data = _load_acquisition_data(dag_config, data_lake) context.log.info("Parameters for device id %s: %s", data.device_id, data.device_parameter) @@ -55,7 +84,20 @@ def acquisition_data_op( dag_config: DAGConfiguration, data_lake: DataLakeResource, ) -> AcquisitionData: - """Op version of the loader so the job can run end-to-end without assets.""" + """Execute the acquisition data operation. + + Loads acquisition data from the data lake based on the provided DAG configuration. + Logs the MRD file path and device parameters. + + Args: + context (OpExecutionContext): The execution context for the operation, used for logging and runtime information. + dag_config (DAGConfiguration): The configuration object for the DAG, containing parameters for the acquisition. + data_lake (DataLakeResource): The data lake resource used to access acquisition data. + + Returns: + AcquisitionData: The loaded acquisition data object containing MRD file path, device ID, and device parameters. + + """ data = _load_acquisition_data(dag_config, data_lake) context.log.info("MRD file path: %s", str(data.mrd_path)) context.log.info("Parameters for device id %s: %s", data.device_id, data.device_parameter) diff --git a/services/orchestration-engine/orchestrator/io/idata_io_manager.py b/services/orchestration-engine/orchestrator/io/idata_io_manager.py index ea1a6a64..d0bb0922 100644 --- a/services/orchestration-engine/orchestrator/io/idata_io_manager.py +++ b/services/orchestration-engine/orchestrator/io/idata_io_manager.py @@ -1,3 +1,6 @@ +# Copyright (C) 2023, BRAIN-LINK UG (haftungsbeschränkt). All Rights Reserved. +# SPDX-License-Identifier: GPL-3.0-only OR LicenseRef-ScanHub-Commercial + """Dagster IO Manager for mrpro IData object.""" from dataclasses import dataclass from pathlib import Path @@ -19,7 +22,20 @@ class IDataIOManager(ConfigurableIOManager): """IO manager for mrpro IData object.""" def handle_output(self, context: OutputContext, obj: IDataContext) -> None: - """Write idata to dicom folder.""" + """Handle the output of a data processing step by saving the data as DICOM files. + + Args: + context (OutputContext): The context object providing logging and metadata methods. + obj (IDataContext): The data context containing the data and DAG configuration. + + Raises: + AttributeError: If the output directory is not defined in the DAG configuration. + + Side Effects: + - Saves the data as DICOM files in the specified output directory. + - Adds output metadata including the output directory path and list of stored files. + + """ # Decide where to write based on asset key if not obj.dag_config.output_directory: context.log.error("Output directory not defined") @@ -36,7 +52,22 @@ def handle_output(self, context: OutputContext, obj: IDataContext) -> None: }) def load_input(self, context: InputContext) -> IData: - """Read idata from dicom folder.""" + """Load input DICOM data from a specified directory using metadata from the provided context. + + Args: + context (InputContext): The context containing metadata and logging utilities. + + Returns: + IData: An IData instance loaded from the DICOM folder. + + Raises: + AttributeError: If required metadata or the output directory is missing. + FileNotFoundError: If the specified DICOM folder does not exist. + + Logs: + Errors are logged if metadata or directory information is missing or invalid. + + """ if (meta := context.metadata) is None: context.log.error("No metadata, cannot save result") raise AttributeError diff --git a/services/orchestration-engine/orchestrator/jobs/mri_frequency_calibration.py b/services/orchestration-engine/orchestrator/jobs/mri_frequency_calibration.py index 1376ada6..b649afd6 100644 --- a/services/orchestration-engine/orchestrator/jobs/mri_frequency_calibration.py +++ b/services/orchestration-engine/orchestrator/jobs/mri_frequency_calibration.py @@ -1,3 +1,6 @@ +# Copyright (C) 2023, BRAIN-LINK UG (haftungsbeschränkt). All Rights Reserved. +# SPDX-License-Identifier: GPL-3.0-only OR LicenseRef-ScanHub-Commercial + """MRI calibration assets.""" import ismrmrd import numpy as np @@ -14,13 +17,28 @@ def frequency_calibration_op( context: OpExecutionContext, data: AcquisitionData, dag_config: DAGConfiguration, - notigier_device_manager: DeviceManagerNotifier, + notifier_device_manager: DeviceManagerNotifier, ) -> None: - """Perform frequency calibration.""" - context.log.info("AcquisitionData: %s", data) + """Perform frequency calibration on MRI acquisition data. + + This operation reads an ISMRMRD file containing MRI acquisition data, computes the frequency offset using FFT, + calculates the signal-to-noise ratio (SNR) and full width at half maximum (FWHM), and updates the device's + larmor frequency parameter accordingly. The updated parameter is sent to the device manager notifier. + + Args: + context (OpExecutionContext): The execution context for logging and operation metadata. + data (AcquisitionData): The acquisition data containing device parameters and the path to the ISMRMRD file. + dag_config (DAGConfiguration): The DAG configuration, including user access token. + notifier_device_manager (DeviceManagerNotifier): The notifier for sending device parameter updates. + + Returns: + None + + """ parameter = data.device_parameter context.log.info(f"Received device parameters (device id: {data.device_id}): {parameter}") + # Read ismrmrd file provided by acquisition data asset with ismrmrd.File(data.mrd_path, "r") as f: ds = f[list(f.keys())[0]] header = ds.header @@ -58,24 +76,25 @@ def frequency_calibration_op( # Correct device parameter parameter["larmor_frequency"] = f0 + freq_offset - notigier_device_manager.send_device_parameter_update( + notifier_device_manager.send_device_parameter_update( device_id=data.device_id, access_token=dag_config.user_access_token, parameter=parameter, ) -# @op(ins={"start": In(Nothing), }) -# def success_notify_op(context: OpExecutionContext, scanhub_notifier: BackendNotifier) -> None: -# """Single, end-of-run success signal.""" -# scanhub_notifier.send_dag_success() -# context.log.info("Backend notification send.") - @job( name="frequency_calibration", description="Read ismrmrd data and adjust Larmor frequency.", ) def frequency_calibration_job() -> None: - """Calibrate Larmor frequency.""" + """Perform the Larmor frequency calibration job. + + This function loads acquisition data and then runs the frequency calibration + operation using loaded ISMRMRD data. It is intended to be used as a job in the + orchestration engine for MRI frequency calibration. + + Returns: + None + + """ data = acquisition_data_op() frequency_calibration_op(data) - # done = frequency_calibration_op(data) - # success_notify_op(start=done) diff --git a/services/orchestration-engine/orchestrator/repository.py b/services/orchestration-engine/orchestrator/repository.py index e8b45ad3..3d466715 100644 --- a/services/orchestration-engine/orchestrator/repository.py +++ b/services/orchestration-engine/orchestrator/repository.py @@ -1,3 +1,6 @@ +# Copyright (C) 2023, BRAIN-LINK UG (haftungsbeschränkt). All Rights Reserved. +# SPDX-License-Identifier: GPL-3.0-only OR LicenseRef-ScanHub-Commercial + """Define dagster repository.""" import os @@ -5,7 +8,7 @@ from scanhub_libraries.resources import DAG_CONFIG_KEY, DATA_LAKE_KEY, IDATA_IO_KEY, NOTIFIER_DM_KEY, NOTIFIER_WM_KEY from scanhub_libraries.resources.dag_config import DAGConfiguration from scanhub_libraries.resources.data_lake import DataLakeResource -from scanhub_libraries.resources.notifier import WorkflowManagerNotifier, DeviceManagerNotifier +from scanhub_libraries.resources.notifier import DeviceManagerNotifier, WorkflowManagerNotifier from orchestrator.assets.mrpro_direct_reconstruction import mrpro_direct_reconstruction from orchestrator.io.acquisition_data import acquisition_data_asset @@ -17,6 +20,7 @@ DEVICE_MANAGER_URI = "http://device-manager:8000/api/v1/device" WORKFLOW_MANAGER_URI = "http://workflow-manager:8000/api/v1/workflowmanager" + assets = [ acquisition_data_asset, mrpro_direct_reconstruction, diff --git a/services/orchestration-engine/orchestrator/sensors.py b/services/orchestration-engine/orchestrator/sensors.py index cfa7cce2..88bc06aa 100644 --- a/services/orchestration-engine/orchestrator/sensors.py +++ b/services/orchestration-engine/orchestrator/sensors.py @@ -1,3 +1,6 @@ +# Copyright (C) 2023, BRAIN-LINK UG (haftungsbeschränkt). All Rights Reserved. +# SPDX-License-Identifier: GPL-3.0-only OR LicenseRef-ScanHub-Commercial + """Sensors to notify workflow manager dependent on run status.""" from dagster import DagsterRunStatus, DefaultSensorStatus, RunStatusSensorContext, run_status_sensor from scanhub_libraries.resources import DAG_CONFIG_KEY @@ -5,7 +8,15 @@ def _get_dag_config_from_run(context: RunStatusSensorContext) -> dict: - """Get DAG configuration from run status sensor conrext.""" + """Extract the DAG configuration dictionary from a given RunStatusSensorContext. + + Args: + context (RunStatusSensorContext): The context containing the Dagster run information. + + Returns: + dict: The DAG configuration found under the run's resources, or an empty dictionary if not present or invalid. + + """ run_config = getattr(context.dagster_run, "run_config", None) if not isinstance(run_config, dict): return {} @@ -19,7 +30,18 @@ def _get_dag_config_from_run(context: RunStatusSensorContext) -> dict: minimum_interval_seconds=5, ) def on_run_success(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier) -> None: - """Notify workflow manager about successful run using a sensor.""" + """Handle successful DAG run events by notifying the workflow manager if required information is available. + + This function retrieves the DAG configuration from the provided context, extracts the user access token + and output result ID, and attempts to notify the workflow manager of the success. + If either the access token or result ID is missing, it logs an + informational message indicating that the DAG status could not be reported. + + Args: + context (RunStatusSensorContext): The context object containing information about the DAG run. + notifier_workflow_manager (WorkflowManagerNotifier): The notifier used to report DAG run success. + + """ dag_config = _get_dag_config_from_run(context) access_token = dag_config.get("user_access_token", "") result_id = dag_config.get("output_result_id", "") @@ -39,7 +61,21 @@ def on_run_success(context: RunStatusSensorContext, notifier_workflow_manager: W minimum_interval_seconds=5, ) def on_run_failure(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier) -> None: - """Notify workflow manager about failed run using a sensor.""" + """Handle the failure of a DAG run by notifying the workflow manager if possible. + + This function retrieves the DAG configuration from the provided context, extracts the user access token + and output result ID, and attempts to notify the workflow manager of the failure. + If either the access token or result ID is missing, it logs an + informational message indicating that the DAG status could not be reported. + + Args: + context (RunStatusSensorContext): The context object containing information about the DAG run and logging utilities. + notifier_workflow_manager (WorkflowManagerNotifier): The notifier used to send DAG status updates. + + Returns: + None + + """ dag_config = _get_dag_config_from_run(context) access_token = dag_config.get("user_access_token", "") result_id = dag_config.get("output_result_id", "") @@ -58,7 +94,18 @@ def on_run_failure(context: RunStatusSensorContext, notifier_workflow_manager: W minimum_interval_seconds=5, ) def on_run_canceled(context: RunStatusSensorContext, notifier_workflow_manager: WorkflowManagerNotifier) -> None: - """Notify workflow manager about cancelled run using a sensor.""" + """Handle the cancellation of a DAG run by notifying the workflow manager and logging the event. + + This function retrieves the DAG configuration from the provided context, extracts the user access token + and output result ID, and attempts to notify the workflow manager of the cancellation. + If either the access token or result ID is missing, it logs an + informational message indicating that the DAG status could not be reported. + + Args: + context (RunStatusSensorContext): The context object containing information about the current DAG run. + notifier_workflow_manager (WorkflowManagerNotifier): The notifier used to send DAG status updates. + + """ dag_config = _get_dag_config_from_run(context) access_token = dag_config.get("user_access_token", "") result_id = dag_config.get("output_result_id", "") From 579138fe7b32af7c571467f2d9e91e327760b9ee Mon Sep 17 00:00:00 2001 From: David Schote Date: Fri, 5 Sep 2025 08:55:42 +0200 Subject: [PATCH 08/57] Fixed sensors --- services/orchestration-engine/orchestrator/sensors.py | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/services/orchestration-engine/orchestrator/sensors.py b/services/orchestration-engine/orchestrator/sensors.py index 88bc06aa..afc31637 100644 --- a/services/orchestration-engine/orchestrator/sensors.py +++ b/services/orchestration-engine/orchestrator/sensors.py @@ -55,7 +55,7 @@ def on_run_success(context: RunStatusSensorContext, notifier_workflow_manager: W @run_status_sensor( - run_status=DagsterRunStatus.SUCCESS, + run_status=DagsterRunStatus.FAILURE, default_status=DefaultSensorStatus.RUNNING, monitor_all_code_locations=True, minimum_interval_seconds=5, @@ -69,8 +69,10 @@ def on_run_failure(context: RunStatusSensorContext, notifier_workflow_manager: W informational message indicating that the DAG status could not be reported. Args: - context (RunStatusSensorContext): The context object containing information about the DAG run and logging utilities. - notifier_workflow_manager (WorkflowManagerNotifier): The notifier used to send DAG status updates. + context (RunStatusSensorContext): + The context object containing information about the DAG run and logging utilities. + notifier_workflow_manager (WorkflowManagerNotifier): + The notifier used to send DAG status updates. Returns: None @@ -88,7 +90,7 @@ def on_run_failure(context: RunStatusSensorContext, notifier_workflow_manager: W context.log.info("Run failed, but can not report DAG status, missing access_token and/or result_id.") @run_status_sensor( - run_status=DagsterRunStatus.SUCCESS, + run_status=DagsterRunStatus.CANCELED, default_status=DefaultSensorStatus.RUNNING, monitor_all_code_locations=True, minimum_interval_seconds=5, From d5fa369c21b16e660d7cedb334ca3a86f20dffbc Mon Sep 17 00:00:00 2001 From: David Schote Date: Sun, 7 Sep 2025 21:03:44 +0200 Subject: [PATCH 09/57] Fixed device parameter notifier endpoint --- .../src/scanhub_libraries/resources/notifier.py | 2 +- services/workflow-manager/app/api/manager_endpoints.py | 7 ------- 2 files changed, 1 insertion(+), 8 deletions(-) diff --git a/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py b/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py index 416030a3..b09d90b6 100644 --- a/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py +++ b/services/base/shared_libs/src/scanhub_libraries/resources/notifier.py @@ -28,7 +28,7 @@ class DeviceManagerNotifier(ConfigurableResource): def send_device_parameter_update(self, device_id: str, access_token: str, parameter: dict) -> None: """Notify backend about device parameter update and send parameters.""" headers = {"Authorization": "Bearer " + access_token} - url = self.base_url.rstrip("/") + f"/{device_id}" + url = self.base_url.rstrip("/") + f"/parameter/{device_id}" with httpx.Client(timeout=self.timeout) as client: response = client.put(url, json=parameter, headers=headers) response.raise_for_status() diff --git a/services/workflow-manager/app/api/manager_endpoints.py b/services/workflow-manager/app/api/manager_endpoints.py index 1e5ea429..ebb4b759 100644 --- a/services/workflow-manager/app/api/manager_endpoints.py +++ b/services/workflow-manager/app/api/manager_endpoints.py @@ -11,7 +11,6 @@ """ import json -import logging import os from pathlib import Path from typing import Annotated, Any @@ -38,7 +37,6 @@ ) from scanhub_libraries.resources import DAG_CONFIG_KEY from scanhub_libraries.resources.dag_config import DAGConfiguration -# from scanhub_libraries.resources.notifier import BackendNotifier from scanhub_libraries.security import get_current_user from scanhub_libraries.utils import calc_age_from_date @@ -229,11 +227,6 @@ def handle_dag_task_trigger( user_access_token=access_token, output_result_id=str(new_result_out.id), ), - # NOTIFIER_KEY: BackendNotifier( - # success_callback_url=f"{WORKFLOW_MANAGER_URI}/result_ready/{new_result_out.id}", - # devicemanager_url=f"{DEVICE_MANAGER_URI}/parameter/", - # access_token=access_token, - # ), } try: From e8e4224f8cad23a18caa022f55c5c184e7208f80 Mon Sep 17 00:00:00 2001 From: David Schote Date: Wed, 8 Oct 2025 08:40:48 +0200 Subject: [PATCH 10/57] Added helper function to compare keys of json object --- scanhub-ui/src/utils/Compare.tsx | 15 +++++++++++++++ 1 file changed, 15 insertions(+) create mode 100644 scanhub-ui/src/utils/Compare.tsx diff --git a/scanhub-ui/src/utils/Compare.tsx b/scanhub-ui/src/utils/Compare.tsx new file mode 100644 index 00000000..e74dff8e --- /dev/null +++ b/scanhub-ui/src/utils/Compare.tsx @@ -0,0 +1,15 @@ +function getAllKeys(obj: any, prefix = ''): string[] { + return Object.keys(obj).flatMap((key) => { + const path = prefix ? `${prefix}.${key}` : key; + if (obj[key] && typeof obj[key] === 'object' && !Array.isArray(obj[key])) { + return getAllKeys(obj[key], path); + } + return path; + }); +} + +export function haveSameKeys(obj1: any, obj2: any): boolean { + const keys1 = getAllKeys(obj1).sort(); + const keys2 = getAllKeys(obj2).sort(); + return JSON.stringify(keys1) === JSON.stringify(keys2); +} \ No newline at end of file From ecd84643d9275cc4848419abdddadb47636afdc0 Mon Sep 17 00:00:00 2001 From: David Schote Date: Wed, 8 Oct 2025 08:41:04 +0200 Subject: [PATCH 11/57] Editing of device parameter --- scanhub-ui/package.json | 1 + scanhub-ui/src/pages/DeviceView.tsx | 127 +++++++++++++++++++--------- scanhub-ui/yarn.lock | 30 +++++++ 3 files changed, 119 insertions(+), 39 deletions(-) diff --git a/scanhub-ui/package.json b/scanhub-ui/package.json index 4d283bd6..9922a51c 100644 --- a/scanhub-ui/package.json +++ b/scanhub-ui/package.json @@ -17,6 +17,7 @@ "@emotion/styled": "^11.11.0", "@fontsource/roboto": "^5.0.5", "@kitware/vtk.js": "32.12.1", + "@monaco-editor/react": "^4.7.0", "@mui/icons-material": "^7.1.1", "@mui/joy": "^5.0.0-beta.52", "@mui/material": "^7.1.1", diff --git a/scanhub-ui/src/pages/DeviceView.tsx b/scanhub-ui/src/pages/DeviceView.tsx index aec5785e..b47722a7 100644 --- a/scanhub-ui/src/pages/DeviceView.tsx +++ b/scanhub-ui/src/pages/DeviceView.tsx @@ -20,6 +20,10 @@ import Sheet from '@mui/joy/Sheet' import LaunchOutlinedIcon from '@mui/icons-material/LaunchOutlined'; import DeleteSharpIcon from '@mui/icons-material/DeleteSharp' import ModalClose from '@mui/joy/ModalClose' +import Textarea from '@mui/joy/Textarea' +import Button from '@mui/joy/Button' +import Editor from "@monaco-editor/react"; +import { Card } from "@mui/joy"; import NotificationContext from '../NotificationContext' import { deviceApi } from '../api' @@ -28,6 +32,8 @@ import { Alerts } from '../interfaces/components.interface' import { DeviceOut } from '../openapi/generated-client/device/api' import DeviceCreateModal from '../components/DeviceCreateModal' import ConfirmDeleteModal from '../components/ConfirmDeleteModal' +import { haveSameKeys } from '../utils/Compare' +import { flexGrow } from '@mui/system' export default function DeviceView() { @@ -36,7 +42,8 @@ export default function DeviceView() { // const [isUpdating, setIsUpdating] = React.useState(false) const [deviceToDelete, setDeviceToDelete] = React.useState(undefined) const [deviceOpen, setDeviceOpen] = React.useState(undefined) - + const [deviceParameterString, setDeviceParameterString] = React.useState(''); + const { data: devices, @@ -73,26 +80,49 @@ export default function DeviceView() { } }) - // const updateMutation = useMutation(async (user) => { - // await userApi - // .updateUserApiV1UserloginUpdateuserPut(user) - // .then(() => { - // console.log('Modified user:', user.username) - // setIsUpdating(false) - // refetch() - // }) - // .catch((err) => { - // let errorMessage = null - // if (err?.response?.data?.detail) { - // errorMessage = 'Could not update user. Detail: ' + err.response.data.detail - // } else { - // errorMessage = 'Could not update user.' - // } - // setIsUpdating(false) - // refetch() - // showNotification({message: errorMessage, type: 'warning'}) - // }) - // }) + const deviceParameterMutation = useMutation({ + mutationFn: async ({ deviceId, parameter }) => { + await deviceApi.updateDeviceParameterApiV1DeviceParameterDeviceIdPut(deviceId, parameter) + .then(() => { + console.log('Modified device parameter:', parameter) + refetch() + }) + .catch((err) => { + let errorMessage = null + if (err?.response?.data?.detail) { + errorMessage = 'Could not update device parameter. Detail: ' + err.response.data.detail + } else { + errorMessage = 'Could not update user.' + } + refetch() + showNotification({ message: errorMessage, type: 'warning' }) + }) + } + }) + + // Set device parameter string, when deviceOpen is not undefined + React.useEffect(() => { + if (deviceOpen !== undefined) { + setDeviceParameterString(JSON.stringify(deviceOpen?.parameter, null, 4)) + } + }, [deviceOpen]) + + const handleDeviceParameterSave = () => { + try { + const parsed = JSON.parse(deviceParameterString); + if (parsed && deviceOpen !== undefined) { + if (haveSameKeys(parsed, deviceOpen?.parameter)) { + deviceParameterMutation.mutate({deviceId: deviceOpen.id, parameter: parsed}) + setDeviceOpen(undefined) + } else { + showNotification({ message: 'Invalid key(s): Modify existing values only.', type: 'warning' }) + } + } + } catch (err) { + showNotification({ message: 'Invalid JSON: Please correct the format before saving.', type: 'warning' }) + } + }; + if (isLoading) { return ( @@ -212,8 +242,6 @@ export default function DeviceView() { /> setDeviceOpen(undefined)} sx={{ display: 'flex', justifyContent: 'center', alignItems: 'center' }} @@ -231,23 +259,44 @@ export default function DeviceView() { }} > - Device Parameter - Device Parameter + {/*