Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions src/services/dicom/c_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
from services.dicom.image_compressor import ImageCompressor
from services.dicom.validation_failure_notifier import ValidationFailureNotifier
from services.dicom.validator import DicomValidationError, DicomValidator
from services.mwl import MWLStatus
from services.storage import InstanceExistsError, MWLStorage, PACSStorage

logger = logging.getLogger(__name__)
Expand Down Expand Up @@ -79,6 +80,7 @@ def call(self, event: Event) -> int:
},
event.assoc.requestor.ae_title,
)
self._mark_in_progress(accession_number)
Comment thread
steventux marked this conversation as resolved.
return SUCCESS

except InstanceExistsError:
Expand All @@ -97,6 +99,14 @@ def dataset_to_bytes(self, ds: Dataset) -> bytes:
buffer.seek(0)
return buffer.read()

def _mark_in_progress(self, accession_number: str) -> None:
if not self.mwl_storage or not accession_number:
return
try:
self.mwl_storage.update_status(accession_number, MWLStatus.IN_PROGRESS.value)
except Exception as e:
logger.error(f"Failed to mark worklist item in progress: {e}", exc_info=True)

def _notify_failure(self, accession_number: str, error: str) -> None:
if not self.mwl_storage or not self.notifier:
return
Expand Down
34 changes: 28 additions & 6 deletions src/services/storage.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
from typing import Dict, List, Optional

from models import WorklistItem
from services.mwl import MWLStatus

logger = logging.getLogger(__name__)

Expand Down Expand Up @@ -284,13 +285,25 @@ class WorklistItemNotFoundError(Exception):
pass


class InvalidStatusTransitionError(Exception):
"""Raised when a requested status transition is not permitted."""

pass


class DuplicateWorklistItemError(Exception):
"""Raised when a worklist item with the same accession number already exists."""

pass


class MWLStorage(Storage):
_STATUS_TRANSITIONS: dict[MWLStatus, MWLStatus] = {
MWLStatus.IN_PROGRESS: MWLStatus.SCHEDULED,
MWLStatus.COMPLETED: MWLStatus.IN_PROGRESS,
MWLStatus.DISCONTINUED: MWLStatus.IN_PROGRESS,
}

def __init__(self, db_path: str = "/var/lib/pacs/worklist.db"):
"""
Initialize Worklist storage.
Expand Down Expand Up @@ -459,16 +472,24 @@ def update_status(
self, accession_number: str, status: str, mpps_instance_uid: Optional[str] = None
) -> Optional[str]:
"""
Update the status of a worklist item.
Transition a worklist item to a new status, enforcing valid state transitions.

Args:
accession_number: The accession number to update
status: New status (SCHEDULED, IN PROGRESS, COMPLETED, DISCONTINUED)
status: Target status
mpps_instance_uid: Optional MPPS instance UID

Returns:
source_message_id if item was updated, None if not found

Raises:
InvalidStatusTransitionError: If the transition is not permitted
"""
target = MWLStatus(status)
if target not in self._STATUS_TRANSITIONS:
raise InvalidStatusTransitionError(f"Cannot transition to '{status}'")
from_status = self._STATUS_TRANSITIONS[target]

with self._get_connection() as conn:
cursor = conn.execute(
"""
Expand All @@ -477,18 +498,19 @@ def update_status(
mpps_instance_uid = COALESCE(?, mpps_instance_uid),
updated_at = CURRENT_TIMESTAMP
WHERE accession_number = ?
""",
(status, mpps_instance_uid, accession_number),
AND status = ?
""",
(status, mpps_instance_uid, accession_number, from_status.value),
)
conn.commit()

if cursor.rowcount == 0:
return None

result = conn.execute(
"SELECT source_message_id FROM worklist_items WHERE accession_number = ?", (accession_number,)
"SELECT source_message_id FROM worklist_items WHERE accession_number = ?",
(accession_number,),
).fetchone()

return result["source_message_id"] if result is not None else None

def update_study_instance_uid(self, accession_number: str, study_instance_uid: str) -> bool:
Expand Down
2 changes: 0 additions & 2 deletions tests/integration/test_c_find_returns_worklist_items.py
Original file line number Diff line number Diff line change
Expand Up @@ -265,8 +265,6 @@ def test_cfind_filters_by_modality(self, event, storage):
source_message_id="MSGID123456",
)
)
storage.update_status("ACC234567", "SCHEDULED")

event.identifier.ScheduledProcedureStepSequence[0].Modality = "MG"

results = list(CFind(storage).call(event))
Expand Down
25 changes: 24 additions & 1 deletion tests/integration/test_c_store_saves_metadata.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,9 @@
DigitalMammographyXRayImageStorageForProcessing,
)

from models import WorklistItem
from services.dicom.c_store import SUCCESS, CStore
from services.storage import PACSStorage
from services.storage import MWLStorage, PACSStorage


@pytest.mark.integration
Expand All @@ -36,6 +37,10 @@ def mock_event(self):
def storage(self, tmp_dir):
return PACSStorage(f"{tmp_dir}/test.db", tmp_dir)

@pytest.fixture
def mwl_storage(self, tmp_dir):
return MWLStorage(f"{tmp_dir}/worklist.db")

def test_existing_sop_instance_uid(self, storage, mock_event):
sop_instance_uid = "1.2.3.4.5.6" # gitleaks:allow
subject = CStore(storage)
Expand Down Expand Up @@ -85,6 +90,24 @@ def test_valid_event_is_stored(self, storage, mock_event):
assert storage_path == "ff/af/ffaff041ab509297.dcm"
assert Path(f"{storage.storage_root}/{storage_path}").is_file()

def test_c_store_marks_worklist_in_progress(self, storage, mwl_storage, mock_event):
item = WorklistItem(
accession_number="ABC123",
modality="MG",
patient_birth_date="19800101",
patient_id="9990001112",
patient_name="JANE^SMITH",
scheduled_date="20240101",
scheduled_time="090000",
)
mwl_storage.store_worklist_item(item)

subject = CStore(storage, mwl_storage=mwl_storage)
assert subject.call(mock_event) == SUCCESS

fetched = mwl_storage.get_worklist_item("ABC123")
assert fetched.status == "IN PROGRESS"

def test_compressed_image_stored_on_filesystem(self, storage, dataset_with_pixels):
"""Verify compressed images are stored with JPEG 2000 transfer syntax."""
# Customize the shared dataset for this test
Expand Down
3 changes: 0 additions & 3 deletions tests/integration/test_end_to_end_relay_to_upload.py
Original file line number Diff line number Diff line change
Expand Up @@ -134,9 +134,6 @@ async def test_full_flow_relay_to_upload(
assert worklist_items[0].accession_number == TEST_ACCESSION_NUMBER
assert worklist_items[0].patient_id == TEST_PATIENT_ID

# Update status to SCHEDULED so it appears in C-FIND results
mwl_storage.update_status(TEST_ACCESSION_NUMBER, "SCHEDULED")

# ===== STEP 2: Query worklist via C-FIND =====
mwl_server.start()
try:
Expand Down
3 changes: 0 additions & 3 deletions tests/integration/test_request_cfind_on_worklist.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,6 @@ def with_pacs_server(self, tmp_dir):
source_message_id="MSGID123456",
)
)
storage.update_status("ACC123456", "SCHEDULED")
storage.store_worklist_item(
WorklistItem(
accession_number="ACC234567",
Expand All @@ -48,8 +47,6 @@ def with_pacs_server(self, tmp_dir):
source_message_id="MSGID234567",
)
)
storage.update_status("ACC234567", "SCHEDULED")

server.start()

yield
Expand Down
24 changes: 24 additions & 0 deletions tests/services/dicom/test_c_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,30 @@ def test_validation_failure_notifies_manage(self, mock_storage, mock_event):
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_worklist_marked_in_progress_on_success(self, mock_storage, mock_event):
mock_mwl = Mock(spec=MWLStorage)
subject = CStore(mock_storage, mwl_storage=mock_mwl)

assert subject.call(mock_event) == SUCCESS

mock_mwl.update_status.assert_called_once_with("ABC123", "IN PROGRESS")

def test_worklist_not_updated_on_store_failure(self, mock_storage, mock_event):
mock_storage.store_instance.side_effect = Exception("store failed")
mock_mwl = Mock(spec=MWLStorage)
subject = CStore(mock_storage, mwl_storage=mock_mwl)

assert subject.call(mock_event) == FAILURE

mock_mwl.update_status.assert_not_called()

def test_worklist_update_error_does_not_fail_store(self, mock_storage, mock_event):
mock_mwl = Mock(spec=MWLStorage)
mock_mwl.update_status.side_effect = Exception("db error")
subject = CStore(mock_storage, mwl_storage=mock_mwl)

assert subject.call(mock_event) == SUCCESS

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)
Expand Down
53 changes: 34 additions & 19 deletions tests/services/test_storage.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
from pydicom.uid import generate_uid

from models import WorklistItem
from services.storage import MWLStorage, PACSStorage, WorklistItemNotFoundError
from services.storage import InvalidStatusTransitionError, MWLStorage, PACSStorage, WorklistItemNotFoundError


@pytest.fixture
Expand Down Expand Up @@ -228,26 +228,19 @@ def test_get_worklist_item_returns_none(self, mwl_storage):

def test_update_status(self, mwl_storage, result):
item = self._insert_item(mwl_storage, result)
mwl_storage.update_status(item.accession_number, "IN PROGRESS")

returned = mwl_storage.update_status(item.accession_number, "COMPLETED")

assert returned == item.source_message_id
assert mwl_storage.get_worklist_item(item.accession_number).status == "COMPLETED"

with mwl_storage._get_connection() as conn:
row = conn.execute(
"SELECT status FROM worklist_items WHERE accession_number = ?",
(item.accession_number,),
).fetchone()

assert row["status"] == "COMPLETED"

def test_update_status_with_no_update(self, mwl_storage):
result = mwl_storage.update_status("DOES_NOT_EXIST", "COMPLETED")

assert result is None
def test_update_status_returns_none_when_not_found(self, mwl_storage):
assert mwl_storage.update_status("DOES_NOT_EXIST", "IN PROGRESS") is None

def test_update_status_with_mpps(self, mwl_storage, result):
item = self._insert_item(mwl_storage, result)
mwl_storage.update_status(item.accession_number, "IN PROGRESS")

returned = mwl_storage.update_status(
item.accession_number,
Expand All @@ -256,14 +249,19 @@ def test_update_status_with_mpps(self, mwl_storage, result):
)

assert returned == item.source_message_id
assert mwl_storage.get_worklist_item(item.accession_number).mpps_instance_uid == "some-uid"

with mwl_storage._get_connection() as conn:
row = conn.execute(
"SELECT mpps_instance_uid FROM worklist_items WHERE accession_number = ?",
(item.accession_number,),
).fetchone()
def test_update_status_raises_on_invalid_target(self, mwl_storage, result):
item = self._insert_item(mwl_storage, result)

assert row["mpps_instance_uid"] == "some-uid"
with pytest.raises(InvalidStatusTransitionError):
mwl_storage.update_status(item.accession_number, "SCHEDULED") # SCHEDULED is never a valid target

def test_update_status_returns_none_on_wrong_state(self, mwl_storage, result):
item = self._insert_item(mwl_storage, result)
mwl_storage.update_status(item.accession_number, "IN PROGRESS")

assert mwl_storage.update_status(item.accession_number, "IN PROGRESS") is None

def test_update_study_instance_uid(self, mwl_storage, result):
item = self._insert_item(mwl_storage, result)
Expand Down Expand Up @@ -294,6 +292,7 @@ def test_delete_worklist_item_raises(self, mwl_storage):
def test_mpps_instance_exists(self, mwl_storage, result):
uid = generate_uid()
item = self._insert_item(mwl_storage, result)
mwl_storage.update_status(item.accession_number, "IN PROGRESS")

mwl_storage.update_status(item.accession_number, "COMPLETED", mpps_instance_uid=uid)

Expand All @@ -305,6 +304,7 @@ def test_mpps_instance_not_exists(self, mwl_storage):
def test_get_worklist_item_by_mpps_instance_uid(self, mwl_storage, result):
uid = generate_uid()
item = self._insert_item(mwl_storage, result)
mwl_storage.update_status(item.accession_number, "IN PROGRESS")

mwl_storage.update_status(item.accession_number, "COMPLETED", mpps_instance_uid=uid)

Expand All @@ -318,3 +318,18 @@ def test_get_worklist_item_by_mpps_instance_uid(self, mwl_storage, result):

def test_get_worklist_item_by_mpps_instance_uid_returns_none(self, mwl_storage):
assert mwl_storage.get_worklist_item_by_mpps_instance_uid("nope") is None

def test_update_status_scheduled_to_in_progress(self, mwl_storage, result):
item = self._insert_item(mwl_storage, result)

mwl_storage.update_status(item.accession_number, "IN PROGRESS")

assert mwl_storage.get_worklist_item(item.accession_number).status == "IN PROGRESS"

def test_update_status_in_progress_to_discontinued(self, mwl_storage, result):
item = self._insert_item(mwl_storage, result)
mwl_storage.update_status(item.accession_number, "IN PROGRESS")

mwl_storage.update_status(item.accession_number, "DISCONTINUED")

assert mwl_storage.get_worklist_item(item.accession_number).status == "DISCONTINUED"
Loading