From 2ec717110b18945147fa11ec651e0829b56ecc15 Mon Sep 17 00:00:00 2001 From: Alex Luft Date: Sun, 27 Sep 2026 06:34:43 +0000 Subject: [PATCH 1/2] fix(dimse): assign sequential-mode instances by StudyInstanceUID (#71) In `sequential` retrieve mode every received instance went to whichever retrieve was active for the AET. When a retrieve ended early (client disconnect, timeout) the C-MOVE was not cancelled and the per-AET lock was released, so the PACS kept sending the previous study and the next retrieve for that AET received its instances. Sequential-mode subscriptions are now keyed by (AET, StudyInstanceUID), and an incoming instance is delivered only to the retrieve of its own study. Concurrent mode (matching by Move Originator Message ID) is unchanged. The per-AET semaphore becomes a per-(AET, study) one: C-MOVEs of the same study are still serialized, C-MOVEs of different studies no longer wait for each other. Instances without a StudyInstanceUID cannot be attributed in sequential mode and are dropped with the existing warning. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 7 + docs/topics/configuration.md | 10 +- src/backend/dimse/cmove/mediator.rs | 291 ++++++++++++++++++++++++---- src/backend/dimse/cmove/mod.rs | 18 +- src/backend/dimse/wado.rs | 11 +- 5 files changed, 288 insertions(+), 49 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index b9365d0..7c79dff 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,9 +13,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- `sequential` retrieve mode could return instances of a different study: when a retrieve ended early (client + disconnect, timeout), the PACS kept sending the previous study, and the next retrieve for the same AET received + those instances. Incoming instances are now assigned by their `StudyInstanceUID`, and only C-MOVEs of the same + study are serialized, so retrieves of different studies no longer wait for each other (#71) + ### Changed - Updated `dicom-rs` dependency to 0.10.0 +- `sequential` retrieve mode now only serializes C-MOVEs of the same study; C-MOVEs of different studies run + concurrently, limited by the association pool size (#71) ## [0.3.1] - 2026-09-14 diff --git a/docs/topics/configuration.md b/docs/topics/configuration.md index 547443f..5fcde5f 100644 --- a/docs/topics/configuration.md +++ b/docs/topics/configuration.md @@ -218,12 +218,14 @@ Each AET (regardless of the backend) has additional settings specific to the DIC DIMSE-backend only: Some PACS do not include the MoveOriginatorMessageId attribute in their C-STORE-RQ messages. This makes it hard to assign incoming C-STORE-RQ responses to an active C-MOVE operation. - As a workaround, you can set the receive mode to sequential to disable concurrent C-MOVEs. - Throughput will be limited, but it will work reliably. Consider increasing the timeouts when using the sequential mode. + As a workaround, you can set the receive mode to sequential, which assigns incoming instances by their + StudyInstanceUID instead. C-MOVEs of the same study run one at a time; C-MOVEs of different studies run concurrently. + Instances without a StudyInstanceUID cannot be assigned in this mode and are dropped with a warning. + Consider increasing the timeouts when using the sequential mode. Most PACS can and should use the concurrent mode. -
  • concurrent: C-MOVE requests are processed concurrently.
  • -
  • sequential: C-MOVE requests are processed sequentially.
  • +
  • concurrent: C-MOVE requests are processed concurrently; instances are assigned by MoveOriginatorMessageId.
  • +
  • sequential: C-MOVE requests for the same study are processed sequentially; instances are assigned by StudyInstanceUID.
  • diff --git a/src/backend/dimse/cmove/mediator.rs b/src/backend/dimse/cmove/mediator.rs index 053bfd3..2dfa7de 100644 --- a/src/backend/dimse/cmove/mediator.rs +++ b/src/backend/dimse/cmove/mediator.rs @@ -1,17 +1,18 @@ use crate::backend::dimse::cmove::movescu::MoveError; use crate::backend::dimse::cmove::MoveSubOperation; use crate::config::{AppConfig, RetrieveMode}; -use crate::types::{AE, US}; +use crate::types::{AE, UI, US}; use std::collections::HashMap; use std::sync::{Arc, Weak}; use thiserror::Error; use tokio::sync::mpsc::Sender; -use tokio::sync::{OwnedSemaphorePermit, RwLock, Semaphore}; +use tokio::sync::{Mutex, OwnedSemaphorePermit, RwLock, Semaphore}; use tracing::info; pub type Callback = Sender>; /// A mediator for the communication between the MOVE-SCU and STORE-SCP. +#[derive(Default)] pub struct MoveMediator { inner: Arc, } @@ -26,41 +27,42 @@ impl Clone for MoveMediator { #[derive(Default)] struct InnerMoveMediator { - semaphores: RwLock>>, + /// Sequential mode: one lock per (AE, `StudyInstanceUID`), so that at most one + /// C-MOVE per study is in flight for an AE. Entries are held weakly; one dies + /// with its last permit or waiter and is pruned on the next subscription. + study_locks: Mutex>>, callbacks: RwLock>, } impl MoveMediator { pub fn new(config: &AppConfig) -> Self { - let mut semaphores = HashMap::new(); for ae in &config.aets { if matches!(ae.wado.mode, RetrieveMode::Sequential) { info!( - "Using Sequential Retrieve Mode for {} - Reduced performance is expected.", + "Using Sequential Retrieve Mode for {} - instances are assigned by StudyInstanceUID; C-MOVEs of the same study run one at a time.", ae.aet ); - semaphores.insert(AE::from(&ae.aet), Arc::new(Semaphore::new(1))); } } - Self { - inner: Arc::new(InnerMoveMediator { - semaphores: RwLock::new(semaphores), - callbacks: RwLock::new(HashMap::new()), - }), - } + Self::default() } + /// Registers `callback` for the sub-operations of `topic`. + /// + /// A study-keyed topic (sequential mode) first waits until no other + /// subscription for the same AE and study is active: instances are attributed + /// by study alone, so two retrieves of one study must not overlap. + /// + /// # Panics + /// Never in practice: the study locks are never closed. pub async fn subscribe(&self, topic: SubscriptionTopic, callback: Callback) -> Subscription { - let semaphore: Option> = { - let semaphores = self.inner.semaphores.read().await; - let semaphore = semaphores.get(&topic.originator).cloned(); - drop(semaphores); - semaphore - }; - - let permit = if let Some(semaphore) = semaphore { - let permit = semaphore.acquire_owned().await.unwrap(); - Some(permit) + let permit = if let Some(study_instance_uid) = &topic.study_instance_uid { + let lock = self.study_lock(&topic.originator, study_instance_uid).await; + Some( + lock.acquire_owned() + .await + .expect("study locks are never closed"), + ) } else { None }; @@ -75,27 +77,47 @@ impl MoveMediator { } } + async fn study_lock(&self, originator: &str, study_instance_uid: &str) -> Arc { + let mut locks = self.inner.study_locks.lock().await; + locks.retain(|_, lock| lock.strong_count() > 0); + let key = (AE::from(originator), UI::from(study_instance_uid)); + if let Some(lock) = locks.get(&key).and_then(Weak::upgrade) { + return lock; + } + let lock = Arc::new(Semaphore::new(1)); + locks.insert(key, Arc::downgrade(&lock)); + lock + } + pub async fn unsubscribe(&self, topic: &SubscriptionTopic) { let mut callbacks = self.inner.callbacks.write().await; callbacks.remove(topic); } + /// Delivers a sub-operation published by the STORE-SCP. + /// + /// A topic with a Move Originator Message ID goes to the retrieve that issued + /// that C-MOVE (concurrent mode). Otherwise, and for peers whose message IDs + /// match no concurrent retrieve, the instance goes to the sequential-mode + /// retrieve of its OWN study on that AE, or nowhere: an instance of another + /// study (e.g. still arriving from an abandoned C-MOVE) is never attributed + /// to the retrieve that happens to be active. pub async fn publish( &self, topic: &SubscriptionTopic, sub_operation: Result, ) -> Result<(), MediatorError> { + let study_topic = sub_operation + .as_ref() + .ok() + .and_then(MoveSubOperation::study_instance_uid) + .map(|study| SubscriptionTopic::for_study(topic.originator.clone(), study)); + let callbacks = self.inner.callbacks.read().await; - let callback = if topic.message_id.is_some() { - callbacks.get(topic).or_else(|| { - callbacks.get(&SubscriptionTopic { - originator: topic.originator.clone(), - message_id: None, - }) - }) - } else { - callbacks.get(topic) - }; + let callback = topic + .message_id + .and_then(|_| callbacks.get(topic)) + .or_else(|| study_topic.as_ref().and_then(|study| callbacks.get(study))); if let Some(callback) = callback { callback .send(sub_operation) @@ -122,6 +144,8 @@ pub enum MediatorError { pub struct SubscriptionTopic { pub originator: AE, pub message_id: Option, + /// Sequential mode: the study the retrieve asked for. + pub study_instance_uid: Option, } pub struct Subscription { @@ -146,30 +170,213 @@ impl Drop for Subscription { impl SubscriptionTopic { pub const fn new(originator: AE, message_id: Option) -> Self { - if let Some(message_id) = message_id { - Self::identified(originator, message_id) - } else { - Self::unidentified(originator) - } - } - pub const fn identified(originator: AE, message_id: US) -> Self { Self { originator, - message_id: Some(message_id), + message_id, + study_instance_uid: None, } } - pub const fn unidentified(originator: AE) -> Self { + pub const fn identified(originator: AE, message_id: US) -> Self { + Self::new(originator, Some(message_id)) + } + + /// Sequential mode: sub-operations are attributed by the study they belong to. + pub const fn for_study(originator: AE, study_instance_uid: UI) -> Self { Self { originator, message_id: None, + study_instance_uid: Some(study_instance_uid), } } pub fn without_message_id(self) -> Self { Self { - originator: self.originator, message_id: None, + ..self + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use dicom::core::{DataElement, VR}; + use dicom::dictionary_std::tags; + use dicom::object::{FileDicomObject, FileMetaTableBuilder, InMemDicomObject}; + use std::time::Duration; + use tokio::sync::mpsc; + + const PACS: &str = "PACS"; + const STUDY_A: &str = "2.25.1001"; + const STUDY_B: &str = "2.25.1002"; + + fn instance(study_instance_uid: Option<&str>) -> MoveSubOperation { + let mut dataset = InMemDicomObject::new_empty(); + if let Some(uid) = study_instance_uid { + dataset.put(DataElement::new(tags::STUDY_INSTANCE_UID, VR::UI, uid)); + } + let file: FileDicomObject = dataset.with_exact_meta( + FileMetaTableBuilder::new() + .media_storage_sop_class_uid("1.2.840.10008.5.1.4.1.1.7") + .media_storage_sop_instance_uid("2.25.4242") + .transfer_syntax("1.2.840.10008.1.2.1") + .build() + .expect("FileMetaTableBuilder should contain required data"), + ); + MoveSubOperation::Pending(Arc::new(file)) + } + + fn received_study(operation: &MoveSubOperation) -> Option { + operation.study_instance_uid() + } + + /// What the STORE-SCP publishes for a C-STORE without Move Originator Message ID. + fn store_topic() -> SubscriptionTopic { + SubscriptionTopic::new(AE::from(PACS), None) + } + + fn study_topic(study: &str) -> SubscriptionTopic { + SubscriptionTopic::for_study(AE::from(PACS), UI::from(study)) + } + + /// Regression test for #71: after a retrieve is abandoned, the peer keeps pushing + /// its study; those instances must not reach the next retrieve on the same AE. + #[tokio::test(flavor = "multi_thread")] + async fn an_abandoned_retrieve_does_not_leak_into_the_next_one() { + let mediator = MoveMediator::default(); + + let (tx_a, _rx_a) = mpsc::channel(8); + let retrieve_a = mediator.subscribe(study_topic(STUDY_A), tx_a).await; + drop(retrieve_a); // client disconnected or timed out + + let (tx_b, mut rx_b) = mpsc::channel(8); + let _retrieve_b = mediator.subscribe(study_topic(STUDY_B), tx_b).await; + + let straggler = mediator + .publish(&store_topic(), Ok(instance(Some(STUDY_A)))) + .await; + assert!(matches!( + straggler, + Err(MediatorError::MissingCallback { .. }) + )); + assert!( + rx_b.try_recv().is_err(), + "an instance of study A reached the retrieve of study B" + ); + + mediator + .publish(&store_topic(), Ok(instance(Some(STUDY_B)))) + .await + .expect("the instance of study B should be delivered"); + let received = rx_b.recv().await.expect("channel open").expect("pending"); + assert_eq!(received_study(&received).as_deref(), Some(STUDY_B)); + } + + #[tokio::test(flavor = "multi_thread")] + async fn sequential_mode_drops_instances_without_study_instance_uid() { + let mediator = MoveMediator::default(); + let (tx, mut rx) = mpsc::channel(8); + let _retrieve = mediator.subscribe(study_topic(STUDY_A), tx).await; + + let result = mediator.publish(&store_topic(), Ok(instance(None))).await; + assert!(matches!(result, Err(MediatorError::MissingCallback { .. }))); + assert!(rx.try_recv().is_err()); + } + + #[tokio::test(flavor = "multi_thread")] + async fn sequential_mode_accepts_instances_that_carry_a_message_id() { + // A peer that does send Move Originator Message IDs, configured sequential. + let mediator = MoveMediator::default(); + let (tx, mut rx) = mpsc::channel(8); + let _retrieve = mediator.subscribe(study_topic(STUDY_A), tx).await; + + let topic = SubscriptionTopic::new(AE::from(PACS), Some(7)); + mediator + .publish(&topic, Ok(instance(Some(STUDY_A)))) + .await + .expect("delivered by study"); + assert!(rx.try_recv().is_ok()); + } + + #[tokio::test(flavor = "multi_thread")] + async fn concurrent_mode_still_matches_by_message_id() { + let mediator = MoveMediator::default(); + let (tx_7, mut rx_7) = mpsc::channel(8); + let (tx_8, mut rx_8) = mpsc::channel(8); + let _retrieve_7 = mediator + .subscribe(SubscriptionTopic::identified(AE::from(PACS), 7), tx_7) + .await; + let _retrieve_8 = mediator + .subscribe(SubscriptionTopic::identified(AE::from(PACS), 8), tx_8) + .await; + + let topic = SubscriptionTopic::new(AE::from(PACS), Some(8)); + mediator + .publish(&topic, Ok(instance(Some(STUDY_A)))) + .await + .expect("delivered by message id"); + assert!(rx_8.try_recv().is_ok()); + assert!(rx_7.try_recv().is_err()); + + // Without a message ID there is nothing to match in concurrent mode. + let result = mediator + .publish(&store_topic(), Ok(instance(Some(STUDY_A)))) + .await; + assert!(matches!(result, Err(MediatorError::MissingCallback { .. }))); + } + + #[tokio::test(flavor = "multi_thread")] + async fn retrieves_of_different_studies_run_concurrently() { + let mediator = MoveMediator::default(); + let (tx_a, _rx_a) = mpsc::channel(8); + let (tx_b, _rx_b) = mpsc::channel(8); + let _retrieve_a = mediator.subscribe(study_topic(STUDY_A), tx_a).await; + tokio::time::timeout( + Duration::from_secs(5), + mediator.subscribe(study_topic(STUDY_B), tx_b), + ) + .await + .expect("a retrieve of another study must not wait"); + } + + // `first` is held on purpose: the second retrieve must wait until it is dropped. + #[allow(clippy::significant_drop_tightening)] + #[tokio::test(flavor = "multi_thread")] + async fn retrieves_of_the_same_study_run_one_at_a_time() { + let mediator = MoveMediator::default(); + let (tx_1, _rx_1) = mpsc::channel(8); + let (tx_2, _rx_2) = mpsc::channel(8); + let first = mediator.subscribe(study_topic(STUDY_A), tx_1).await; + + let second = tokio::spawn({ + let mediator = mediator.clone(); + async move { mediator.subscribe(study_topic(STUDY_A), tx_2).await } + }); + tokio::time::sleep(Duration::from_millis(100)).await; + assert!( + !second.is_finished(), + "a second retrieve of the same study must wait" + ); + + drop(first); + let second = tokio::time::timeout(Duration::from_secs(5), second) + .await + .expect("the waiting retrieve should start once the first ends") + .expect("task should not panic"); + drop(second); + } + + #[tokio::test(flavor = "multi_thread")] + async fn study_locks_are_pruned_once_unused() { + let mediator = MoveMediator::default(); + for study in [STUDY_A, STUDY_B] { + let (tx, _rx) = mpsc::channel(8); + drop(mediator.subscribe(study_topic(study), tx).await); } + // The next subscription prunes every lock nobody holds or waits for. + let (tx, _rx) = mpsc::channel(8); + let _retrieve = mediator.subscribe(study_topic(STUDY_A), tx).await; + assert_eq!(mediator.inner.study_locks.lock().await.len(), 1); } } diff --git a/src/backend/dimse/cmove/mod.rs b/src/backend/dimse/cmove/mod.rs index dffa47c..a55c8b6 100644 --- a/src/backend/dimse/cmove/mod.rs +++ b/src/backend/dimse/cmove/mod.rs @@ -1,5 +1,5 @@ use crate::backend::dimse::{DicomMessage, DATA_SET_EXISTS}; -use crate::types::{AE, US}; +use crate::types::{AE, UI, US}; use dicom::core::{DataElement, VR}; use dicom::dicom_value; use dicom::dictionary_std::{tags, uids}; @@ -45,3 +45,19 @@ pub enum MoveSubOperation { Completed, Pending(Arc>), } + +impl MoveSubOperation { + /// The `StudyInstanceUID` of a received instance, if it carries a non-empty one. + pub fn study_instance_uid(&self) -> Option { + match self { + Self::Pending(file) => file + .element(tags::STUDY_INSTANCE_UID) + .ok()? + .to_str() + .ok() + .filter(|uid| !uid.is_empty()) + .map(|uid| UI::from(uid.as_ref())), + Self::Completed => None, + } + } +} diff --git a/src/backend/dimse/wado.rs b/src/backend/dimse/wado.rs index c0d8a47..f76c51f 100644 --- a/src/backend/dimse/wado.rs +++ b/src/backend/dimse/wado.rs @@ -11,7 +11,7 @@ use crate::backend::dimse::{next_message_id, WriteError}; use crate::config::{RetrieveMode, WadoConfig}; use crate::rendering::render_instances; use crate::types::{Priority, US}; -use crate::types::{QueryRetrieveLevel, AE}; +use crate::types::{QueryRetrieveLevel, AE, UI}; use association::pool::AssociationPool; use async_stream::stream; use async_trait::async_trait; @@ -63,6 +63,7 @@ impl WadoService for DimseWadoService { .retrieve_instances( &request.query.aet, storescp_aet, + &request.query.study_instance_uid, Self::create_identifier(Some(&request.query.study_instance_uid), None, None), ) .await; @@ -91,6 +92,7 @@ impl WadoService for DimseWadoService { .retrieve_instances( &request.query.aet, storescp_aet, + &request.query.study_instance_uid, Self::create_identifier(Some(&request.query.study_instance_uid), None, None), ) .await @@ -168,6 +170,7 @@ impl DimseWadoService { &self, aet: &str, storescp_aet: &str, + study_instance_uid: &str, identifier: InMemDicomObject, ) -> BoxStream<'static, Result>, MoveError>> { let message_id = next_message_id(); @@ -175,7 +178,11 @@ impl DimseWadoService { let subscription_topic = match self.config.mode { RetrieveMode::Concurrent => SubscriptionTopic::identified(AE::from(aet), message_id), - RetrieveMode::Sequential => SubscriptionTopic::unidentified(AE::from(aet)), + // The C-MOVE identifier is study-level, and the peer does not tell us + // which C-MOVE an instance answers: attribute instances by their study. + RetrieveMode::Sequential => { + SubscriptionTopic::for_study(AE::from(aet), UI::from(study_instance_uid)) + } }; let subscription = self .mediator From ce5372e4613f78eff2775014ff3d147bc7dd1fca Mon Sep 17 00:00:00 2001 From: Alex Luft Date: Sun, 27 Sep 2026 07:14:32 +0000 Subject: [PATCH 2/2] test(dimse): cover the study a sequential retrieve subscribes to The choice of mediator topic moves into `subscription_topic()` and gets a test that routes through the mediator: if a sequential retrieve stopped subscribing to the study it requested, every sequential retrieve would 404 with all other tests green. Plus a check that concurrent retrieves still subscribe to their message ID. No behaviour change. Co-Authored-By: Claude Opus 5.5 --- src/backend/dimse/wado.rs | 73 ++++++++++++++++++++++++++++++++++----- 1 file changed, 65 insertions(+), 8 deletions(-) diff --git a/src/backend/dimse/wado.rs b/src/backend/dimse/wado.rs index f76c51f..e36b26a 100644 --- a/src/backend/dimse/wado.rs +++ b/src/backend/dimse/wado.rs @@ -176,14 +176,8 @@ impl DimseWadoService { let message_id = next_message_id(); let (tx, mut rx) = mpsc::channel::>(1); - let subscription_topic = match self.config.mode { - RetrieveMode::Concurrent => SubscriptionTopic::identified(AE::from(aet), message_id), - // The C-MOVE identifier is study-level, and the peer does not tell us - // which C-MOVE an instance answers: attribute instances by their study. - RetrieveMode::Sequential => { - SubscriptionTopic::for_study(AE::from(aet), UI::from(study_instance_uid)) - } - }; + let subscription_topic = + subscription_topic(self.config.mode, aet, study_instance_uid, message_id); let subscription = self .mediator .subscribe(subscription_topic, tx.clone()) @@ -232,6 +226,23 @@ impl DimseWadoService { } } +/// The mediator topic a retrieve subscribes to. +fn subscription_topic( + mode: RetrieveMode, + aet: &str, + study_instance_uid: &str, + message_id: US, +) -> SubscriptionTopic { + match mode { + RetrieveMode::Concurrent => SubscriptionTopic::identified(AE::from(aet), message_id), + // The C-MOVE identifier is study-level, and the peer does not tell us + // which C-MOVE an instance answers: attribute instances by their study. + RetrieveMode::Sequential => { + SubscriptionTopic::for_study(AE::from(aet), UI::from(study_instance_uid)) + } + } +} + /// Stream that takes ownership of a value. /// Especially useful for keeping semaphore permits until the stream is completed. #[pin_project] @@ -342,6 +353,7 @@ impl Stream for DicomMultipartStream<'_> { #[cfg(test)] mod tests { use super::*; + use dicom::core::DataElement; use dicom::object::FileMetaTableBuilder; use futures::TryStreamExt; @@ -356,6 +368,51 @@ mod tests { ) } + fn instance_of_study(study_instance_uid: &str) -> MoveSubOperation { + let mut file = test_file(); + file.put(DataElement::new( + tags::STUDY_INSTANCE_UID, + VR::UI, + study_instance_uid, + )); + MoveSubOperation::Pending(Arc::new(file)) + } + + /// A sequential retrieve must subscribe to the study it requested: that is the + /// only key a C-STORE without Move Originator Message ID can be matched by. + #[tokio::test(flavor = "multi_thread")] + async fn a_sequential_retrieve_receives_its_requested_study_only() { + let mediator = MoveMediator::default(); + let (tx, mut rx) = mpsc::channel(8); + let _subscription = mediator + .subscribe( + subscription_topic(RetrieveMode::Sequential, "PACS", "2.25.1001", 7), + tx, + ) + .await; + + // What the STORE-SCP publishes for a peer that omits the message ID. + let store_topic = SubscriptionTopic::new(AE::from("PACS"), None); + mediator + .publish(&store_topic, Ok(instance_of_study("2.25.1001"))) + .await + .expect("an instance of the requested study must be delivered"); + assert!(rx.try_recv().is_ok()); + + let other = mediator + .publish(&store_topic, Ok(instance_of_study("2.25.1002"))) + .await; + assert!(other.is_err(), "an instance of another study was delivered"); + } + + #[test] + fn a_concurrent_retrieve_subscribes_to_its_message_id() { + assert_eq!( + subscription_topic(RetrieveMode::Concurrent, "PACS", "2.25.1001", 7), + SubscriptionTopic::identified(AE::from("PACS"), 7) + ); + } + #[tokio::test] async fn multipart_stream_uses_the_provided_boundary() { let stream = futures::stream::iter(vec![Ok(Arc::new(test_file()))]);