#![cfg(all(target_os = "linux", feature = "std"))]
#![allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
use std::sync::Arc;
use std::time::{Duration, Instant};
use zerodds_dcps::dds_type::RawBytes;
use zerodds_dcps::runtime::{DcpsRuntime, RuntimeConfig, UserReaderConfig, UserWriterConfig};
use zerodds_qos::{
DeadlineQosPolicy, DurabilityKind, LifespanQosPolicy, LivelinessQosPolicy, OwnershipKind,
};
use zerodds_rtps::wire_types::GuidPrefix;
fn wait_for_matched(
rt: &Arc<DcpsRuntime>,
eid: zerodds_rtps::wire_types::EntityId,
n: usize,
timeout: Duration,
) -> bool {
let deadline = Instant::now() + timeout;
while Instant::now() < deadline {
if rt.user_reader_matched_count(eid) >= n {
return true;
}
std::thread::sleep(Duration::from_millis(20));
}
false
}
fn wait_for_writer_matched(
rt: &Arc<DcpsRuntime>,
eid: zerodds_rtps::wire_types::EntityId,
n: usize,
timeout: Duration,
) -> bool {
let deadline = Instant::now() + timeout;
while Instant::now() < deadline {
if rt.user_writer_matched_count(eid) >= n {
return true;
}
std::thread::sleep(Duration::from_millis(20));
}
false
}
fn wait_for_peers(rt: &Arc<DcpsRuntime>, n: usize, timeout: Duration) -> bool {
let deadline = Instant::now() + timeout;
while Instant::now() < deadline {
if rt.discovered_participants().len() >= n {
return true;
}
std::thread::sleep(Duration::from_millis(20));
}
false
}
fn make_writer_config(topic: &str) -> UserWriterConfig {
UserWriterConfig {
topic_name: topic.into(),
type_name: "RawBytes".into(),
reliable: true,
durability: DurabilityKind::Volatile,
deadline: DeadlineQosPolicy::default(),
lifespan: LifespanQosPolicy::default(),
liveliness: LivelinessQosPolicy::default(),
ownership: OwnershipKind::Shared,
ownership_strength: 0,
partition: vec![],
user_data: vec![],
topic_data: vec![],
group_data: vec![],
type_identifier: zerodds_types::TypeIdentifier::None,
data_representation_offer: None,
}
}
fn make_reader_config(topic: &str) -> UserReaderConfig {
UserReaderConfig {
topic_name: topic.into(),
type_name: "RawBytes".into(),
reliable: true,
durability: DurabilityKind::Volatile,
deadline: DeadlineQosPolicy::default(),
liveliness: LivelinessQosPolicy::default(),
ownership: OwnershipKind::Shared,
partition: vec![],
user_data: vec![],
topic_data: vec![],
group_data: vec![],
type_identifier: zerodds_types::TypeIdentifier::None,
type_consistency: zerodds_types::qos::TypeConsistencyEnforcement::default(),
data_representation_offer: None,
}
}
#[test]
fn e2e_two_runtimes_udp_roundtrip() {
let domain = 117;
let prefix_a = GuidPrefix::from_bytes([0xA1; 12]);
let prefix_b = {
let mut p = [0xB2; 12];
p[..4].copy_from_slice(&prefix_a.to_bytes()[..4]);
GuidPrefix::from_bytes(p)
};
let rt_a = Arc::new(
DcpsRuntime::start(domain, prefix_a, RuntimeConfig::default()).expect("rt_a start"),
);
let rt_b = Arc::new(
DcpsRuntime::start(domain, prefix_b, RuntimeConfig::default()).expect("rt_b start"),
);
assert!(
wait_for_peers(&rt_a, 1, Duration::from_secs(10)),
"rt_a did not see rt_b via SPDP"
);
assert!(
wait_for_peers(&rt_b, 1, Duration::from_secs(10)),
"rt_b did not see rt_a via SPDP"
);
let writer_eid = rt_a
.register_user_writer(make_writer_config("E2EUdp"))
.unwrap();
let (reader_eid, rx) = rt_b
.register_user_reader(make_reader_config("E2EUdp"))
.unwrap();
assert!(
wait_for_matched(&rt_b, reader_eid, 1, Duration::from_secs(10)),
"reader did not see writer via SEDP"
);
assert!(
wait_for_writer_matched(&rt_a, writer_eid, 1, Duration::from_secs(10)),
"writer did not see reader via SEDP (rt_a side asymmetric match)"
);
let drops_before = rt_b.user_reader_unknown_src_count(reader_eid);
let payload = b"udp-roundtrip-1234567890";
rt_a.write_user_sample(writer_eid, payload.to_vec())
.unwrap();
let sample = rx
.recv_timeout(Duration::from_secs(2))
.expect("UDP path: a single write must suffice (no retry needed)");
let drops_after = rt_b.user_reader_unknown_src_count(reader_eid);
assert_eq!(
drops_after - drops_before,
0,
"unknown_src_count must stay 0 with symmetric wait"
);
match sample {
zerodds_dcps::runtime::UserSample::Alive { payload: bytes, .. } => {
assert_eq!(bytes.as_slice(), payload);
}
other => panic!(
"expected Alive sample, got {:?}",
std::mem::discriminant(&other)
),
}
}
#[cfg(feature = "same-host-shm")]
#[test]
fn e2e_two_runtimes_shm_roundtrip() {
let _ = std::fs::remove_dir_all("/tmp/zerodds-shm");
let domain = 118;
let prefix_a = GuidPrefix::from_bytes([0xA3; 12]);
let prefix_b = {
let mut p = [0xB4; 12];
p[..4].copy_from_slice(&prefix_a.to_bytes()[..4]);
GuidPrefix::from_bytes(p)
};
let rt_a = Arc::new(
DcpsRuntime::start(domain, prefix_a, RuntimeConfig::default()).expect("rt_a start"),
);
let rt_b = Arc::new(
DcpsRuntime::start(domain, prefix_b, RuntimeConfig::default()).expect("rt_b start"),
);
assert!(wait_for_peers(&rt_a, 1, Duration::from_secs(10)));
assert!(wait_for_peers(&rt_b, 1, Duration::from_secs(10)));
let writer_eid = rt_a
.register_user_writer(make_writer_config("E2EShm"))
.unwrap();
let (reader_eid, rx) = rt_b
.register_user_reader(make_reader_config("E2EShm"))
.unwrap();
assert!(wait_for_matched(
&rt_b,
reader_eid,
1,
Duration::from_secs(10)
));
assert!(wait_for_writer_matched(
&rt_a,
writer_eid,
1,
Duration::from_secs(10)
));
let payload = b"shm-roundtrip-9876543210";
rt_a.write_user_sample(writer_eid, payload.to_vec())
.unwrap();
let sample = rx
.recv_timeout(Duration::from_secs(2))
.expect("SHM path: a single write must suffice (no retry needed)");
match sample {
zerodds_dcps::runtime::UserSample::Alive { payload: bytes, .. } => {
assert_eq!(bytes.as_slice(), payload);
}
other => panic!(
"expected Alive sample, got {:?}",
std::mem::discriminant(&other)
),
}
}
#[test]
fn spdp_works_after_runtime_drop() {
let domain = 119;
let prefix = GuidPrefix::from_bytes([0xC5; 12]);
{
let rt = Arc::new(
DcpsRuntime::start(domain, prefix, RuntimeConfig::default()).expect("rt1 start"),
);
std::thread::sleep(Duration::from_millis(500));
drop(rt);
}
std::thread::sleep(Duration::from_millis(500));
let rt_a =
Arc::new(DcpsRuntime::start(domain, prefix, RuntimeConfig::default()).expect("rt_a start"));
let prefix_b = {
let mut p = [0xD6; 12];
p[..4].copy_from_slice(&prefix.to_bytes()[..4]);
GuidPrefix::from_bytes(p)
};
let rt_b = Arc::new(
DcpsRuntime::start(domain, prefix_b, RuntimeConfig::default()).expect("rt_b start"),
);
assert!(
wait_for_peers(&rt_a, 1, Duration::from_secs(15)),
"rt_a did not see rt_b after a prior runtime drop — SPDP-Membership-Leak"
);
assert!(
wait_for_peers(&rt_b, 1, Duration::from_secs(15)),
"rt_b did not see rt_a after a prior runtime drop — SPDP-Membership-Leak"
);
}
#[cfg(feature = "same-host-shm")]
#[test]
fn e2e_cross_vendor_different_host_id_no_shm_bind() {
use zerodds_dcps::same_host::SameHostState;
let _ = std::fs::remove_dir_all("/tmp/zerodds-shm");
let domain = 120;
let prefix_a = GuidPrefix::from_bytes([
0x5A, 0xE7, 0x0D, 0xD5, 0xAA, 0xBB, 0xCC, 0xDD, 0x11, 0x22, 0x33, 0x44, ]);
let prefix_b = GuidPrefix::from_bytes([
0xC1, 0xC1, 0x0E, 0x00, 0x55, 0x66, 0x77, 0x88, 0x99, 0xAA, 0xBB, 0xCC, ]);
assert!(
!prefix_a.is_same_host(prefix_b),
"test precondition: prefixes must have different host_id"
);
let rt_a = Arc::new(
DcpsRuntime::start(domain, prefix_a, RuntimeConfig::default()).expect("rt_a start"),
);
let rt_b = Arc::new(
DcpsRuntime::start(domain, prefix_b, RuntimeConfig::default()).expect("rt_b start"),
);
assert!(wait_for_peers(&rt_a, 1, Duration::from_secs(10)));
assert!(wait_for_peers(&rt_b, 1, Duration::from_secs(10)));
let writer_eid = rt_a
.register_user_writer(make_writer_config("CrossVendor"))
.unwrap();
let (reader_eid, rx) = rt_b
.register_user_reader(make_reader_config("CrossVendor"))
.unwrap();
assert!(wait_for_matched(
&rt_b,
reader_eid,
1,
Duration::from_secs(10)
));
assert!(wait_for_writer_matched(
&rt_a,
writer_eid,
1,
Duration::from_secs(10)
));
let snapshot_a = rt_a.same_host.snapshot();
for (_, _, state) in &snapshot_a {
assert!(
!matches!(state, SameHostState::Bound { .. }),
"rt_a must have no bound state with a cross-host-id peer"
);
}
let snapshot_b = rt_b.same_host.snapshot();
for (_, _, state) in &snapshot_b {
assert!(
!matches!(state, SameHostState::Bound { .. }),
"rt_b must have no bound state with a cross-host-id peer"
);
}
let payload = b"cross-vendor-udp-only";
rt_a.write_user_sample(writer_eid, payload.to_vec())
.unwrap();
let sample = rx
.recv_timeout(Duration::from_secs(2))
.expect("UDP path must keep working with different host_ids");
match sample {
zerodds_dcps::runtime::UserSample::Alive { payload: bytes, .. } => {
assert_eq!(bytes.as_slice(), payload);
}
other => panic!(
"expected Alive sample, got {:?}",
std::mem::discriminant(&other)
),
}
}
#[allow(dead_code)]
fn _silence(_: RawBytes) {}