Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
46 commits
Select commit Hold shift + click to select a range
26ab1ef
Add anon_queue to pixl_core with different message and no priority
ruaridhg Aug 21, 2026
735d67b
Remove async functionality from anon_queue subscriber
ruaridhg Aug 21, 2026
91b891b
Replace POST with anon queue in pixl_imaging/_orthanc.py and update n…
ruaridhg Aug 21, 2026
6c5b362
AnonymisationProducer publishes [message] instead of message
ruaridhg Aug 21, 2026
0e0ee07
Clean up anon_queue subscriber
ruaridhg Aug 21, 2026
48cee58
Add RabbitMQ anon_queue in orthanc-anon plugin
ruaridhg Aug 21, 2026
250f775
Add anonymisation queue to _message_count()
ruaridhg Aug 21, 2026
c485ea2
Merge branch 'main' into 650-anonymisation-reads-from-queue
ruaridhg Aug 21, 2026
0e3c3d1
Added tests for anon_queue and modified subscriber to move away from …
ruaridhg Aug 21, 2026
050ddfb
Remove SERVICE_SETTINGS[rabbitmq] from AnonymisationProducer
ruaridhg Aug 21, 2026
e093a01
Add anonymisation queue to cli tests
ruaridhg Aug 21, 2026
47830a7
Remove anonymisation queue being populated from cli
ruaridhg Aug 21, 2026
124246e
Added __init__.py to both queues tests so differentiate same names
ruaridhg Aug 21, 2026
2dd0829
Renamed anon_queue test queues to fix name conflict
ruaridhg Aug 21, 2026
2f22ef9
Add test to increase coverage including anonymisation in queue_names
ruaridhg Aug 21, 2026
1eb4a9d
Add test to cover new orthanc queue connection
ruaridhg Aug 21, 2026
af4a6e6
Merge branch 'main' into 650-anonymisation-reads-from-queue
ruaridhg Sep 3, 2026
6586344
Combine 2 queues into core.queue
ruaridhg Sep 9, 2026
d5bbd00
Update imports to match new core.queue
ruaridhg Sep 9, 2026
d3e4b3f
Remove ImportStudiesFromRaw and API call
ruaridhg Sep 9, 2026
2e61304
Create models.py for ImagingRequest message and AnonymisationMessage
ruaridhg Sep 9, 2026
93c608f
Add parent_context back to process_anon_message
ruaridhg Sep 9, 2026
5c0d0e1
Combine 2 queues testing into 1 test suite
ruaridhg Sep 9, 2026
323b5ec
Merge branch 'main' into 650-anonymisation-reads-from-queue
ruaridhg Sep 9, 2026
82e161f
Remove redundant anon_queue refs
ruaridhg Sep 10, 2026
36c0e6c
Ensure pixl_core tests pass locally
ruaridhg Sep 10, 2026
8312af5
Update tests that need Message replace w/ ImagingRequestMessage
ruaridhg Sep 10, 2026
4b7ce8f
Update tests that need Message replace w/ ImagingRequestMessage
ruaridhg Sep 10, 2026
942a15f
Update pixl_core/src/core/queue/subscriber.py
ruaridhg Sep 16, 2026
24e118c
Address PR review comments - reduce redundancy
ruaridhg Sep 16, 2026
c145c6b
Rename tests
ruaridhg Sep 16, 2026
fe71ad3
Fix unit tests
ruaridhg Sep 16, 2026
09a171d
Try to fix system test w import studies from raw message context
ruaridhg Sep 17, 2026
5855db6
Merge branch 'main' into 650-anonymisation-reads-from-queue
ruaridhg Sep 17, 2026
664b269
Fix queue-declaration argument mismatch in cli
ruaridhg Sep 17, 2026
b163d1c
Add retry to AnonymisationPixlConsumer
ruaridhg Sep 17, 2026
7d96ef9
Add copilot suggestions
ruaridhg Sep 18, 2026
7dffc13
Remove redundant input for anon message type from async consumer
ruaridhg Sep 18, 2026
9ce7679
Improve test coverage for subscriber covering error handling
ruaridhg Sep 18, 2026
e97ae06
Merge remote-tracking branch 'origin/main' into 650-anonymisation-rea…
ruaridhg Sep 18, 2026
2126628
Merge branch 'main' into 650-anonymisation-reads-from-queue
ruaridhg Sep 22, 2026
1ad2a0d
Simplify max-priority config args and give Anon subscriber max-priority
ruaridhg Sep 22, 2026
7807637
Simplify max-priority config args and give Anon subscriber max-priority
ruaridhg Sep 22, 2026
b03f507
Remove None as an option for parent context
ruaridhg Sep 22, 2026
fb914cb
Replace Message with ImagingRequestMessage in docs
ruaridhg Sep 22, 2026
a79d946
AnonymisationPixlConsumer changed to async so ack does not get sent b…
ruaridhg Sep 23, 2026
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
14 changes: 8 additions & 6 deletions cli/src/pixl_cli/_message_processing.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,9 +20,9 @@

import pandas as pd
import tqdm
from core.patient_queue._base import PixlBlockingInterface
from core.patient_queue.message import Message
from core.patient_queue.producer import PixlProducer
from core.queue._base import PixlBlockingInterface
from core.queue.models import ImagingRequestMessage
from core.queue.producer import PixlProducer
from decouple import config
from loguru import logger

Expand All @@ -35,15 +35,15 @@

def messages_from_df(
df: pd.DataFrame,
) -> list[Message]:
) -> list[ImagingRequestMessage]:
"""
Reads patient information from a DataFrame and transforms that into messages.

:param messages_df: DataFrame containing patient information
"""
messages = []
for _, row in df.iterrows():
message = Message(
message = ImagingRequestMessage(
mrn=row["mrn"],
accession_number=row["accession_number"],
study_uid=row["study_uid"],
Expand Down Expand Up @@ -129,6 +129,8 @@ def _message_count(queues_to_populate: list[str]) -> int:
if "imaging-primary" in queues_to_populate:
queues_to_count.add("imaging-secondary")

queues_to_count.add("anonymisation")

messages_in_queues = 0
for queue in queues_to_count:
with PixlBlockingInterface(queue_name=queue, **SERVICE_SETTINGS["rabbitmq"]) as rabbitmq:
Expand All @@ -139,7 +141,7 @@ def _message_count(queues_to_populate: list[str]) -> int:

def populate_queue_and_db(
queues: list[str], messages_df: pd.DataFrame, messages_priority: int
) -> list[Message]:
) -> list[ImagingRequestMessage]:
"""
Populate queues with messages,
for imaging queue update the database and filter out exported or skipped studies.
Expand Down
2 changes: 1 addition & 1 deletion cli/src/pixl_cli/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@
import click
import requests
from core.exports import ParquetExport
from core.patient_queue.producer import PixlProducer
from core.queue.producer import PixlProducer
from core.telemetry import configure_logging, configure_tracing, telemetry_is_enabled
from decouple import RepositoryEnv, UndefinedValueError
from loguru import logger
Expand Down
12 changes: 6 additions & 6 deletions cli/tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,8 @@
import pandas as pd
import pytest
from core.db.models import Base, Extract, Image
from core.patient_queue.message import Message
from core.patient_queue.producer import PixlProducer
from core.queue.models import ImagingRequestMessage
from core.queue.producer import PixlProducer
from sqlalchemy import Engine, create_engine
from sqlalchemy.orm import Session, sessionmaker

Expand Down Expand Up @@ -136,8 +136,8 @@ def _make_message(
accession_number: str,
mrn: str,
study_uid: str,
) -> Message:
return Message(
) -> ImagingRequestMessage:
return ImagingRequestMessage(
project_name=project_name,
accession_number=accession_number,
mrn=mrn,
Expand All @@ -150,7 +150,7 @@ def _make_message(


@pytest.fixture
def example_messages() -> list[Message]:
def example_messages() -> list[ImagingRequestMessage]:
"""Test input data."""
return [
_make_message(
Expand All @@ -174,7 +174,7 @@ def example_messages_df(example_messages):


@pytest.fixture
def example_messages_multiple_projects() -> list[Message]:
def example_messages_multiple_projects() -> list[ImagingRequestMessage]:
"""Test input data."""
return [
_make_message(
Expand Down
70 changes: 67 additions & 3 deletions cli/tests/test_message_processing.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,12 +16,17 @@

import os
from collections.abc import Generator
from unittest.mock import Mock
from unittest.mock import AsyncMock, Mock

import pytest
from _pytest.monkeypatch import MonkeyPatch
from core.patient_queue.producer import PixlProducer
from pixl_cli._message_processing import retry_until_export_count_is_unchanged
from core.queue.models import AnonymisationMessage
from core.queue.producer import PixlProducer
from pixl_cli._message_processing import (
_message_count,
retry_until_export_count_is_unchanged,
)
from pixl_imaging._orthanc import PIXLAnonOrthanc


@pytest.fixture
Expand Down Expand Up @@ -98,3 +103,62 @@ def test_retry_with_image_exported_and_no_change_multiple_projects(
)

mock_publisher.assert_called_once()


def test_message_count_includes_anonymisation(mocker) -> None:
"""Checks that the anonymisation queue is included when counting messages."""
mock_rabbitmq = Mock()
mock_rabbitmq.message_count = 0

mock_interface = mocker.patch("pixl_cli._message_processing.PixlBlockingInterface")
mock_interface.return_value.__enter__.return_value = mock_rabbitmq

_message_count(["imaging-primary"])

queue_names = {call.kwargs["queue_name"] for call in mock_interface.call_args_list}

assert queue_names == {
"imaging-primary",
"imaging-secondary",
"anonymisation",
}


@pytest.mark.asyncio
async def test_notify_anon_publishes_anonymisation_message(monkeypatch) -> None:
"""Checks that anonymisation requests are published to RabbitMQ."""
orthanc_raw = AsyncMock()
orthanc_raw.get_local_study.side_effect = [
{"MainDicomTags": {"StudyInstanceUID": "1.2.3"}},
{"MainDicomTags": {"StudyInstanceUID": "4.5.6"}},
]

producer = Mock()
producer_context = Mock()
producer_context.__enter__ = Mock(return_value=producer)
producer_context.__exit__ = Mock(return_value=None)

monkeypatch.setattr(
"pixl_imaging._orthanc.AnonymisationProducer",
Mock(return_value=producer_context),
)

orthanc_anon = PIXLAnonOrthanc()

await orthanc_anon.notify_anon_to_retrieve_study_resources(
orthanc_raw=orthanc_raw,
resource_ids=["resource-1", "resource-2"],
series_uid="1.2.3.1\\1.2.3.2",
project_name="test project",
)

producer.publish.assert_called_once_with(
[
AnonymisationMessage(
resource_ids=["resource-1", "resource-2"],
series_uids=["1.2.3.1", "1.2.3.2"],
study_uids=["1.2.3", "4.5.6"],
project_name="test project",
)
]
)
28 changes: 14 additions & 14 deletions cli/tests/test_messages_from_files.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@

import pytest
from core.db.models import Image
from core.patient_queue.message import Message
from core.queue.models import ImagingRequestMessage
from pixl_cli._io import read_patient_info
from pixl_cli._message_processing import messages_from_df, populate_queue_and_db

Expand All @@ -40,10 +40,10 @@ def test_messages_from_csv(omop_resources: Path) -> None:
# Act
messages = messages_from_df(messages_df)
# Assert
assert all(isinstance(msg, Message) for msg in messages)
assert all(isinstance(msg, ImagingRequestMessage) for msg in messages)

expected_messages = [
Message(
ImagingRequestMessage(
procedure_occurrence_id=0,
mrn="patient_identifier",
accession_number="123456789",
Expand Down Expand Up @@ -71,7 +71,7 @@ def test_whitespace_and_na_processing(omop_resources: Path) -> None:
messages = messages_from_df(messages_df)
# Assert
assert messages == [
Message(
ImagingRequestMessage(
procedure_occurrence_id=0,
mrn="patient_identifier",
accession_number="123456789",
Expand Down Expand Up @@ -117,10 +117,10 @@ def test_messages_from_parquet(omop_resources: Path) -> None:
# Act
messages = messages_from_df(messages_df)
# Assert
assert all(isinstance(msg, Message) for msg in messages)
assert all(isinstance(msg, ImagingRequestMessage) for msg in messages)

expected_messages = [
Message(
ImagingRequestMessage(
mrn="987654321",
accession_number="AA12345601",
study_uid="1.3.6.1.4.1.14519.5.2.1.99.1071.12985477682660597455732044031486",
Expand All @@ -130,7 +130,7 @@ def test_messages_from_parquet(omop_resources: Path) -> None:
project_name="test-extract-uclh-omop-cdm",
extract_generated_timestamp=datetime.datetime.fromisoformat("2023-12-07T14:08:58"),
),
Message(
ImagingRequestMessage(
mrn="987654321",
accession_number="AA12345605",
study_uid="1.2.276.0.7230010.3.1.2.929116473.1.1710754859.579485",
Expand All @@ -157,10 +157,10 @@ def test_messages_from_batched_parquet(omop_resources: Path) -> None:
# Act
messages = messages_from_df(messages_df)
# Assert
assert all(isinstance(msg, Message) for msg in messages)
assert all(isinstance(msg, ImagingRequestMessage) for msg in messages)

expected_messages = [
Message(
ImagingRequestMessage(
mrn="5020765",
accession_number="MIG0234560",
study_uid="1.2.840.114350.2.525.2.798268.2.110000014.1",
Expand All @@ -170,7 +170,7 @@ def test_messages_from_batched_parquet(omop_resources: Path) -> None:
project_name="test-extract-uclh-omop-cdm",
extract_generated_timestamp=datetime.datetime.fromisoformat("2023-12-07T14:08:58"),
),
Message(
ImagingRequestMessage(
mrn="987654321",
accession_number="ABC1234560",
study_uid="1.2.840.114350.2.525.2.798268.2.190000013.1",
Expand All @@ -180,7 +180,7 @@ def test_messages_from_batched_parquet(omop_resources: Path) -> None:
project_name="test-extract-uclh-omop-cdm",
extract_generated_timestamp=datetime.datetime.fromisoformat("2023-12-07T14:08:58"),
),
Message(
ImagingRequestMessage(
mrn="987654321",
accession_number="AA12345601",
study_uid="1.2.840.114350.2.525.2.798268.2.190000015.1",
Expand All @@ -190,7 +190,7 @@ def test_messages_from_batched_parquet(omop_resources: Path) -> None:
project_name="test-extract-uclh-omop-cdm",
extract_generated_timestamp=datetime.datetime.fromisoformat("2023-12-07T14:08:58"),
),
Message(
ImagingRequestMessage(
mrn="987654321",
accession_number="AA12345605",
study_uid="1.2.840.114350.2.525.2.798268.2.190000016.1",
Expand All @@ -200,7 +200,7 @@ def test_messages_from_batched_parquet(omop_resources: Path) -> None:
project_name="test-extract-uclh-omop-cdm",
extract_generated_timestamp=datetime.datetime.fromisoformat("2023-12-07T14:08:58"),
),
Message(
ImagingRequestMessage(
mrn="12345678",
accession_number="12345678",
study_uid="1.2.840.114350.2.525.2.798268.2.190000011.1",
Expand All @@ -210,7 +210,7 @@ def test_messages_from_batched_parquet(omop_resources: Path) -> None:
project_name="test-extract-uclh-omop-cdm",
extract_generated_timestamp=datetime.datetime.fromisoformat("2023-12-07T14:08:58"),
),
Message(
ImagingRequestMessage(
mrn="12345678",
accession_number="ABC1234567",
study_uid="1.2.840.114350.2.525.2.798268.2.190000012.1",
Expand Down
6 changes: 3 additions & 3 deletions cli/tests/test_populate.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,13 +19,13 @@

import pixl_cli._message_processing
from click.testing import CliRunner
from core.patient_queue.producer import PixlProducer
from core.queue.producer import PixlProducer
from pixl_cli.main import populate

if TYPE_CHECKING:
from pathlib import Path

from core.patient_queue.message import Message
from core.queue.models import ImagingRequestMessage


class MockProducer(PixlProducer):
Expand All @@ -39,7 +39,7 @@ def __exit__(self, *args: object, **kwargs) -> None:
"""Context exit point."""
return

def publish(self, messages: list[Message], priority: int) -> None: # noqa: ARG002 don't access messages or priority
def publish(self, messages: list[ImagingRequestMessage], priority: int) -> None: # noqa: ARG002 don't access messages or priority
"""Dummy method for publish."""
return

Expand Down
11 changes: 8 additions & 3 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,9 @@ services:
args:
PIXL_PACKAGE_DIR: hasher
<<: *build-args-common
depends_on:
queue:
condition: service_healthy
environment:
<<: [*proxy-common, *pixl-common-env, *otel-common]
OTEL_SERVICE_NAME: "hasher-api"
Expand Down Expand Up @@ -130,7 +133,7 @@ services:
command: /run/secrets
restart: always
environment:
<<: [*pixl-db, *proxy-common, *pixl-common-env, *azure-keyvault, *otel-common]
<<: [*pixl-db, *proxy-common, *pixl-common-env, *azure-keyvault, *otel-common, *pixl-rabbit-mq]
OTEL_SERVICE_NAME: "orthanc-anon"
ORTHANC_NAME: "PIXL: Anon"
ORTHANC_USERNAME: ${ORTHANC_ANON_USERNAME}
Expand Down Expand Up @@ -178,13 +181,15 @@ services:
depends_on:
postgres:
condition: service_healthy
queue:
condition: service_healthy
healthcheck:
test:
[
"CMD-SHELL",
"/probes/test-aliveness.py --user=$ORTHANC_USERNAME --pwd=$ORTHANC_PASSWORD",
]
start_period: 10s
start_period: 90s
retries: 10
interval: 3s
timeout: 2s
Expand Down Expand Up @@ -243,7 +248,7 @@ services:
"CMD-SHELL",
"/probes/test-aliveness.py --user=$ORTHANC_USERNAME --pwd=$ORTHANC_PASSWORD",
]
start_period: 10s
start_period: 90s
retries: 10
interval: 3s
timeout: 2s
Expand Down
Loading
Loading