1extern crate alloc;
54
55use alloc::boxed::Box;
56
57use crate::entity::StatusMask;
58use crate::instance_handle::InstanceHandle;
59use crate::psm_constants::status as status_bits;
60use crate::status::{
61 InconsistentTopicStatus, LivelinessChangedStatus, LivelinessLostStatus,
62 OfferedDeadlineMissedStatus, OfferedIncompatibleQosStatus, PublicationMatchedStatus,
63 RequestedDeadlineMissedStatus, RequestedIncompatibleQosStatus, SampleLostStatus,
64 SampleRejectedStatus, SubscriptionMatchedStatus,
65};
66
67pub trait TopicListener: Send + Sync {
76 fn on_inconsistent_topic(&self, _topic: InstanceHandle, _status: InconsistentTopicStatus) {}
79}
80
81pub trait DataWriterListener: Send + Sync {
90 fn on_offered_deadline_missed(
93 &self,
94 _writer: InstanceHandle,
95 _status: OfferedDeadlineMissedStatus,
96 ) {
97 }
98
99 fn on_offered_incompatible_qos(
102 &self,
103 _writer: InstanceHandle,
104 _status: OfferedIncompatibleQosStatus,
105 ) {
106 }
107
108 fn on_liveliness_lost(&self, _writer: InstanceHandle, _status: LivelinessLostStatus) {}
111
112 fn on_publication_matched(&self, _writer: InstanceHandle, _status: PublicationMatchedStatus) {}
115}
116
117pub trait PublisherListener: Send + Sync {
127 fn on_offered_deadline_missed(
129 &self,
130 _writer: InstanceHandle,
131 _status: OfferedDeadlineMissedStatus,
132 ) {
133 }
134
135 fn on_offered_incompatible_qos(
137 &self,
138 _writer: InstanceHandle,
139 _status: OfferedIncompatibleQosStatus,
140 ) {
141 }
142
143 fn on_liveliness_lost(&self, _writer: InstanceHandle, _status: LivelinessLostStatus) {}
145
146 fn on_publication_matched(&self, _writer: InstanceHandle, _status: PublicationMatchedStatus) {}
148}
149
150pub trait DataReaderListener: Send + Sync {
159 fn on_data_available(&self, _reader: InstanceHandle) {}
161
162 fn on_sample_lost(&self, _reader: InstanceHandle, _status: SampleLostStatus) {}
165
166 fn on_sample_rejected(&self, _reader: InstanceHandle, _status: SampleRejectedStatus) {}
168
169 fn on_requested_deadline_missed(
172 &self,
173 _reader: InstanceHandle,
174 _status: RequestedDeadlineMissedStatus,
175 ) {
176 }
177
178 fn on_requested_incompatible_qos(
181 &self,
182 _reader: InstanceHandle,
183 _status: RequestedIncompatibleQosStatus,
184 ) {
185 }
186
187 fn on_liveliness_changed(&self, _reader: InstanceHandle, _status: LivelinessChangedStatus) {}
190
191 fn on_subscription_matched(&self, _reader: InstanceHandle, _status: SubscriptionMatchedStatus) {
193 }
194}
195
196pub trait SubscriberListener: Send + Sync {
204 fn on_data_on_readers(&self, _subscriber: InstanceHandle) {}
207
208 fn on_data_available(&self, _reader: InstanceHandle) {}
210
211 fn on_sample_lost(&self, _reader: InstanceHandle, _status: SampleLostStatus) {}
213
214 fn on_sample_rejected(&self, _reader: InstanceHandle, _status: SampleRejectedStatus) {}
216
217 fn on_requested_deadline_missed(
219 &self,
220 _reader: InstanceHandle,
221 _status: RequestedDeadlineMissedStatus,
222 ) {
223 }
224
225 fn on_requested_incompatible_qos(
227 &self,
228 _reader: InstanceHandle,
229 _status: RequestedIncompatibleQosStatus,
230 ) {
231 }
232
233 fn on_liveliness_changed(&self, _reader: InstanceHandle, _status: LivelinessChangedStatus) {}
235
236 fn on_subscription_matched(&self, _reader: InstanceHandle, _status: SubscriptionMatchedStatus) {
238 }
239}
240
241pub trait DomainParticipantListener: Send + Sync {
259 fn on_inconsistent_topic(&self, _topic: InstanceHandle, _status: InconsistentTopicStatus) {}
263
264 fn on_offered_deadline_missed(
268 &self,
269 _writer: InstanceHandle,
270 _status: OfferedDeadlineMissedStatus,
271 ) {
272 }
273
274 fn on_offered_incompatible_qos(
276 &self,
277 _writer: InstanceHandle,
278 _status: OfferedIncompatibleQosStatus,
279 ) {
280 }
281
282 fn on_liveliness_lost(&self, _writer: InstanceHandle, _status: LivelinessLostStatus) {}
284
285 fn on_publication_matched(&self, _writer: InstanceHandle, _status: PublicationMatchedStatus) {}
287
288 fn on_data_on_readers(&self, _subscriber: InstanceHandle) {}
292
293 fn on_data_available(&self, _reader: InstanceHandle) {}
295
296 fn on_sample_lost(&self, _reader: InstanceHandle, _status: SampleLostStatus) {}
298
299 fn on_sample_rejected(&self, _reader: InstanceHandle, _status: SampleRejectedStatus) {}
301
302 fn on_requested_deadline_missed(
304 &self,
305 _reader: InstanceHandle,
306 _status: RequestedDeadlineMissedStatus,
307 ) {
308 }
309
310 fn on_requested_incompatible_qos(
312 &self,
313 _reader: InstanceHandle,
314 _status: RequestedIncompatibleQosStatus,
315 ) {
316 }
317
318 fn on_liveliness_changed(&self, _reader: InstanceHandle, _status: LivelinessChangedStatus) {}
320
321 fn on_subscription_matched(&self, _reader: InstanceHandle, _status: SubscriptionMatchedStatus) {
323 }
324}
325
326pub type BoxedTopicListener = Box<dyn TopicListener>;
333pub type BoxedDataWriterListener = Box<dyn DataWriterListener>;
335pub type BoxedPublisherListener = Box<dyn PublisherListener>;
337pub type BoxedDataReaderListener = Box<dyn DataReaderListener>;
339pub type BoxedSubscriberListener = Box<dyn SubscriberListener>;
341pub type BoxedDomainParticipantListener = Box<dyn DomainParticipantListener>;
343
344pub type ArcTopicListener = alloc::sync::Arc<dyn TopicListener>;
349pub type ArcDataWriterListener = alloc::sync::Arc<dyn DataWriterListener>;
351pub type ArcPublisherListener = alloc::sync::Arc<dyn PublisherListener>;
353pub type ArcDataReaderListener = alloc::sync::Arc<dyn DataReaderListener>;
355pub type ArcSubscriberListener = alloc::sync::Arc<dyn SubscriberListener>;
357pub type ArcDomainParticipantListener = alloc::sync::Arc<dyn DomainParticipantListener>;
359
360#[inline]
368#[must_use]
369pub fn listener_handles(listener_present: bool, mask: StatusMask, status_bit: u32) -> bool {
370 listener_present && (mask & status_bit) != 0
371}
372
373#[must_use]
377pub fn status_bit_for_inconsistent_topic() -> u32 {
378 status_bits::INCONSISTENT_TOPIC
379}
380
381#[cfg(test)]
382#[allow(clippy::expect_used, clippy::unwrap_used)]
383mod tests {
384 use super::*;
385 use core::sync::atomic::{AtomicU32, Ordering};
386
387 #[test]
390 fn topic_listener_is_object_safe() {
391 struct Counter(AtomicU32);
392 impl TopicListener for Counter {
393 fn on_inconsistent_topic(
394 &self,
395 _topic: InstanceHandle,
396 _status: InconsistentTopicStatus,
397 ) {
398 self.0.fetch_add(1, Ordering::Relaxed);
399 }
400 }
401 let _: BoxedTopicListener = Box::new(Counter(AtomicU32::new(0)));
402 }
403
404 #[test]
405 fn datawriter_listener_is_object_safe() {
406 struct L;
407 impl DataWriterListener for L {}
408 let _: BoxedDataWriterListener = Box::new(L);
409 }
410
411 #[test]
412 fn publisher_listener_is_object_safe() {
413 struct L;
414 impl PublisherListener for L {}
415 let _: BoxedPublisherListener = Box::new(L);
416 }
417
418 #[test]
419 fn datareader_listener_is_object_safe() {
420 struct L;
421 impl DataReaderListener for L {}
422 let _: BoxedDataReaderListener = Box::new(L);
423 }
424
425 #[test]
426 fn subscriber_listener_is_object_safe() {
427 struct L;
428 impl SubscriberListener for L {}
429 let _: BoxedSubscriberListener = Box::new(L);
430 }
431
432 #[test]
433 fn participant_listener_is_object_safe() {
434 struct L;
435 impl DomainParticipantListener for L {}
436 let _: BoxedDomainParticipantListener = Box::new(L);
437 }
438
439 #[test]
442 fn default_callbacks_do_not_panic() {
443 struct Noop;
445 impl TopicListener for Noop {}
446 impl DataWriterListener for Noop {}
447 impl PublisherListener for Noop {}
448 impl DataReaderListener for Noop {}
449 impl SubscriberListener for Noop {}
450 impl DomainParticipantListener for Noop {}
451 let _: BoxedDomainParticipantListener = Box::new(Noop);
454 }
455
456 #[test]
457 fn listener_handles_respects_mask_and_presence() {
458 let bit = status_bit_for_inconsistent_topic();
459 assert!(listener_handles(true, bit, bit));
460 assert!(!listener_handles(false, bit, bit));
461 assert!(!listener_handles(true, 0, bit));
462 assert!(!listener_handles(true, status_bits::SAMPLE_LOST, bit));
464 }
465
466 #[test]
467 fn status_bit_for_inconsistent_topic_matches_psm() {
468 assert_eq!(
469 status_bit_for_inconsistent_topic(),
470 status_bits::INCONSISTENT_TOPIC
471 );
472 }
473
474 #[test]
475 fn all_listener_traits_default_methods_invoke_safely() {
476 struct Noop;
480 impl TopicListener for Noop {}
481 impl DataWriterListener for Noop {}
482 impl PublisherListener for Noop {}
483 impl DataReaderListener for Noop {}
484 impl SubscriberListener for Noop {}
485 impl DomainParticipantListener for Noop {}
486
487 let h = InstanceHandle::from_raw(1);
488 let n = Noop;
489 TopicListener::on_inconsistent_topic(&n, h, InconsistentTopicStatus::default());
490
491 DataWriterListener::on_offered_deadline_missed(
492 &n,
493 h,
494 OfferedDeadlineMissedStatus::default(),
495 );
496 DataWriterListener::on_offered_incompatible_qos(
497 &n,
498 h,
499 OfferedIncompatibleQosStatus::default(),
500 );
501 DataWriterListener::on_liveliness_lost(&n, h, LivelinessLostStatus::default());
502 DataWriterListener::on_publication_matched(&n, h, PublicationMatchedStatus::default());
503
504 PublisherListener::on_offered_deadline_missed(
505 &n,
506 h,
507 OfferedDeadlineMissedStatus::default(),
508 );
509 PublisherListener::on_offered_incompatible_qos(
510 &n,
511 h,
512 OfferedIncompatibleQosStatus::default(),
513 );
514 PublisherListener::on_liveliness_lost(&n, h, LivelinessLostStatus::default());
515 PublisherListener::on_publication_matched(&n, h, PublicationMatchedStatus::default());
516
517 DataReaderListener::on_data_available(&n, h);
518 DataReaderListener::on_sample_lost(&n, h, SampleLostStatus::default());
519 DataReaderListener::on_sample_rejected(&n, h, SampleRejectedStatus::default());
520 DataReaderListener::on_requested_deadline_missed(
521 &n,
522 h,
523 RequestedDeadlineMissedStatus::default(),
524 );
525 DataReaderListener::on_requested_incompatible_qos(
526 &n,
527 h,
528 RequestedIncompatibleQosStatus::default(),
529 );
530 DataReaderListener::on_liveliness_changed(&n, h, LivelinessChangedStatus::default());
531 DataReaderListener::on_subscription_matched(&n, h, SubscriptionMatchedStatus::default());
532
533 SubscriberListener::on_data_on_readers(&n, h);
534 SubscriberListener::on_data_available(&n, h);
535 SubscriberListener::on_sample_lost(&n, h, SampleLostStatus::default());
536 SubscriberListener::on_sample_rejected(&n, h, SampleRejectedStatus::default());
537 SubscriberListener::on_requested_deadline_missed(
538 &n,
539 h,
540 RequestedDeadlineMissedStatus::default(),
541 );
542 SubscriberListener::on_requested_incompatible_qos(
543 &n,
544 h,
545 RequestedIncompatibleQosStatus::default(),
546 );
547 SubscriberListener::on_liveliness_changed(&n, h, LivelinessChangedStatus::default());
548 SubscriberListener::on_subscription_matched(&n, h, SubscriptionMatchedStatus::default());
549
550 DomainParticipantListener::on_inconsistent_topic(&n, h, InconsistentTopicStatus::default());
551 DomainParticipantListener::on_offered_deadline_missed(
552 &n,
553 h,
554 OfferedDeadlineMissedStatus::default(),
555 );
556 DomainParticipantListener::on_offered_incompatible_qos(
557 &n,
558 h,
559 OfferedIncompatibleQosStatus::default(),
560 );
561 DomainParticipantListener::on_liveliness_lost(&n, h, LivelinessLostStatus::default());
562 DomainParticipantListener::on_publication_matched(
563 &n,
564 h,
565 PublicationMatchedStatus::default(),
566 );
567 DomainParticipantListener::on_data_on_readers(&n, h);
568 DomainParticipantListener::on_data_available(&n, h);
569 DomainParticipantListener::on_sample_lost(&n, h, SampleLostStatus::default());
570 DomainParticipantListener::on_sample_rejected(&n, h, SampleRejectedStatus::default());
571 DomainParticipantListener::on_requested_deadline_missed(
572 &n,
573 h,
574 RequestedDeadlineMissedStatus::default(),
575 );
576 DomainParticipantListener::on_requested_incompatible_qos(
577 &n,
578 h,
579 RequestedIncompatibleQosStatus::default(),
580 );
581 DomainParticipantListener::on_liveliness_changed(&n, h, LivelinessChangedStatus::default());
582 DomainParticipantListener::on_subscription_matched(
583 &n,
584 h,
585 SubscriptionMatchedStatus::default(),
586 );
587 }
588
589 #[test]
590 fn datareader_listener_call_runs_default_methods() {
591 struct Counters {
592 avail: AtomicU32,
593 lost: AtomicU32,
594 }
595 impl DataReaderListener for Counters {
596 fn on_data_available(&self, _r: InstanceHandle) {
597 self.avail.fetch_add(1, Ordering::Relaxed);
598 }
599 fn on_sample_lost(&self, _r: InstanceHandle, _s: SampleLostStatus) {
600 self.lost.fetch_add(1, Ordering::Relaxed);
601 }
602 }
603 let c = Counters {
604 avail: AtomicU32::new(0),
605 lost: AtomicU32::new(0),
606 };
607 let h = InstanceHandle::from_raw(1);
608 c.on_data_available(h);
609 c.on_data_available(h);
610 c.on_sample_lost(h, SampleLostStatus::default());
611 c.on_subscription_matched(h, SubscriptionMatchedStatus::default());
613 assert_eq!(c.avail.load(Ordering::Relaxed), 2);
614 assert_eq!(c.lost.load(Ordering::Relaxed), 1);
615 }
616}