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}