From 9513cc7421f79bb66146c0aa2ffcdcaa52abb40d Mon Sep 17 00:00:00 2001 From: Carlos Martinez Date: Tue, 10 Mar 2026 15:18:07 +0000 Subject: [PATCH 1/3] C-STORE validation failure notifies Manage API Via ValidationFailureNotifier#notify --- src/pacs_main.py | 3 +- src/server.py | 10 +++- src/services/dicom/c_store.py | 32 ++++++++++-- .../dicom/validation_failure_notifier.py | 52 +++++++++++++++++++ .../test_end_to_end_relay_to_upload.py | 5 +- .../test_send_c_store_to_gateway.py | 4 +- tests/services/dicom/test_c_store.py | 34 ++++++++++++ tests/test_server.py | 23 +++++--- 8 files changed, 146 insertions(+), 17 deletions(-) create mode 100644 src/services/dicom/validation_failure_notifier.py diff --git a/src/pacs_main.py b/src/pacs_main.py index b0df9320..5dfc7640 100644 --- a/src/pacs_main.py +++ b/src/pacs_main.py @@ -29,8 +29,9 @@ def main(): pacs_port = int(os.getenv("PACS_PORT", "4244")) pacs_storage_path = os.getenv("PACS_STORAGE_PATH", "/var/lib/pacs/storage") pacs_db_path = os.getenv("PACS_DB_PATH", "/var/lib/pacs/pacs.db") + mwl_db_path = os.getenv("MWL_DB_PATH", "/var/lib/pacs/worklist.db") - pacs_server = PACSServer(pacs_aet, pacs_port, pacs_storage_path, pacs_db_path, block=True) + pacs_server = PACSServer(pacs_aet, pacs_port, pacs_storage_path, pacs_db_path, block=True, mwl_db_path=mwl_db_path) try: pacs_server.start() diff --git a/src/server.py b/src/server.py index ece83807..4d33480c 100644 --- a/src/server.py +++ b/src/server.py @@ -15,6 +15,7 @@ from services.dicom.c_echo import CEcho from services.dicom.c_store import CStore +from services.dicom.validation_failure_notifier import ValidationFailureNotifier from services.mwl.c_find import CFind from services.mwl.n_create import NCreate from services.mwl.n_set import NSet @@ -33,6 +34,7 @@ def __init__( storage_path: str = "/var/lib/pacs/storage", db_path: str = "/var/lib/pacs/pacs.db", block: bool = True, + mwl_db_path: str = "/var/lib/pacs/worklist.db", ): """ Initialize PACS server. @@ -42,10 +44,13 @@ def __init__( port: Port to listen on storage_path: Directory for DICOM file storage db_path: Path to SQLite database + mwl_db_path: Path to the MWL SQLite database (for failure notification lookups) """ self.ae_title = ae_title self.port = port self.storage = PACSStorage(db_path, storage_path) + self.mwl_storage = MWLStorage(mwl_db_path) + self.notifier = ValidationFailureNotifier() self.ae = None self.block = block @@ -56,7 +61,10 @@ def start(self): self.ae = AE(ae_title=self.ae_title) self.ae.supported_contexts = StoragePresentationContexts - handlers = [(evt.EVT_C_ECHO, CEcho().call), (evt.EVT_C_STORE, CStore(self.storage).call)] + handlers = [ + (evt.EVT_C_ECHO, CEcho().call), + (evt.EVT_C_STORE, CStore(self.storage, mwl_storage=self.mwl_storage, notifier=self.notifier).call), + ] logger.info(f"PACS server listening on 0.0.0.0:{self.port}") logger.info(f"Storage: {self.storage.storage_root}") diff --git a/src/services/dicom/c_store.py b/src/services/dicom/c_store.py index 899d102d..03562033 100644 --- a/src/services/dicom/c_store.py +++ b/src/services/dicom/c_store.py @@ -10,8 +10,9 @@ from services.dicom import FAILURE, SUCCESS from services.dicom.image_compressor import ImageCompressor +from services.dicom.validation_failure_notifier import ValidationFailureNotifier from services.dicom.validator import DicomValidationError, DicomValidator -from services.storage import InstanceExistsError, PACSStorage +from services.storage import InstanceExistsError, MWLStorage, PACSStorage logger = logging.getLogger(__name__) @@ -27,10 +28,14 @@ def __init__( storage: PACSStorage, compressor: ImageCompressor | None = None, validator: DicomValidator | None = None, + mwl_storage: MWLStorage | None = None, + notifier: ValidationFailureNotifier | None = None, ): self.storage = storage self.compressor = compressor or ImageCompressor() self.validator = validator or DicomValidator() + self.mwl_storage = mwl_storage + self.notifier = notifier def call(self, event: Event) -> int: try: @@ -42,24 +47,27 @@ def call(self, event: Event) -> int: return FAILURE sop_instance_uid = ds.get("SOPInstanceUID", "") + accession_number = ds.get("AccessionNumber", "") + patient_id = ds.get("PatientID") + patient_name = str(ds.get("PatientName", "")) + if not sop_instance_uid: logger.error("Missing SOPInstanceUID") + self._notify_failure(accession_number, "Missing SOPInstanceUID") return FAILURE - patient_id = ds.get("PatientID") if not patient_id: logger.error("Missing PatientID") + self._notify_failure(accession_number, "Missing PatientID") return FAILURE - accession_number = ds.get("AccessionNumber", "") - patient_name = str(ds.get("PatientName", "")) - # Validate dataset before compression try: self.validator.validate_dataset(ds) self.validator.validate_pixel_data(ds) except DicomValidationError as e: logger.error(f"DICOM validation failed: {e}") + self._notify_failure(accession_number, f"DICOM validation failed: {e}") return FAILURE # Compress dataset before storing @@ -71,6 +79,7 @@ def call(self, event: Event) -> int: self.validator.validate_bytes(dicom_bytes) except DicomValidationError as e: logger.error(f"Serialized DICOM invalid: {e}") + self._notify_failure(accession_number, f"Serialized DICOM invalid: {e}") return FAILURE self.storage.store_instance( @@ -100,3 +109,16 @@ def dataset_to_bytes(self, ds: Dataset) -> bytes: dcmwrite(buffer, ds, enforce_file_format=True) buffer.seek(0) return buffer.read() + + def _notify_failure(self, accession_number: str, error: str) -> None: + if not self.mwl_storage or not self.notifier: + return + + source_message_id = self.mwl_storage.get_source_message_id(accession_number) + if not source_message_id: + logger.warning( + f"Cannot report validation failure: no worklist item found for accession {accession_number!r}" + ) + return + + self.notifier.notify(source_message_id, error) diff --git a/src/services/dicom/validation_failure_notifier.py b/src/services/dicom/validation_failure_notifier.py new file mode 100644 index 00000000..62128db4 --- /dev/null +++ b/src/services/dicom/validation_failure_notifier.py @@ -0,0 +1,52 @@ +"""Notifier for DICOM C-STORE validation failures. + +Reports validation failures to the Manage Breast Screening HTTP API. +""" + +import logging +import os + +import requests + +logger = logging.getLogger(__name__) + + +class ValidationFailureNotifier: + def __init__(self, api_endpoint: str | None = None, timeout: int = 30, verify_ssl: bool = True): + self.api_endpoint = api_endpoint or os.getenv("CLOUD_API_ENDPOINT", "http://localhost:8000/api/v1/dicom") + self.timeout = timeout + self.verify_ssl = verify_ssl + + def headers(self) -> dict: + return { + "Authorization": f"Bearer {os.getenv('CLOUD_API_TOKEN', '')}", + } + + def notify(self, source_message_id: str, error: str) -> bool: + try: + logger.info(f"Reporting validation failure for action {source_message_id}") + + response = requests.post( + f"{self.api_endpoint}/{source_message_id}/failure", + json={"error": error}, + timeout=self.timeout, + verify=self.verify_ssl, + headers=self.headers(), + ) + + if response.status_code == 200: + logger.info(f"Validation failure reported for action {source_message_id}") + return True + else: + logger.error( + f"Failed to report validation failure for {source_message_id}: " + f"status {response.status_code}, body: {response.text}" + ) + return False + + except requests.exceptions.Timeout: + logger.error(f"Timeout reporting validation failure for {source_message_id} after {self.timeout}s") + return False + except requests.exceptions.RequestException as e: + logger.error(f"Error reporting validation failure for {source_message_id}: {e}", exc_info=True) + return False diff --git a/tests/integration/test_end_to_end_relay_to_upload.py b/tests/integration/test_end_to_end_relay_to_upload.py index 2828ea2c..89df791f 100644 --- a/tests/integration/test_end_to_end_relay_to_upload.py +++ b/tests/integration/test_end_to_end_relay_to_upload.py @@ -19,6 +19,7 @@ from services.dicom import PENDING, SUCCESS from services.dicom.dicom_uploader import DICOMUploader from services.dicom.upload_processor import UploadProcessor +from services.dicom.validation_failure_notifier import ValidationFailureNotifier from services.storage import MWLStorage, PACSStorage TEST_ACCESSION_NUMBER = "ACC-E2E-12345" # gitleaks:allow @@ -85,12 +86,14 @@ def mwl_server(self, mwl_storage): return server @pytest.fixture - def pacs_server(self, pacs_storage): + def pacs_server(self, pacs_storage, mwl_storage): """PACS server using the shared storage.""" server = PACSServer.__new__(PACSServer) server.ae_title = "SCREENING_PACS" server.port = 4244 server.storage = pacs_storage + server.mwl_storage = mwl_storage + server.notifier = ValidationFailureNotifier() server.ae = None server.block = False return server diff --git a/tests/integration/test_send_c_store_to_gateway.py b/tests/integration/test_send_c_store_to_gateway.py index 707e0e5f..8e7c69ac 100644 --- a/tests/integration/test_send_c_store_to_gateway.py +++ b/tests/integration/test_send_c_store_to_gateway.py @@ -12,7 +12,9 @@ class TestSendCStoreToGateway: @pytest.fixture(autouse=True) def with_pacs_server(self, tmp_dir): - server = PACSServer("SCREENING_PACS", 4244, tmp_dir, f"{tmp_dir}/test.db", block=False) + server = PACSServer( + "SCREENING_PACS", 4244, tmp_dir, f"{tmp_dir}/test.db", block=False, mwl_db_path=f"{tmp_dir}/worklist.db" + ) server.start() yield diff --git a/tests/services/dicom/test_c_store.py b/tests/services/dicom/test_c_store.py index 679cd0fd..11e0df5d 100644 --- a/tests/services/dicom/test_c_store.py +++ b/tests/services/dicom/test_c_store.py @@ -8,6 +8,9 @@ from services.dicom import FAILURE, SUCCESS from services.dicom.c_store import CStore from services.dicom.image_compressor import ImageCompressor +from services.dicom.validation_failure_notifier import ValidationFailureNotifier +from services.dicom.validator import DicomValidationError, DicomValidator +from services.storage import MWLStorage class TestCStore: @@ -104,3 +107,34 @@ def test_compression_applied_on_storage(self, mock_storage, mock_event): stored_bytes = mock_storage.store_instance.call_args[0][1] stored_ds = pydicom.dcmread(BytesIO(stored_bytes), force=True) assert stored_ds.file_meta.TransferSyntaxUID == JPEG2000 + + def test_validation_failure_notifies_manage(self, mock_storage, mock_event): + """When validation fails and accession is in MWL, notify manage.""" + mock_validator = Mock(spec=DicomValidator) + mock_validator.validate_dataset.side_effect = DicomValidationError("Missing required tag") + + mock_mwl = Mock(spec=MWLStorage) + mock_mwl.get_source_message_id.return_value = "action-uuid-123" + + mock_notifier = Mock(spec=ValidationFailureNotifier) + + subject = CStore(mock_storage, validator=mock_validator, mwl_storage=mock_mwl, notifier=mock_notifier) + assert subject.call(mock_event) == FAILURE + + mock_notifier.notify.assert_called_once_with("action-uuid-123", "DICOM validation failed: Missing required tag") + mock_mwl.get_source_message_id.assert_called_once_with("ABC123") + + def test_validation_failure_accession_not_in_mwl(self, mock_storage, mock_event): + """When accession is not in MWL, validation failure returns FAILURE without calling notify.""" + mock_validator = Mock(spec=DicomValidator) + mock_validator.validate_dataset.side_effect = DicomValidationError("Missing required tag") + + mock_mwl = Mock(spec=MWLStorage) + mock_mwl.get_source_message_id.return_value = None + + mock_notifier = Mock(spec=ValidationFailureNotifier) + + subject = CStore(mock_storage, validator=mock_validator, mwl_storage=mock_mwl, notifier=mock_notifier) + assert subject.call(mock_event) == FAILURE + + mock_notifier.notify.assert_not_called() diff --git a/tests/test_server.py b/tests/test_server.py index 173a492a..23a76197 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -10,34 +10,41 @@ from server import MWLServer, PACSServer +@patch(f"{PACSServer.__module__}.MWLStorage") @patch(f"{PACSServer.__module__}.PACSStorage") class TestPACSServer: - def test_init(self, mock_storage, tmp_dir): - subject = PACSServer("Custom AE Title", 2222, tmp_dir, f"{tmp_dir}/test.db", False) + def test_init(self, mock_pacs_storage, mock_mwl_storage, tmp_dir): + subject = PACSServer( + "Custom AE Title", 2222, tmp_dir, f"{tmp_dir}/test.db", False, mwl_db_path=f"{tmp_dir}/worklist.db" + ) assert subject.ae_title == "Custom AE Title" assert subject.port == 2222 - assert subject.storage == mock_storage.return_value + assert subject.storage == mock_pacs_storage.return_value + assert subject.mwl_storage == mock_mwl_storage.return_value assert subject.ae is None assert subject.block is False - mock_storage.assert_called_once_with(f"{tmp_dir}/test.db", tmp_dir) + mock_pacs_storage.assert_called_once_with(f"{tmp_dir}/test.db", tmp_dir) + mock_mwl_storage.assert_called_once_with(f"{tmp_dir}/worklist.db") - def test_init_defaults(self, mock_storage): + def test_init_defaults(self, mock_pacs_storage, mock_mwl_storage): subject = PACSServer() assert subject.ae_title == "SCREENING_PACS" assert subject.port == 4244 - assert subject.storage == mock_storage.return_value + assert subject.storage == mock_pacs_storage.return_value + assert subject.mwl_storage == mock_mwl_storage.return_value assert subject.ae is None assert subject.block is True - mock_storage.assert_called_once_with("/var/lib/pacs/pacs.db", "/var/lib/pacs/storage") + mock_pacs_storage.assert_called_once_with("/var/lib/pacs/pacs.db", "/var/lib/pacs/storage") + mock_mwl_storage.assert_called_once_with("/var/lib/pacs/worklist.db") @patch(f"{PACSServer.__module__}.AE") @patch(f"{PACSServer.__module__}.CEcho") @patch(f"{PACSServer.__module__}.CStore") - def test_start(self, mock_c_store, mock_c_echo, mock_ae, _): + def test_start(self, mock_c_store, mock_c_echo, mock_ae, _mock_pacs_storage, _mock_mwl_storage): subject = PACSServer() subject.start() From a7f0aa7a3354d90b886b2fa827262f3e5804aaf5 Mon Sep 17 00:00:00 2001 From: Carlos Martinez Date: Wed, 11 Mar 2026 16:23:31 +0000 Subject: [PATCH 2/3] Use PATCH not POST --- src/services/dicom/validation_failure_notifier.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/services/dicom/validation_failure_notifier.py b/src/services/dicom/validation_failure_notifier.py index 62128db4..b8bc2822 100644 --- a/src/services/dicom/validation_failure_notifier.py +++ b/src/services/dicom/validation_failure_notifier.py @@ -26,7 +26,7 @@ def notify(self, source_message_id: str, error: str) -> bool: try: logger.info(f"Reporting validation failure for action {source_message_id}") - response = requests.post( + response = requests.patch( f"{self.api_endpoint}/{source_message_id}/failure", json={"error": error}, timeout=self.timeout, From 0245d213618a7eb0a30b5db6a168e8fd9df19071 Mon Sep 17 00:00:00 2001 From: Carlos Martinez Date: Wed, 11 Mar 2026 16:40:21 +0000 Subject: [PATCH 3/3] Instantiate ValidationFailureNotifier at point of use --- src/server.py | 4 +--- src/services/dicom/c_store.py | 2 +- tests/integration/test_end_to_end_relay_to_upload.py | 2 -- 3 files changed, 2 insertions(+), 6 deletions(-) diff --git a/src/server.py b/src/server.py index 4d33480c..1f327a60 100644 --- a/src/server.py +++ b/src/server.py @@ -15,7 +15,6 @@ from services.dicom.c_echo import CEcho from services.dicom.c_store import CStore -from services.dicom.validation_failure_notifier import ValidationFailureNotifier from services.mwl.c_find import CFind from services.mwl.n_create import NCreate from services.mwl.n_set import NSet @@ -50,7 +49,6 @@ def __init__( self.port = port self.storage = PACSStorage(db_path, storage_path) self.mwl_storage = MWLStorage(mwl_db_path) - self.notifier = ValidationFailureNotifier() self.ae = None self.block = block @@ -63,7 +61,7 @@ def start(self): handlers = [ (evt.EVT_C_ECHO, CEcho().call), - (evt.EVT_C_STORE, CStore(self.storage, mwl_storage=self.mwl_storage, notifier=self.notifier).call), + (evt.EVT_C_STORE, CStore(self.storage, mwl_storage=self.mwl_storage).call), ] logger.info(f"PACS server listening on 0.0.0.0:{self.port}") diff --git a/src/services/dicom/c_store.py b/src/services/dicom/c_store.py index 03562033..acfe45ad 100644 --- a/src/services/dicom/c_store.py +++ b/src/services/dicom/c_store.py @@ -35,7 +35,7 @@ def __init__( self.compressor = compressor or ImageCompressor() self.validator = validator or DicomValidator() self.mwl_storage = mwl_storage - self.notifier = notifier + self.notifier = notifier or ValidationFailureNotifier() def call(self, event: Event) -> int: try: diff --git a/tests/integration/test_end_to_end_relay_to_upload.py b/tests/integration/test_end_to_end_relay_to_upload.py index 89df791f..a699041c 100644 --- a/tests/integration/test_end_to_end_relay_to_upload.py +++ b/tests/integration/test_end_to_end_relay_to_upload.py @@ -19,7 +19,6 @@ from services.dicom import PENDING, SUCCESS from services.dicom.dicom_uploader import DICOMUploader from services.dicom.upload_processor import UploadProcessor -from services.dicom.validation_failure_notifier import ValidationFailureNotifier from services.storage import MWLStorage, PACSStorage TEST_ACCESSION_NUMBER = "ACC-E2E-12345" # gitleaks:allow @@ -93,7 +92,6 @@ def pacs_server(self, pacs_storage, mwl_storage): server.port = 4244 server.storage = pacs_storage server.mwl_storage = mwl_storage - server.notifier = ValidationFailureNotifier() server.ae = None server.block = False return server