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..e36b26a 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,15 +170,14 @@ impl DimseWadoService { &self, aet: &str, storescp_aet: &str, + study_instance_uid: &str, identifier: InMemDicomObject, ) -> BoxStream<'static, Result>, MoveError>> { 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), - RetrieveMode::Sequential => SubscriptionTopic::unidentified(AE::from(aet)), - }; + let subscription_topic = + subscription_topic(self.config.mode, aet, study_instance_uid, message_id); let subscription = self .mediator .subscribe(subscription_topic, tx.clone()) @@ -225,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] @@ -335,6 +353,7 @@ impl Stream for DicomMultipartStream<'_> { #[cfg(test)] mod tests { use super::*; + use dicom::core::DataElement; use dicom::object::FileMetaTableBuilder; use futures::TryStreamExt; @@ -349,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()))]);