From 9da459930a604fd76a09bc4a36447594da54beb9 Mon Sep 17 00:00:00 2001 From: Johannes Schrimpf Date: Fri, 4 Sep 2026 23:46:05 +0200 Subject: [PATCH] Fix DiscoveryDB and participant lock inversion during endpoint creation Concurrent endpoint creation and discovery snapshot queries on the same DomainParticipant could deadlock. Writer and reader creation held the DiscoveryDB write lock while building records that acquired the participant mutex (dpi). The discovered_* accessors acquired dpi first and then read DiscoveryDB, producing an ABBA lock cycle. Build writer and reader discovery records before taking the database write lock. Split local reader record construction from insertion so participant access happens outside the database critical section. In the three snapshot accessors, clone the database handle under dpi and release the participant mutex before reading the database. Also release the writer's database guard before the blocking send to add_writer_sender. That bounded channel is drained by the participant event loop, which itself reads DiscoveryDB; holding the write guard while waiting for channel capacity could prevent the consumer from making progress. Document these non-overlapping lock scopes and identify find_topic() as the remaining path that holds dpi while accessing DiscoveryDB. find_topic() continues to hold the participant mutex during its lookup and can delay other participant calls until it returns. Add a regression test that races endpoint creation against discovery snapshots in a child process. Complete publisher/subscriber setup before spawning barrier workers, exercise every accessor at least once, and fail after a 30-second deadline instead of hanging the test harness. Require a dedicated child completion exit code so an unmatched test filter cannot pass silently. Report child failures separately from a possible deadlock; the exact interleavings remain scheduler-dependent. A separate sanity test checks that a local writer's topic is visible. Reuse discovery_db_read for the snapshot accessors' poison handling. Validation on macOS with separate unfixed/fixed Cargo build directories: - The updated race test reaches the 30-second deadlock deadline on the unchanged pre-fix base. - Both focused integration tests pass with default features and security. - Temporary unmatched-filter and creator-setup-panic probes fail promptly through the child-failure path, rather than passing or timing out. - The integration test passes nightly-2026-09-08 rustfmt; git diff whitespace checks pass. Existing crate-wide formatting failures are outside this lock-order fix. Assisted-by: Claude Fable 5.1 Assisted-by: OpenAI Codex (gpt-6-astra) --- src/dds/participant.rs | 59 +++---- src/dds/pubsub.rs | 64 +++++--- src/discovery/discovery_db.rs | 49 ++++-- tests/lock_order_discovery_db.rs | 269 +++++++++++++++++++++++++++++++ 4 files changed, 368 insertions(+), 73 deletions(-) create mode 100644 tests/lock_order_discovery_db.rs diff --git a/src/dds/participant.rs b/src/dds/participant.rs index 0e212a1f..aa48bd73 100644 --- a/src/dds/participant.rs +++ b/src/dds/participant.rs @@ -34,7 +34,7 @@ use crate::{ }, discovery::{ discovery::{Discovery, DiscoveryCommand}, - discovery_db::DiscoveryDB, + discovery_db::{discovery_db_read, DiscoveryDB}, sedp_messages::{DiscoveredReaderData, DiscoveredTopicData, DiscoveredWriterData}, }, network::{constant::*, udp_listener::UDPListener}, @@ -531,7 +531,12 @@ impl DomainParticipant { /// } /// ``` pub fn discovered_topics(&self) -> Vec { - self.dpi.lock().unwrap().discovered_topics() + // Clone the Discovery DB handle under `dpi`, then release `dpi` before + // reading the DB. These locks are not held together, so waiting for the + // DB does not block other participant calls that only need `dpi`. + let db = self.discovery_db(); + let db = discovery_db_read(&db); + db.all_user_topics().cloned().collect() } /// Gets a snapshot of all Readers discovered over the DDS network. @@ -553,7 +558,12 @@ impl DomainParticipant { /// } /// ``` pub fn discovered_readers(&self) -> Vec { - self.dpi.lock().unwrap().discovered_readers() + // Clone the Discovery DB handle under `dpi`, then release `dpi` before + // reading the DB. These locks are not held together, so waiting for the + // DB does not block other participant calls that only need `dpi`. + let db = self.discovery_db(); + let db = discovery_db_read(&db); + db.get_all_external_topic_readers().cloned().collect() } /// Gets a snapshot of all Writers discovered over the DDS network. @@ -575,7 +585,12 @@ impl DomainParticipant { /// } /// ``` pub fn discovered_writers(&self) -> Vec { - self.dpi.lock().unwrap().discovered_writers() + // Clone the Discovery DB handle under `dpi`, then release `dpi` before + // reading the DB. These locks are not held together, so waiting for the + // DB does not block other participant calls that only need `dpi`. + let db = self.discovery_db(); + let db = discovery_db_read(&db); + db.get_all_external_topic_writers().cloned().collect() } /// Manually asserts liveliness, affecting all writers with @@ -992,18 +1007,6 @@ impl DomainParticipantDisc { self.dpi.participant_id() } - pub fn discovered_topics(&self) -> Vec { - self.dpi.discovered_topics() - } - - pub fn discovered_readers(&self) -> Vec { - self.dpi.discovered_readers() - } - - pub fn discovered_writers(&self) -> Vec { - self.dpi.discovered_writers() - } - pub(crate) fn dds_cache(&self) -> Arc> { self.dpi.dds_cache() } @@ -1598,30 +1601,6 @@ impl DomainParticipantInner { self.domain_info.participant_id } - pub fn discovered_topics(&self) -> Vec { - let db = self.discovery_db.read().unwrap_or_else(|e| { - panic!("RustDDS internal bug: DiscoveryDB is poisoned after a prior panic: {e:?}") - }); - - db.all_user_topics().cloned().collect() - } - - pub fn discovered_readers(&self) -> Vec { - let db = self.discovery_db.read().unwrap_or_else(|e| { - panic!("RustDDS internal bug: DiscoveryDB is poisoned after a prior panic: {e:?}") - }); - - db.get_all_external_topic_readers().cloned().collect() - } - - pub fn discovered_writers(&self) -> Vec { - let db = self.discovery_db.read().unwrap_or_else(|e| { - panic!("RustDDS internal bug: DiscoveryDB is poisoned after a prior panic: {e:?}") - }); - - db.get_all_external_topic_writers().cloned().collect() - } - pub(crate) fn status_channel_receiver( &self, ) -> &StatusChannelReceiver { diff --git a/src/dds/pubsub.rs b/src/dds/pubsub.rs index f0aa6794..64e4e3f1 100644 --- a/src/dds/pubsub.rs +++ b/src/dds/pubsub.rs @@ -646,29 +646,43 @@ impl InnerPublisher { None }; - // Add the topic & writer to Discovery DB - let mut db = self - .discovery_db - .write() - .map_err(|e| CreateError::Poisoned { - reason: format!("Discovery DB: {e}"), - })?; - + // Build the discovery record before taking the Discovery DB lock: + // `DiscoveredWriterData::new` locks the participant (`dpi`) for its domain + // id, participant id, networks and GUID. This keeps the two locks from + // overlapping. Taking `dpi` while holding `discovery_db` could deadlock + // against `DomainParticipant::find_topic()`, which holds `dpi` while + // reading the DB. let dwd = DiscoveredWriterData::new(&data_writer, topic, &dp, security_info); - db.update_local_topic_writer(dwd); - db.update_topic_data_p(topic); - // Inform Discovery about the topic - if let Err(e) = self.discovery_command.try_send(DiscoveryCommand::AddTopic { - topic_name: topic.name(), - }) { - // Log the error but don't quit, failing to inform Discovery about the topic - // shouldn't be that serious - error!( - "Failed send DiscoveryCommand::AddTopic about topic {}: {}", - topic.name(), - e - ); + // Add the topic & writer to Discovery DB. The write guard is scoped so + // that it is released before the blocking `add_writer_sender.send` below: + // that channel is bounded and drained by the DP event loop, which itself + // takes `discovery_db` for reading, so holding the guard across the send + // could deadlock against it (the reader path below already scopes its + // guard the same way). + { + let mut db = self + .discovery_db + .write() + .map_err(|e| CreateError::Poisoned { + reason: format!("Discovery DB: {e}"), + })?; + + db.update_local_topic_writer(dwd); + db.update_topic_data_p(topic); + + // Inform Discovery about the topic + if let Err(e) = self.discovery_command.try_send(DiscoveryCommand::AddTopic { + topic_name: topic.name(), + }) { + // Log the error but don't quit, failing to inform Discovery about the topic + // shouldn't be that serious + error!( + "Failed send DiscoveryCommand::AddTopic about topic {}: {}", + topic.name(), + e + ); + } } // Note: notifying Discovery about the new writer is no longer done here. @@ -1234,13 +1248,19 @@ impl InnerSubscriber { None }; + // Build the discovery record before taking the Discovery DB lock: it locks + // the participant (`dpi`) for its locators and GUID. Keeping the two locks + // from overlapping avoids the reverse of `DomainParticipant::find_topic()`'s + // nested `dpi` -> `discovery_db` order. + let drd = DiscoveryDB::local_topic_reader_data(&dp, topic, &new_reader, security_info); + // Add the topic & reader to Discovery DB { let mut db = self .discovery_db .write() .or_else(|e| create_error_poisoned!("Cannot lock discovery_db. {}", e))?; - db.update_local_topic_reader(&dp, topic, &new_reader, security_info); + db.update_local_topic_reader(drd); db.update_topic_data_p(topic); // Inform Discovery about the topic diff --git a/src/discovery/discovery_db.rs b/src/discovery/discovery_db.rs index f9aa6405..de3a3020 100644 --- a/src/discovery/discovery_db.rs +++ b/src/discovery/discovery_db.rs @@ -635,13 +635,21 @@ impl DiscoveryDB { } // local topic readers - pub fn update_local_topic_reader( - &mut self, + + /// Builds the discovery record of a local reader. + /// + /// This locks the participant (`dpi`, for its locators and GUID), so it is + /// a separate step from [`Self::update_local_topic_reader`]: build the record + /// first, then take the Discovery DB lock and insert it. The two locks are + /// not held together. Calling this while holding `discovery_db` could + /// deadlock against `DomainParticipant::find_topic()`, which holds `dpi` + /// while reading the DB. + pub fn local_topic_reader_data( domain_participant: &DomainParticipant, topic: &Topic, reader: &ReaderIngredients, sec_info_opt: Option, - ) { + ) -> DiscoveredReaderData { let reader_guid = reader.guid; let reader_proxy = RtpsReaderProxy::from_reader(reader, domain_participant); @@ -658,16 +666,20 @@ impl DiscoveryDB { // TODO: possibly change content filter to dynamic value let content_filter = None; - let discovered_reader_data = DiscoveredReaderData { + DiscoveredReaderData { reader_proxy: ReaderProxy::from(reader_proxy), subscription_topic_data: subscription_data, content_filter, user_data: Vec::new(), - }; + } + } - self - .local_topic_readers - .insert(reader_guid, discovered_reader_data); + /// Inserts a record built by [`Self::local_topic_reader_data`]. + pub fn update_local_topic_reader(&mut self, discovered_reader_data: DiscoveredReaderData) { + self.local_topic_readers.insert( + discovered_reader_data.reader_proxy.remote_reader_guid, + discovered_reader_data, + ); } pub fn remove_local_topic_reader(&mut self, guid: GUID) { @@ -1100,12 +1112,22 @@ mod tests { }; // Add the reader to the database and verify the info is updated - discoverydb.update_local_topic_reader(&dp, &topic, &reader1_ing, None); + discoverydb.update_local_topic_reader(DiscoveryDB::local_topic_reader_data( + &dp, + &topic, + &reader1_ing, + None, + )); assert_eq!(discoverydb.local_topic_readers.len(), 1); assert_eq!(discoverydb.get_local_topic_readers(&topic).len(), 1); // Verify that the info does not change if the reader is added a second time - discoverydb.update_local_topic_reader(&dp, &topic, &reader1_ing, None); + discoverydb.update_local_topic_reader(DiscoveryDB::local_topic_reader_data( + &dp, + &topic, + &reader1_ing, + None, + )); assert_eq!(discoverydb.local_topic_readers.len(), 1); assert_eq!(discoverydb.get_local_topic_readers(&topic).len(), 1); @@ -1137,7 +1159,12 @@ mod tests { }; // Add the second reader to the database and verify the info is updated - discoverydb.update_local_topic_reader(&dp, &topic, &reader2_ing, None); + discoverydb.update_local_topic_reader(DiscoveryDB::local_topic_reader_data( + &dp, + &topic, + &reader2_ing, + None, + )); assert_eq!(discoverydb.get_local_topic_readers(&topic).len(), 2); assert_eq!(discoverydb.get_all_local_topic_readers().count(), 2); } diff --git a/tests/lock_order_discovery_db.rs b/tests/lock_order_discovery_db.rs new file mode 100644 index 00000000..cb7be8d2 --- /dev/null +++ b/tests/lock_order_discovery_db.rs @@ -0,0 +1,269 @@ +//! Regression test for an ABBA lock-order inversion between the two +//! per-participant locks of a single `DomainParticipant`: +//! +//! * `dpi` : `Arc>` +//! * `discovery_db` : `Arc>` +//! +//! Before the fix, endpoint creation took `discovery_db` -> `dpi`: +//! `InnerPublisher::create_datawriter` held `discovery_db.write()` across the +//! `DiscoveredWriterData::new(..)` call, which reached back into the +//! participant (`domain_id()`, `participant_id()`, `only_networks()`, +//! `guid()`) and therefore locked `dpi`. `InnerSubscriber::create_datareader` +//! had the same shape via `DiscoveryDB::update_local_topic_reader`. +//! +//! The discovery accessors took them in the opposite order, `dpi` -> +//! `discovery_db`: `DomainParticipant::discovered_readers()` (and +//! `discovered_writers`, `discovered_topics`, `find_topic`) locked `dpi` for +//! the whole statement and took `discovery_db.read()` inside it. +//! +//! Two threads on the *same* participant - one creating endpoints, one polling +//! the discovery accessors - could therefore deadlock: the creating thread held +//! `discovery_db.write()` and waited for `dpi`, while the polling thread held +//! `dpi` and waited for `discovery_db.read()`. +//! +//! Endpoint creation and the three snapshot accessors now avoid holding the +//! two locks together. +//! +//! Because a deadlocked pair would wedge the libtest harness forever, the +//! racing threads run in a *child process*: this same test binary re-executed +//! with `RUSTDDS_LOCK_ORDER_CHILD=1`. The parent waits with a deadline, kills +//! the child if it hangs, and fails instead of hanging. + +use std::{ + env, + process::{Command, Stdio}, + sync::{ + atomic::{AtomicUsize, Ordering}, + Arc, Barrier, + }, + thread, + time::{Duration, Instant}, +}; + +use rustdds::{policy, DomainParticipant, QosPolicies, QosPolicyBuilder, TopicKind}; +use serde::{Deserialize, Serialize}; + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq)] +struct Sample { + seq: u32, +} + +/// Set in the re-executed child process that actually runs the racing threads. +const CHILD_ENV: &str = "RUSTDDS_LOCK_ORDER_CHILD"; +/// Domain id handed down to the child, so parent and child agree on it. +const DOMAIN_ENV: &str = "RUSTDDS_LOCK_ORDER_DOMAIN"; +/// Must match the name of the `#[test]` function below. +const RACE_TEST_NAME: &str = "endpoint_creation_races_discovery_accessors"; +/// Proves the child completed the race; libtest also exits with 0 if no test +/// matched. +const CHILD_SUCCESS_EXIT_CODE: i32 = 42; + +/// How long the parent waits for the child before declaring a deadlock. +const CHILD_DEADLINE: Duration = Duration::from_secs(30); +/// Threads that create endpoints. +const CREATOR_THREADS: usize = 2; +/// Threads that poll discovery snapshots. +const ACCESSOR_THREADS: usize = 2; +/// Writer + reader pairs created per creator thread. +const ENDPOINTS_PER_CREATOR: usize = 150; + +fn test_qos() -> QosPolicies { + QosPolicyBuilder::new() + .reliability(policy::Reliability::BestEffort) + .durability(policy::Durability::Volatile) + .history(policy::History::KeepLast { depth: 1 }) + .build() +} + +/// Pick a domain id that is unlikely to collide with other DDS traffic on this +/// machine, while staying inside the range where the RTPS port arithmetic +/// (`7400 + 250 * domain_id + ..`) still fits in a `u16`. +fn scratch_domain_id(base: u16, span: u32) -> u16 { + base + (std::process::id() % span) as u16 +} + +/// The body that runs in the child process: `CREATOR_THREADS` threads creating +/// endpoints and `ACCESSOR_THREADS` threads polling the discovery accessors, +/// all on one shared `DomainParticipant`. +fn run_race() { + let domain_id: u16 = env::var(DOMAIN_ENV) + .expect("child started without a domain id") + .parse() + .expect("domain id is not a u16"); + + let qos = test_qos(); + let dp = DomainParticipant::new(domain_id).expect("failed to create DomainParticipant"); + + // Complete fallible setup before starting any barrier workers. A setup panic + // must fail the child instead of stranding another worker at the barrier. + let creators: Vec<_> = (0..CREATOR_THREADS) + .map(|_| { + let publisher = dp + .create_publisher(&qos) + .expect("failed to create Publisher"); + let subscriber = dp + .create_subscriber(&qos) + .expect("failed to create Subscriber"); + (publisher, subscriber) + }) + .collect(); + + // Counts the creator threads that have finished, so the accessor threads know + // when to stop. Start only after all workers have completed their setup; + // the lock interleavings themselves are still determined by the scheduler. + let creators_done = Arc::new(AtomicUsize::new(0)); + let start = Arc::new(Barrier::new(CREATOR_THREADS + ACCESSOR_THREADS)); + let mut threads = Vec::with_capacity(CREATOR_THREADS + ACCESSOR_THREADS); + + for (creator, (publisher, subscriber)) in creators.into_iter().enumerate() { + let dp = dp.clone(); + let qos = qos.clone(); + let creators_done = Arc::clone(&creators_done); + let start = Arc::clone(&start); + threads.push(thread::spawn(move || { + start.wait(); + for i in 0..ENDPOINTS_PER_CREATOR { + let topic = dp + .create_topic( + format!("lock_order_topic_{creator}_{i}"), + "Sample".to_string(), + &qos, + TopicKind::NoKey, + ) + .expect("failed to create Topic"); + + // Both constructors used to hold `discovery_db.write()` while taking + // `dpi`. Drop the endpoints at each iteration to bound resource use. + let _writer = publisher + .create_datawriter_no_key_cdr::(&topic, None) + .expect("failed to create DataWriter"); + let _reader = subscriber + .create_datareader_no_key_cdr::(&topic, None) + .expect("failed to create DataReader"); + } + + creators_done.fetch_add(1, Ordering::SeqCst); + })); + } + + for _ in 0..ACCESSOR_THREADS { + let dp = dp.clone(); + let creators_done = Arc::clone(&creators_done); + let start = Arc::clone(&start); + threads.push(thread::spawn(move || { + start.wait(); + // Exercise every accessor at least once, even if this thread is delayed + // until after the creators finish. + loop { + std::hint::black_box(dp.discovered_readers()); + std::hint::black_box(dp.discovered_writers()); + std::hint::black_box(dp.discovered_topics()); + + if creators_done.load(Ordering::SeqCst) == CREATOR_THREADS { + break; + } + } + })); + } + + for t in threads { + t.join().expect("a racing thread panicked"); + } +} + +/// Reproduces the lock-order inversion in a child process and fails - rather +/// than hanging - if the child deadlocks. +#[test] +fn endpoint_creation_races_discovery_accessors() { + if env::var_os(CHILD_ENV).is_some() { + run_race(); + std::process::exit(CHILD_SUCCESS_EXIT_CODE); + } + + let exe = env::current_exe().expect("cannot locate the test executable"); + let domain_id = scratch_domain_id(200, 20); + + let started = Instant::now(); + let mut child = Command::new(&exe) + .arg(RACE_TEST_NAME) + .arg("--exact") + .arg("--nocapture") + .arg("--test-threads=1") + .env(CHILD_ENV, "1") + .env(DOMAIN_ENV, domain_id.to_string()) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .spawn() + .expect("cannot re-execute the test binary"); + + let deadline = started + CHILD_DEADLINE; + let status = loop { + match child.try_wait().expect("cannot poll the child process") { + Some(status) => break Some(status), + None if Instant::now() >= deadline => break None, + // Polling interval only; it does not affect whether the test passes, + // just how promptly the parent notices that the child is done. + None => thread::sleep(Duration::from_millis(20)), + } + }; + + match status { + Some(status) if status.code() == Some(CHILD_SUCCESS_EXIT_CODE) => (), + Some(status) => panic!( + "the child did not report race completion (domain {domain_id}, exit {status}, expected exit \ + code {CHILD_SUCCESS_EXIT_CODE}) after {:.1} s", + started.elapsed().as_secs_f64() + ), + None => { + let _ = child.kill(); + let _ = child.wait(); + panic!( + "The racing child process (domain {domain_id}) did not finish within {} s and was killed. \ + Possible lock-order deadlock between endpoint creation and discovery snapshots.", + CHILD_DEADLINE.as_secs() + ) + } + } +} + +/// Single-threaded sanity check that documents the API exercised above. It +/// passes both before and after the fix, because one thread can never hold both +/// locks in conflicting orders. +#[test] +fn single_thread_create_writer_then_query_discovery() { + let qos = test_qos(); + let dp = + DomainParticipant::new(scratch_domain_id(220, 10)).expect("failed to create participant"); + + let topic = dp + .create_topic( + "lock_order_sanity_topic".to_string(), + "Sample".to_string(), + &qos, + TopicKind::NoKey, + ) + .expect("failed to create Topic"); + let publisher = dp + .create_publisher(&qos) + .expect("failed to create Publisher"); + let _writer = publisher + .create_datawriter_no_key_cdr::(&topic, None) + .expect("failed to create DataWriter"); + + // The accessors must return without deadlocking or panicking. The remote + // sets may legitimately be empty - nothing else is required to be on the + // network - so only the local topic is asserted on. + let _ = dp.discovered_readers(); + let _ = dp.discovered_writers(); + let topics = dp.discovered_topics(); + assert!( + topics + .iter() + .any(|t| t.topic_data.name == "lock_order_sanity_topic"), + "the locally created topic should show up in discovered_topics(), got {:?}", + topics + .iter() + .map(|t| t.topic_data.name.clone()) + .collect::>() + ); +}