Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
59 changes: 19 additions & 40 deletions src/dds/participant.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
Expand Down Expand Up @@ -531,7 +531,12 @@ impl DomainParticipant {
/// }
/// ```
pub fn discovered_topics(&self) -> Vec<DiscoveredTopicData> {
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.
Expand All @@ -553,7 +558,12 @@ impl DomainParticipant {
/// }
/// ```
pub fn discovered_readers(&self) -> Vec<DiscoveredReaderData> {
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.
Expand All @@ -575,7 +585,12 @@ impl DomainParticipant {
/// }
/// ```
pub fn discovered_writers(&self) -> Vec<DiscoveredWriterData> {
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
Expand Down Expand Up @@ -992,18 +1007,6 @@ impl DomainParticipantDisc {
self.dpi.participant_id()
}

pub fn discovered_topics(&self) -> Vec<DiscoveredTopicData> {
self.dpi.discovered_topics()
}

pub fn discovered_readers(&self) -> Vec<DiscoveredReaderData> {
self.dpi.discovered_readers()
}

pub fn discovered_writers(&self) -> Vec<DiscoveredWriterData> {
self.dpi.discovered_writers()
}

pub(crate) fn dds_cache(&self) -> Arc<RwLock<DDSCache>> {
self.dpi.dds_cache()
}
Expand Down Expand Up @@ -1598,30 +1601,6 @@ impl DomainParticipantInner {
self.domain_info.participant_id
}

pub fn discovered_topics(&self) -> Vec<DiscoveredTopicData> {
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<DiscoveredReaderData> {
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<DiscoveredWriterData> {
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<DomainParticipantStatusEvent> {
Expand Down
64 changes: 42 additions & 22 deletions src/dds/pubsub.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down
49 changes: 38 additions & 11 deletions src/discovery/discovery_db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<EndpointSecurityInfo>,
) {
) -> DiscoveredReaderData {
let reader_guid = reader.guid;

let reader_proxy = RtpsReaderProxy::from_reader(reader, domain_participant);
Expand All @@ -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) {
Expand Down Expand Up @@ -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);

Expand Down Expand Up @@ -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);
}
Expand Down
Loading