Skip to main content

zerodds_dcps/
lib.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright 2026 ZeroDDS Contributors
3//! Crate `zerodds-dcps`. Safety classification: **STANDARD**.
4//!
5//! DCPS Public API (OMG DDS 1.4 §2.2.2): `DomainParticipant`,
6//! `Publisher`, `Subscriber`, `Topic`, `DataWriter`, `DataReader`.
7//!
8//! Spec: OMG DDS 1.4 §2.2 (Data-Centric Publish-Subscribe Module) +
9//! DDSI-RTPS 2.5 §8.5 (Discovery + WLP) + XTypes 1.3 §7.6.3
10//! (TypeLookup service wiring).
11//!
12//! ## Layer position
13//!
14//! Layer 4 — Core Services. Built on Layer 1
15//! (foundation/cdr/qos/types/time-service), Layer 2
16//! (rtps/discovery/transport-*), Layer 3 (idl/idl-rust/xml).
17//!
18//! ## Public API (as of 1.0.0-rc.1)
19//!
20//! - [`DomainParticipantFactory`] — singleton factory;
21//!   `create_participant` spawns a live runtime with UDP/SPDP/SEDP/WLP,
22//!   `create_participant_offline` builds an in-process skeleton without
23//!   networking for unit tests.
24//! - [`DomainParticipant`] — top-level entity; creates
25//!   publishers/subscribers/topics, maintains the built-in type
26//!   registry, exposes TypeLookup hooks and ignore filters.
27//! - [`Publisher`] / [`DataWriter`] — typed `Writer<T>` with a `DdsType`
28//!   bound; integrates the RTPS ReliableWriter (live) or an in-memory
29//!   queue (offline) plus the durability backend (DDS 1.4 §2.2.3.5).
30//! - [`Subscriber`] / [`DataReader`] — typed `Reader<T>` with
31//!   `take`/`read`/conditions, sample cache and InstanceState tracker
32//!   (DDS 1.4 §2.2.2.5).
33//! - [`Topic`] / [`ContentFilteredTopic`] / [`MultiTopic`] — the topic
34//!   hierarchy incl. SQL filters (DDS 1.4 §2.2.2.3).
35//! - Builtin topics: [`BuiltinSubscriber`] +
36//!   [`DcpsParticipantBuiltinTopicData`] /
37//!   [`DcpsPublicationBuiltinTopicData`] /
38//!   [`DcpsSubscriptionBuiltinTopicData`] /
39//!   [`DcpsTopicBuiltinTopicData`] (DDS 1.4 §2.2.5).
40//! - Conditions/WaitSet: [`Condition`] / [`ReadCondition`] /
41//!   [`QueryCondition`] / [`GuardCondition`] / [`WaitSet`].
42//! - QoS families: [`DomainParticipantQos`], [`PublisherQos`],
43//!   [`SubscriberQos`], [`TopicQos`], [`DataWriterQos`],
44//!   [`DataReaderQos`].
45//!
46//! ## Example
47//!
48//! ```
49//! use zerodds_dcps::*;
50//! let factory = DomainParticipantFactory::instance();
51//! // Offline mode for the doctest (no UDP multicast needed).
52//! let participant = factory.create_participant_offline(0, DomainParticipantQos::default());
53//! let topic = participant
54//!     .create_topic::<RawBytes>("Chatter", TopicQos::default())
55//!     .expect("create_topic");
56//! let publisher = participant.create_publisher(PublisherQos::default());
57//! let writer = publisher
58//!     .create_datawriter::<RawBytes>(&topic, DataWriterQos::default())
59//!     .expect("create_datawriter");
60//! writer.write(&RawBytes::new(vec![1, 2, 3])).expect("write");
61//! ```
62
63#![cfg_attr(not(feature = "std"), no_std)]
64#![deny(unsafe_code)]
65#![warn(missing_docs)]
66
67#[cfg(feature = "alloc")]
68extern crate alloc;
69
70pub mod builtin_subscriber;
71pub mod builtin_topics;
72pub mod coherent_set;
73#[cfg(feature = "std")]
74pub mod condition;
75/// Cross-vendor same-host zero-copy with Cyclone over iceoryx C++ (POSH).
76#[cfg(all(feature = "std", feature = "cyclone-iox", target_os = "linux"))]
77pub mod cyclone_iox_integration;
78pub mod dds_type;
79#[cfg(feature = "std")]
80pub mod durability_service;
81pub mod entity;
82pub mod error;
83pub mod factory;
84/// ADR-0005: opt-in flatdata integration.
85#[cfg(all(feature = "std", feature = "flatdata-integration"))]
86pub mod flatdata_integration;
87/// In-process discovery fastpath (same-process, same-domain).
88#[cfg(feature = "std")]
89mod inproc;
90pub mod instance_handle;
91#[cfg(feature = "std")]
92pub mod instance_tracker;
93#[cfg(feature = "alloc")]
94pub mod interop;
95/// Preference-ordered multi-transport for user traffic (SHM + UDP fallback).
96#[cfg(feature = "std")]
97pub mod layered_transport;
98pub mod listener;
99#[cfg(feature = "std")]
100pub mod listener_dispatch;
101#[cfg(feature = "metrics")]
102pub mod metrics;
103pub mod participant;
104pub mod psm_constants;
105pub mod publisher;
106pub mod qos;
107#[cfg(feature = "std")]
108pub mod runtime;
109pub mod same_host;
110/// Wave 4b.3: feature-gated bridge between the `same_host` tracker and
111/// `zerodds-transport-shm::PosixShmTransport`. Only compiled when the
112/// `same-host-shm` feature is active.
113#[cfg(all(feature = "std", feature = "same-host-shm"))]
114pub mod same_host_shm;
115/// 4b.5: alternative UDS datagram backend for same-host paths. Not true
116/// zero-copy but container-friendly. Only compiled when the
117/// `same-host-uds` feature is active.
118#[cfg(all(feature = "std", feature = "same-host-uds"))]
119pub mod same_host_uds;
120pub mod sample;
121pub mod sample_bytes;
122pub mod sample_info;
123/// D.5e Phase 3 — deadline-heap scheduler (event-driven replacement for the
124/// fixed-period tick poll). Std-only (mpsc channel + Instant park).
125#[cfg(feature = "std")]
126pub mod scheduler;
127/// Multi-peer SHM adapter for `user_unicast` transport selection. Wraps
128/// multiple `PosixShmTransport` pairs behind the transport trait.
129/// Feature-gated via `same-host-shm` because zerodds-transport-shm is
130/// optional.
131#[cfg(all(feature = "std", feature = "same-host-shm"))]
132pub mod shm_user;
133pub mod status;
134pub mod subscriber;
135pub mod time;
136pub mod topic;
137#[cfg(feature = "std")]
138pub mod wlp;
139
140// Flat re-exports for the typical import line.
141pub use builtin_subscriber::{BuiltinSinks, BuiltinSubscriber, BuiltinTopic, builtin_reader_qos};
142pub use builtin_topics::{
143    ParticipantBuiltinTopicData as DcpsParticipantBuiltinTopicData,
144    PublicationBuiltinTopicData as DcpsPublicationBuiltinTopicData,
145    SubscriptionBuiltinTopicData as DcpsSubscriptionBuiltinTopicData, TOPIC_NAME_DCPS_PARTICIPANT,
146    TOPIC_NAME_DCPS_PUBLICATION, TOPIC_NAME_DCPS_SUBSCRIPTION, TOPIC_NAME_DCPS_TOPIC,
147    TopicBuiltinTopicData as DcpsTopicBuiltinTopicData,
148};
149pub use dds_type::{
150    DdsType, DdsTypeRow, DecodeError, EncodeError, Extensibility, ExtensibilityKind, RawBytes,
151};
152pub use entity::{Entity, EntityState, StatusCondition, StatusMask, immutable_if_enabled};
153
154pub use coherent_set::{CoherentScope, CoherentSetMarker, GroupAccessScope};
155#[cfg(feature = "std")]
156pub use condition::{Condition, GuardCondition, QueryCondition, ReadCondition, WaitSet};
157pub use error::{DdsError, Result};
158pub use factory::DomainParticipantFactory;
159pub use instance_handle::{HANDLE_NIL, InstanceHandle, InstanceHandleAllocator};
160#[cfg(feature = "std")]
161pub use instance_tracker::{InstanceState, InstanceTracker, KeyHash};
162#[cfg(feature = "std")]
163pub use participant::IgnoreFilter;
164pub use participant::{DomainId, DomainParticipant};
165pub use publisher::{DataWriter, Publisher};
166pub use qos::{
167    DataReaderQos, DataWriterQos, DomainParticipantQos, PublisherQos, SubscriberQos, TopicQos,
168};
169pub use sample::Sample;
170pub use sample_info::{
171    InstanceStateKind, SampleInfo, SampleStateKind, ViewStateKind, instance_state_mask,
172    sample_state_mask, view_state_mask,
173};
174pub use subscriber::{DataReader, Subscriber};
175pub use time::{Duration, Time, get_current_time};
176#[cfg(feature = "std")]
177pub use topic::hash_join_two;
178pub use topic::{
179    ContentFilteredTopic, JoinedRow, MultiTopic, Topic, TopicDescription, TopicDescriptionHandle,
180};
181
182#[cfg(test)]
183#[allow(clippy::expect_used, clippy::unwrap_used)]
184mod tests {
185    use super::*;
186
187    #[test]
188    fn end_to_end_in_process_loopback() {
189        // Loopback smoke: Factory → offline participant → Topic → Writer/Reader.
190        // We loop a sample via __push_raw from the DataWriter into the
191        // DataReader. Live transport is covered in the runtime tests in
192        // crates/dcps/src/runtime.rs.
193        let factory = DomainParticipantFactory::instance();
194        let p = factory.create_participant_offline(0, DomainParticipantQos::default());
195        let topic = p
196            .create_topic::<RawBytes>("Chatter", TopicQos::default())
197            .unwrap();
198
199        let pub_ = p.create_publisher(PublisherQos::default());
200        let w = pub_
201            .create_datawriter::<RawBytes>(&topic, DataWriterQos::default())
202            .unwrap();
203
204        let sub = p.create_subscriber(SubscriberQos::default());
205        let r = sub
206            .create_datareader::<RawBytes>(&topic, DataReaderQos::default())
207            .unwrap();
208
209        w.write(&RawBytes::new(vec![1, 2, 3])).unwrap();
210        w.write(&RawBytes::new(vec![4, 5])).unwrap();
211        // Manually drain the queue and push into the reader, simulating
212        // the live transport path.
213        for bytes in w.__drain_pending() {
214            r.__push_raw(bytes).unwrap();
215        }
216        let samples = r.take().unwrap();
217        assert_eq!(samples.len(), 2);
218        assert_eq!(samples[0].data, vec![1, 2, 3]);
219        assert_eq!(samples[1].data, vec![4, 5]);
220    }
221}