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()))]);