Skip to main content

cyclonedds/
subscriber.rs

1use crate::internal::ffi;
2use crate::internal::traits::AsFfi;
3use crate::{Participant, Result};
4
5/// A `Subscriber` groups [`Readers`](crate::Reader) and controls their shared
6/// [`QoS`](crate::QoS). Readers created under a subscriber inherit its
7/// [`QoS`](crate::QoS) where applicable.
8///
9/// Use [`Subscriber::new`] for simple construction or [`Subscriber::builder`]
10/// for [`QoS`](crate::QoS) and
11/// [`listener`](crate::listener::SubscriberListener) configuration.
12///
13/// In most applications a subscriber is created implicitly when constructing a
14/// [`Reader`](crate::Reader) directly. Use an explicit subscriber when you need
15/// coordinated reads across multiple readers.
16#[derive(Debug)]
17pub struct Subscriber<'domain, 'participant> {
18    pub(crate) inner: cyclonedds_sys::dds_entity_t,
19    phantom: std::marker::PhantomData<&'participant Participant<'domain>>,
20}
21
22/// Builder for [`Subscriber`] (accessible via [`Subscriber::builder`]).
23#[derive(Debug)]
24pub struct SubscriberBuilder<'domain, 'participant, 'qos> {
25    participant: &'participant Participant<'domain>,
26    qos: Option<&'qos crate::QoS>,
27    listener: Option<crate::SubscriberListener>,
28}
29
30impl<'d, 'p, 'q> SubscriberBuilder<'d, 'p, 'q> {
31    /// Creates a new [`SubscriberBuilder`] for the given [`Participant`].
32    ///
33    /// # Examples
34    ///
35    /// ```
36    /// use cyclonedds::builder::SubscriberBuilder;
37    /// use cyclonedds::{Domain, Participant};
38    ///
39    /// let domain = Domain::default();
40    /// let participant = Participant::new(&domain)?;
41    /// let subscriber_builder = SubscriberBuilder::new(&participant);
42    /// # Ok::<_, cyclonedds::Error>(())
43    /// ```
44    #[must_use]
45    pub const fn new(participant: &'p Participant<'d>) -> Self {
46        Self {
47            participant,
48            qos: None,
49            listener: None,
50        }
51    }
52
53    /// Sets the [`QoS`](crate::QoS) for this subscriber builder.
54    ///
55    /// # Examples
56    ///
57    /// ```
58    /// use cyclonedds::builder::SubscriberBuilder;
59    /// use cyclonedds::qos::policy;
60    /// use cyclonedds::{Duration, QoS};
61    /// # use cyclonedds::{Domain, Participant};
62    /// # let domain = Domain::default();
63    /// # let participant = Participant::new(&domain)?;
64    ///
65    /// let qos = QoS::new().with_reliability(policy::Reliability::Reliable {
66    ///     max_blocking_time: Duration::from_millis(100),
67    /// });
68    /// let subscriber_builder = SubscriberBuilder::new(&participant).with_qos(&qos);
69    /// # Ok::<_, cyclonedds::Error>(())
70    /// ```
71    #[must_use]
72    pub const fn with_qos(mut self, qos: &'q crate::QoS) -> Self {
73        self.qos = Some(qos);
74        self
75    }
76
77    ///
78    /// Sets the [`Listener`](crate::Listener) on this subscriber builder.
79    ///
80    /// # Examples
81    ///
82    /// ```
83    /// use cyclonedds::Listener;
84    /// use cyclonedds::builder::SubscriberBuilder;
85    /// # use cyclonedds::{Domain, Participant};
86    /// # let domain = Domain::default();
87    /// # let participant = Participant::new(&domain)?;
88    ///
89    /// let subscriber_builder = SubscriberBuilder::new(&participant).with_listener(Listener::new());
90    /// # Ok::<_, cyclonedds::Error>(())
91    /// ```
92    #[must_use]
93    pub fn with_listener<L>(mut self, listener: L) -> Self
94    where
95        L: AsRef<crate::SubscriberListener>,
96    {
97        self.listener = Some(*listener.as_ref());
98        self
99    }
100
101    /// Builds the [`Subscriber`].
102    ///
103    /// # Errors
104    ///
105    /// Returns an [`Error`](crate::Error) if the subscriber failed to create.
106    ///
107    /// # Examples
108    ///
109    /// ```
110    /// use cyclonedds::QoS;
111    /// use cyclonedds::builder::SubscriberBuilder;
112    /// use cyclonedds::qos::policy;
113    /// # use cyclonedds::{Domain, Participant};
114    /// # let domain = Domain::default();
115    /// # let participant = Participant::new(&domain)?;
116    ///
117    /// let qos = QoS::new().with_durability(policy::Durability::TransientLocal);
118    /// let subscriber = SubscriberBuilder::new(&participant)
119    ///     .with_qos(&qos)
120    ///     .build()?;
121    /// # Ok::<_, cyclonedds::Error>(())
122    /// ```
123    pub fn build(self) -> Result<Subscriber<'d, 'p>> {
124        // NOTE: using `and_then` to avoid ? branch on the listener for coverage
125        // since the C lib currently panics on OOM rather than returning null.
126        self.listener
127            .map(|listener| listener.as_ffi())
128            .transpose()
129            .and_then(|listener| {
130                Ok(Subscriber {
131                    inner: ffi::dds_create_subscriber(
132                        self.participant.inner,
133                        self.qos.map(|qos| &qos.inner),
134                        listener.as_ref(),
135                    )?,
136                    phantom: std::marker::PhantomData,
137                })
138            })
139    }
140}
141
142impl<'d, 'p> Subscriber<'d, 'p> {
143    /// Creates a new `Subscriber` under `participant` with default
144    /// [`QoS`](crate::QoS) and no
145    /// [`listener`](crate::listener::SubscriberListener).
146    ///
147    /// # Errors
148    ///
149    /// Returns an [`Error`](crate::Error) if the subscriber fails to create.
150    ///
151    /// # Examples
152    ///
153    /// ```
154    /// use cyclonedds::Subscriber;
155    /// # use cyclonedds::{Domain, Participant};
156    /// # let domain = Domain::default();
157    /// # let participant = Participant::new(&domain)?;
158    ///
159    /// let subscriber = Subscriber::new(&participant)?;
160    /// Ok::<_, cyclonedds::Error>(())
161    /// ```
162    pub fn new(participant: &'p Participant<'d>) -> Result<Self> {
163        Self::builder(participant).build()
164    }
165
166    /// Returns a [`SubscriberBuilder`](crate::builder::SubscriberBuilder) for
167    /// constructing a subscriber with custom [`QoS`](crate::QoS) or a
168    /// [`listener`](crate::listener::SubscriberListener).
169    ///
170    /// # Examples
171    ///
172    /// ```
173    /// use cyclonedds::qos::policy::{Durability, Presentation};
174    /// use cyclonedds::{QoS, Subscriber};
175    /// # use cyclonedds::{Domain, Participant};
176    /// # let domain = Domain::default();
177    /// # let participant = Participant::new(&domain)?;
178    ///
179    /// let qos = QoS::new().with_presentation(Presentation::Topic {
180    ///     coherent_access: true,
181    ///     ordered_access: true,
182    /// });
183    /// let subscriber = Subscriber::builder(&participant).with_qos(&qos).build()?;
184    /// Ok::<_, cyclonedds::Error>(())
185    /// ```
186    #[must_use]
187    pub const fn builder<'q>(participant: &'p Participant<'d>) -> SubscriberBuilder<'d, 'p, 'q> {
188        SubscriberBuilder::new(participant)
189    }
190
191    /// (WARN: unimplemented in C lib): Notifies all readers belonging to this
192    /// subscriber that data is available.
193    ///
194    /// <div class="warning">
195    ///
196    /// This function is currently not implemented by the underlying C library
197    /// and will thus always return an unsupported error.
198    ///
199    /// </div>
200    ///
201    /// Triggers the
202    /// [`DataOnReaders`](crate::listener::SubscriberListener::with_data_on_readers)
203    /// callback on the subscriber's listener and the
204    /// [`DataAvailable`](crate::listener::ReaderListener::with_data_available)
205    /// callback on each reader's listener.
206    ///
207    /// # Errors
208    ///
209    /// Returns an [`Error`](crate::Error) if the subscriber fails to notify the
210    /// readers.
211    ///
212    /// # Examples
213    ///
214    /// ```no_run
215    /// use cyclonedds::Subscriber;
216    /// # use cyclonedds::{Domain, Participant};
217    /// # let domain = Domain::default();
218    /// # let participant = Participant::new(&domain)?;
219    ///
220    /// let subscriber = Subscriber::new(&participant)?;
221    /// subscriber.notify_readers()?;
222    /// # Ok::<_, cyclonedds::Error>(())
223    /// ```
224    pub fn notify_readers(&self) -> Result<()> {
225        ffi::dds_notify_readers(self.inner)
226    }
227
228    pub(crate) const fn from_existing(
229        inner: cyclonedds_sys::dds_entity_t,
230    ) -> std::mem::ManuallyDrop<Self> {
231        std::mem::ManuallyDrop::new(Self {
232            inner,
233            phantom: std::marker::PhantomData,
234        })
235    }
236
237    /// Sets the [`SubscriberListener`](crate::SubscriberListener) on this
238    /// subscriber, replacing any previously set listener.
239    ///
240    /// # Errors
241    ///
242    /// Returns an [`Error`](crate::Error) if the subscriber fails to set the
243    /// listener.
244    ///
245    /// # Examples
246    ///
247    /// ```
248    /// use cyclonedds::SubscriberListener;
249    /// # use cyclonedds::{Domain, Participant, Subscriber};
250    /// # let domain = Domain::default();
251    /// # let participant = Participant::new(&domain)?;
252    ///
253    /// let mut subscriber = Subscriber::new(&participant)?;
254    /// subscriber.set_listener(SubscriberListener::new())?;
255    /// # Ok::<_, cyclonedds::Error>(())
256    /// ```
257    pub fn set_listener<L>(&mut self, listener: L) -> Result<()>
258    where
259        L: AsRef<crate::SubscriberListener>,
260    {
261        listener
262            .as_ref()
263            .as_ffi()
264            .and_then(|listener| ffi::dds_set_listener(self.inner, Some(listener.inner)))
265    }
266
267    /// Removes the listener from this subscriber.
268    ///
269    /// # Errors
270    ///
271    /// Returns an [`Error`](crate::Error) if the subscriber fails to unset the
272    /// listener.
273    ///
274    /// # Examples
275    ///
276    /// ```
277    /// # use cyclonedds::{Domain, Participant, Subscriber};
278    /// # let domain = Domain::default();
279    /// # let participant = Participant::new(&domain)?;
280    /// let mut subscriber = Subscriber::new(&participant)?;
281    /// subscriber.unset_listener()?;
282    /// # Ok::<_, cyclonedds::Error>(())
283    /// ```
284    pub fn unset_listener(&mut self) -> Result<()> {
285        ffi::dds_set_listener(self.inner, None)?;
286        Ok(())
287    }
288
289    /// Sets the [`SubscriberListener`](crate::SubscriberListener) on this
290    /// subscriber, consuming and returning `self`.
291    ///
292    /// # Errors
293    ///
294    /// Returns an [`Error`](crate::Error) if the subscriber fails to set the
295    /// listener.
296    ///
297    /// # Examples
298    ///
299    /// ```
300    /// use cyclonedds::SubscriberListener;
301    /// # use cyclonedds::{Domain, Participant, Subscriber};
302    /// # let domain = Domain::default();
303    /// # let participant = Participant::new(&domain)?;
304    ///
305    /// let subscriber = Subscriber::new(&participant)?.with_listener(SubscriberListener::new())?;
306    /// # Ok::<_, cyclonedds::Error>(())
307    /// ```
308    pub fn with_listener<L>(mut self, listener: L) -> Result<Self>
309    where
310        L: AsRef<crate::SubscriberListener>,
311    {
312        self.set_listener(listener).map(|_err| self)
313    }
314}
315
316impl Drop for Subscriber<'_, '_> {
317    fn drop(&mut self) {
318        let result = ffi::dds_delete(self.inner);
319        debug_assert!(
320            result.is_ok(),
321            "unable to delete {self:?}: failed with {result:?}"
322        );
323    }
324}
325
326#[cfg(test)]
327mod tests {
328    use super::*;
329
330    #[test]
331    fn test_subscriber_create() {
332        let domain_id = crate::tests::domain::unique_id();
333        let domain = crate::Domain::new(domain_id).unwrap();
334        let qos = crate::QoS::new();
335        let participant = Participant::new(&domain).unwrap();
336        let _ = Subscriber::new(&participant).unwrap();
337        let _ = Subscriber::builder(&participant)
338            .with_qos(&qos)
339            .build()
340            .unwrap();
341    }
342
343    #[test]
344    fn test_subscriber_create_with_invalid_participant() {
345        let domain_id = crate::tests::domain::unique_id();
346        let domain = crate::Domain::new(domain_id).unwrap();
347        let qos = crate::QoS::new();
348        let mut participant = Participant::new(&domain).unwrap();
349        let participant_id = participant.inner;
350        participant.inner = 0;
351        let result = Subscriber::new(&participant).unwrap_err();
352        assert_eq!(result, crate::Error::BadParameter);
353        let result = Subscriber::builder(&participant)
354            .with_qos(&qos)
355            .build()
356            .unwrap_err();
357        assert_eq!(result, crate::Error::BadParameter);
358        participant.inner = participant_id;
359    }
360
361    #[test]
362    fn test_subscriber_from_existing_subscriber() {
363        let domain_id = crate::tests::domain::unique_id();
364        let domain = crate::Domain::new(domain_id).unwrap();
365        let participant = crate::Participant::new(&domain).unwrap();
366        let subscriber = Subscriber::new(&participant).unwrap();
367
368        let new_subscriber = Subscriber::from_existing(subscriber.inner);
369
370        assert_eq!(new_subscriber.inner, subscriber.inner);
371    }
372
373    #[test]
374    fn test_subscriber_notify_readers_not_yet_supported_by_c_lib() {
375        let domain_id = crate::tests::domain::unique_id();
376        let domain = crate::Domain::new(domain_id).unwrap();
377        let participant = crate::Participant::new(&domain).unwrap();
378
379        let subscriber = Subscriber::new(&participant).unwrap();
380
381        let result = subscriber.notify_readers();
382        assert_eq!(
383            result,
384            Err(crate::Error::Unsupported),
385            "result was not unsupported (might be implemented now?)"
386        );
387    }
388
389    #[test]
390    fn test_subscriber_with_listener() {
391        let domain_id = crate::tests::domain::unique_id();
392        let domain = crate::Domain::new(domain_id).unwrap();
393        let participant = crate::Participant::new(&domain).unwrap();
394
395        let listener = crate::SubscriberListener::new().with_data_on_readers(|_| ());
396
397        let _ = Subscriber::new(&participant)
398            .unwrap()
399            .with_listener(listener)
400            .unwrap();
401
402        let _ = Subscriber::builder(&participant)
403            .with_listener(listener)
404            .build()
405            .unwrap();
406
407        let mut subscriber = Subscriber::new(&participant).unwrap();
408        subscriber.set_listener(listener).unwrap();
409        subscriber.unset_listener().unwrap();
410    }
411
412    #[test]
413    fn test_subscriber_with_listener_on_invalid_subscriber() {
414        let domain_id = crate::tests::domain::unique_id();
415        let domain = crate::Domain::new(domain_id).unwrap();
416        let participant = crate::Participant::new(&domain).unwrap();
417
418        let listener = crate::SubscriberListener::new().with_data_on_readers(|_| ());
419
420        let mut subscriber = Subscriber::new(&participant).unwrap();
421        let subscriber_id = subscriber.inner;
422        subscriber.inner = 0;
423        let result = subscriber.set_listener(listener).unwrap_err();
424        assert_eq!(result, crate::Error::BadParameter);
425        let result = subscriber.unset_listener().unwrap_err();
426        assert_eq!(result, crate::Error::BadParameter);
427        subscriber.inner = subscriber_id;
428    }
429}