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
7 changes: 5 additions & 2 deletions .github/workflows/1_run_interoperability_tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -6,11 +6,11 @@ on:
publishers:
description: Publishers to use
type: string
default: '["connext_dds","toc_coredx_dds", "eclipse_cyclone", "opendds"]'
default: '["connext_dds","toc_coredx_dds", "eclipse_cyclone", "opendds", "dust_dds"]'
subscribers:
description: Subscribers to use
type: string
default: '["connext_dds","toc_coredx_dds", "eclipse_cyclone", "opendds"]'
default: '["connext_dds","toc_coredx_dds", "eclipse_cyclone", "opendds", "dust_dds"]'
jobs:
generate_timestamp:
runs-on: ubuntu-latest
Expand Down Expand Up @@ -81,6 +81,9 @@ jobs:
if [[ "${publisher_exe,,}" == *opendds* || "${subscriber_exe,,}" == *opendds* ]]; then
xcdr1_enabled=false
fi
if [[ "${publisher_exe,,}" == *dust_dds* || "${subscriber_exe,,}" == *dust_dds* ]]; then
to1_enabled=false
fi
if [[ "$xcdr1_enabled" == true ]]; then
if [[ "$to1_enabled" == true ]]; then
# Run interoperability tests (xcdr1, type_object_version 1)
Expand Down
4 changes: 1 addition & 3 deletions src/rs/DustDDS/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,4 @@ rust-version = "1.87"
[dependencies]
clap = { version = "4.5.47", features = ["derive", "string"] }
ctrlc = "3.4"
dust_dds = { version = "0.16.0", git = "https://github.com/s2e-systems/dust-dds", branch = "main", features = [
"xtypes-xml",
] }
dust_dds = { version = "0.16.0", features = ["xtypes-xml"] }
234 changes: 140 additions & 94 deletions src/rs/DustDDS/src/main.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
use clap::{Parser, ValueEnum};
use ctrlc;
use dust_dds::{
configuration::DustDdsConfigurationBuilder,
dds_async::topic_description::TopicDescriptionAsync,
domain::{
domain_participant::DomainParticipant,
domain_participant_factory::DomainParticipantFactory,
Expand All @@ -14,18 +16,17 @@ use dust_dds::{
self, DataRepresentationQosPolicy, DurabilityQosPolicy, HistoryQosPolicy,
HistoryQosPolicyKind, OwnershipQosPolicy, OwnershipQosPolicyKind,
OwnershipStrengthQosPolicy, PartitionQosPolicy, ReliabilityQosPolicy,
XCDR_DATA_REPRESENTATION, XCDR2_DATA_REPRESENTATION,
TypeConsistencyEnforcementQosPolicy, TypeConsistencyKind, XCDR2_DATA_REPRESENTATION,
XCDR_DATA_REPRESENTATION,
},
sample_info::{ANY_INSTANCE_STATE, ANY_SAMPLE_STATE, ANY_VIEW_STATE},
status::{NO_STATUS, StatusKind},
status::{StatusKind, NO_STATUS},
time::DurationKind,
},
publication::data_writer::DataWriter,
subscription::data_reader::DataReader,
xtypes::{
dynamic_type::DynamicData,
dynamic_type::DynamicDataFactory,
dynamic_type::{DynamicType, DynamicTypeBuilderFactory},
xtypes::dynamic_type::{
DynamicData, DynamicDataFactory, DynamicType, DynamicTypeBuilderFactory,
},
};
use std::{
Expand Down Expand Up @@ -270,16 +271,77 @@ impl Options {
}

struct Listener;
impl DomainParticipantListener for Listener {}
impl DomainParticipantListener for Listener {
fn on_publication_matched(
&mut self,
the_writer: dust_dds::dds_async::data_writer::DataWriterAsync<()>,
status: dust_dds::infrastructure::status::PublicationMatchedStatus,
) -> impl Future<Output = ()> + Send {
let topic_name = the_writer.get_topic().get_name();
let type_name = the_writer.get_topic().get_type_name();
println!(
"on_publication_matched() topic: '{}' type: '{}' : matched readers {} (change = {})",
topic_name, type_name, status.current_count, status.current_count_change
);
core::future::ready(())
}

fn on_subscription_matched(
&mut self,
the_reader: dust_dds::dds_async::data_reader::DataReaderAsync<()>,
status: dust_dds::infrastructure::status::SubscriptionMatchedStatus,
) -> impl Future<Output = ()> + Send {
let topic_name = the_reader.get_topicdescription().get_name();
let type_name = the_reader.get_topicdescription().get_type_name();
println!(
"on_subscription_matched() topic: '{}' type: '{}' : matched writers {} (change = {})",
topic_name, type_name, status.current_count, status.current_count_change
);
core::future::ready(())
}

fn on_liveliness_changed(
&mut self,
the_reader: dust_dds::dds_async::data_reader::DataReaderAsync<()>,
status: dust_dds::infrastructure::status::LivelinessChangedStatus,
) -> impl Future<Output = ()> + Send {
let topic_name = the_reader.get_topicdescription().get_name();
let type_name = the_reader.get_topicdescription().get_type_name();
println!(
"on_liveliness_changed() topic: '{}' type: '{}' : (alive = {}, not_alive = {}",
topic_name, type_name, status.alive_count, status.not_alive_count
);
core::future::ready(())
}

fn on_inconsistent_topic(
&mut self,
the_topic: dust_dds::dds_async::topic::TopicAsync,
_status: dust_dds::infrastructure::status::InconsistentTopicStatus,
) -> impl Future<Output = ()> + Send {
println!(
"on_inconsistent_topic() topic: '{}' type: '{}'",
the_topic.get_name(),
the_topic.get_type_name()
);
core::future::ready(())
}
}

fn init_publisher(
participant: &DomainParticipant,
options: Options,
dynamic_type: DynamicType,
) -> Result<DataWriter<DynamicData>, InitializeError> {
dynamic_type: DynamicType<'static>,
) -> Result<DataWriter<DynamicData<'static>>, InitializeError> {
let topic_name = options.topic_name.clone().unwrap_or("test".to_string());
let type_name = options.type_name.clone().unwrap();

println!("Create topic: {}", topic_name);
println!(
"Create writer for topic: {} type: {}",
topic_name, type_name
);

let topic = participant.create_dynamic_topic(
&topic_name,
&type_name,
Expand All @@ -289,8 +351,6 @@ fn init_publisher(
dynamic_type,
)?;

println!("Create topic: {}", topic_name);

let publisher_qos = QosKind::Specific(PublisherQos {
partition: options.partition_qos_policy(),
..Default::default()
Expand All @@ -313,8 +373,6 @@ fn init_publisher(
data_writer_qos.ownership_strength = options.ownership_strength_qos_policy();
}

println!("Create writer for topic: {} type: {}", topic_name, type_name);

let data_writer = publisher.create_datawriter::<DynamicData>(
&topic,
QosKind::Specific(data_writer_qos),
Expand All @@ -326,9 +384,9 @@ fn init_publisher(
}

fn run_publisher(
data_writer: &DataWriter<DynamicData>,
data_writer: &DataWriter<DynamicData<'static>>,
options: Options,
dynamic_type: DynamicType,
dynamic_type: DynamicType<'static>,
all_done: Receiver<()>,
) -> Result<(), RunningError> {
let mut dd = DynamicDataFactory::create_data(dynamic_type);
Expand All @@ -345,6 +403,7 @@ fn run_publisher(
while all_done.try_recv().is_err() {
if options.print_writer_samples {
println!(" Wrote:");
println!("{:?}", dd);
}
// Write dynamic data
data_writer.write(dd.clone(), None).ok();
Expand All @@ -356,11 +415,14 @@ fn run_publisher(
fn init_subscriber(
participant: &DomainParticipant,
options: Options,
dynamic_type: DynamicType,
) -> Result<DataReader<DynamicData>, InitializeError> {
dynamic_type: DynamicType<'static>,
) -> Result<DataReader<DynamicData<'static>>, InitializeError> {
let topic_name = options.topic_name.clone().unwrap_or("test".to_string());
let type_name = options.type_name.clone().unwrap();

println!("Create topic: {}", topic_name);
println!("Create reader for topic: {}", topic_name);

let topic = participant
.create_dynamic_topic(
&topic_name,
Expand All @@ -372,8 +434,6 @@ fn init_subscriber(
)
.unwrap();

println!("Create topic: {}", topic_name);

let subscriber_qos = QosKind::Specific(SubscriberQos {
partition: options.partition_qos_policy(),
..Default::default()
Expand All @@ -398,74 +458,57 @@ fn init_subscriber(
);
}

// // Set type consistency enforcement based on arguments
// let mut type_consistency = TypeConsistencyEnforcementQosPolicy::default();

// if let Some(allow_type_coercion) = options.allow_type_coercion {
// match allow_type_coercion {
// TypeConsistencyArg::T => type_consistency.kind = TypeConsistencyKind::AllowTypeCoercion,
// TypeConsistencyArg::F => {
// type_consistency.kind = TypeConsistencyKind::DisallowTypeCoercion
// }
// TypeConsistencyArg::D => {
// type_consistency.kind = TypeConsistencyEnforcementQosPolicy::default().kind
// }
// }
// }
// if let Some(force_type_validation) = options.force_type_validation {
// match force_type_validation {
// TypeConsistencyArg::T => type_consistency.force_type_validation = true,
// TypeConsistencyArg::F => type_consistency.force_type_validation = false,
// TypeConsistencyArg::D => {
// type_consistency.force_type_validation =
// TypeConsistencyEnforcementQosPolicy::default().force_type_validation
// }
// }
// }
// if let Some(ignore_member_names) = options.ignore_member_names {
// match ignore_member_names {
// TypeConsistencyArg::T => type_consistency.ignore_member_names = true,
// TypeConsistencyArg::F => type_consistency.ignore_member_names = false,
// TypeConsistencyArg::D => {
// type_consistency.ignore_member_names =
// TypeConsistencyEnforcementQosPolicy::default().ignore_member_names
// }
// }
// }
// if let Some(ignore_seq_bounds) = options.ignore_seq_bounds {
// match ignore_seq_bounds {
// TypeConsistencyArg::T => type_consistency.ignore_sequence_bounds = true,
// TypeConsistencyArg::F => type_consistency.ignore_sequence_bounds = false,
// TypeConsistencyArg::D => {
// type_consistency.ignore_sequence_bounds =
// TypeConsistencyEnforcementQosPolicy::default().ignore_sequence_bounds
// }
// }
// }
// if let Some(ignore_str_bounds) = options.ignore_str_bounds {
// match ignore_str_bounds {
// TypeConsistencyArg::T => type_consistency.ignore_string_bounds = true,
// TypeConsistencyArg::F => type_consistency.ignore_string_bounds = false,
// TypeConsistencyArg::D => {
// type_consistency.ignore_string_bounds =
// TypeConsistencyEnforcementQosPolicy::default().ignore_string_bounds
// }
// }
// }
// if let Some(prevent_type_widening) = options.prevent_type_widening {
// match prevent_type_widening {
// TypeConsistencyArg::T => type_consistency.prevent_type_widening = true,
// TypeConsistencyArg::F => type_consistency.prevent_type_widening = false,
// TypeConsistencyArg::D => {
// type_consistency.prevent_type_widening =
// TypeConsistencyEnforcementQosPolicy::default().prevent_type_widening
// }
// }
// }

// data_reader_qos.type_consistency = type_consistency;
// Set type consistency enforcement based on arguments
let mut type_consistency = TypeConsistencyEnforcementQosPolicy::default();
// Note: The default of the DDS XTypes standard is false
type_consistency.ignore_member_names = true;

println!("Create reader for topic: {}", topic_name);
if let Some(allow_type_coercion) = options.allow_type_coercion {
match allow_type_coercion {
TypeConsistencyArg::T => type_consistency.kind = TypeConsistencyKind::AllowTypeCoercion,
TypeConsistencyArg::F => {
type_consistency.kind = TypeConsistencyKind::DisallowTypeCoercion
}
TypeConsistencyArg::D => (),
}
}
if let Some(force_type_validation) = options.force_type_validation {
match force_type_validation {
TypeConsistencyArg::T => type_consistency.force_type_validation = true,
TypeConsistencyArg::F => type_consistency.force_type_validation = false,
TypeConsistencyArg::D => (),
}
}
if let Some(ignore_member_names) = options.ignore_member_names {
match ignore_member_names {
TypeConsistencyArg::T => type_consistency.ignore_member_names = true,
TypeConsistencyArg::F => type_consistency.ignore_member_names = false,
TypeConsistencyArg::D => type_consistency.ignore_member_names = false,
}
}
if let Some(ignore_seq_bounds) = options.ignore_seq_bounds {
match ignore_seq_bounds {
TypeConsistencyArg::T => type_consistency.ignore_sequence_bounds = true,
TypeConsistencyArg::F => type_consistency.ignore_sequence_bounds = false,
TypeConsistencyArg::D => (),
}
}
if let Some(ignore_str_bounds) = options.ignore_str_bounds {
match ignore_str_bounds {
TypeConsistencyArg::T => type_consistency.ignore_string_bounds = true,
TypeConsistencyArg::F => type_consistency.ignore_string_bounds = false,
TypeConsistencyArg::D => (),
}
}
if let Some(prevent_type_widening) = options.prevent_type_widening {
match prevent_type_widening {
TypeConsistencyArg::T => type_consistency.prevent_type_widening = true,
TypeConsistencyArg::F => type_consistency.prevent_type_widening = false,
TypeConsistencyArg::D => (),
}
}

data_reader_qos.type_consistency = type_consistency;

let data_reader = subscriber.create_datareader::<DynamicData>(
&topic,
Expand All @@ -478,9 +521,9 @@ fn init_subscriber(
}

fn run_subscriber(
data_reader: &DataReader<DynamicData>,
data_reader: &DataReader<DynamicData<'static>>,
options: Options,
dynamic_type: DynamicType,
dynamic_type: DynamicType<'static>,
all_done: Receiver<()>,
) -> Result<(), RunningError> {
let mut expected_data = None;
Expand Down Expand Up @@ -508,7 +551,7 @@ fn run_subscriber(
match read_result {
Ok(samples) => {
for sample in samples {
if sample.sample_info.valid_data {
if sample.sample_info.valid_data && sample.data.is_some() {
println!("sample_received()");
if let Some(expected) = &expected_data {
if sample.data.as_ref() == Some(expected) {
Expand All @@ -520,10 +563,10 @@ fn run_subscriber(
}
previous_handle = Some(sample.sample_info.instance_handle);
}
std::thread::sleep(std::time::Duration::from_millis(100));
}
Err(_) => break,
}
std::thread::sleep(std::time::Duration::from_millis(100));
}
}
Ok(())
Expand All @@ -532,7 +575,12 @@ fn run_subscriber(
fn initialize(options: &Options) -> Result<DomainParticipant, InitializeError> {
let participant_factory = DomainParticipantFactory::get_instance();

// Set domain participant factory QoS if needed
if options.disable_type_info {
let configuration = DustDdsConfigurationBuilder::new()
.enable_type_information(false)
.build()?;
*participant_factory.get_mut_configuration() = configuration;
}

let participant = participant_factory.create_participant(
options.domain_id,
Expand Down Expand Up @@ -616,14 +664,12 @@ fn main() -> Result<(), Return> {
})?;

// Create the type
let mut dt: Option<DynamicType> = None;
let mut dt = None;
if let (Some(type_folder), Some(type_file), Some(type_name)) =
(&options.type_folder, &options.type_file, &options.type_name)
{
let file_path = format!("{}/xml/{}.xml", type_folder, type_file);
let type_xml = std::fs::read_to_string(file_path).unwrap();
// This function is unimplemented in dust_dds currently
println!("type_name {type_name}");
let type_builder =
DynamicTypeBuilderFactory::create_type_w_document(&type_xml, type_name, vec![])
.unwrap();
Expand Down
Loading
Loading