1use core::time::Duration;
35
36extern crate alloc;
37use alloc::vec::Vec;
38
39use alloc::rc::Rc;
40
41use crate::error::WireError;
42use crate::fragment_assembler::{AssemblerCaps, FragmentAssembler};
43use crate::header::RtpsHeader;
44use crate::history_cache::{CacheChange, ChangeKind, HistoryCache};
45use crate::message_builder::OutboundDatagram;
46use crate::submessage_header::{FLAG_E_LITTLE_ENDIAN, SubmessageHeader, SubmessageId};
47use crate::submessages::{
48 AckNackSubmessage, DataFragSubmessage, DataSubmessage, GapSubmessage, HeartbeatSubmessage,
49 NackFragSubmessage, SequenceNumberSet,
50};
51use crate::wire_types::{Guid, GuidPrefix, SequenceNumber, VendorId};
52use crate::writer_proxy::WriterProxy;
53
54pub const DEFAULT_HEARTBEAT_RESPONSE_DELAY: Duration = Duration::from_millis(0);
70
71#[derive(Debug, Clone)]
78pub struct WriterProxyState {
79 pub proxy: WriterProxy,
81 pub received_cache: HistoryCache,
83 pub delivered_up_to: SequenceNumber,
85 pub assembler: FragmentAssembler,
87 pub pending_acknack_since: Option<Duration>,
90}
91
92impl WriterProxyState {
93 fn new(proxy: WriterProxy, max_samples: usize, caps: AssemblerCaps) -> Self {
94 Self {
95 proxy,
96 received_cache: HistoryCache::new(max_samples),
97 delivered_up_to: SequenceNumber(0),
98 assembler: FragmentAssembler::new(caps),
99 pending_acknack_since: None,
100 }
101 }
102}
103
104#[derive(Debug, Clone)]
106pub struct ReliableReader {
107 guid: Guid,
108 vendor_id: VendorId,
109 writer_proxies: Vec<WriterProxyState>,
110 heartbeat_response_delay: Duration,
111 acknack_count: i32,
112 nackfrag_count: i32,
113 duplicate_frag_count: u64,
114 max_samples_per_proxy: usize,
116 assembler_caps: AssemblerCaps,
117 unknown_src_count: u64,
119 best_effort: bool,
128}
129
130#[derive(Debug, Clone)]
132pub struct ReliableReaderConfig {
133 pub guid: Guid,
135 pub vendor_id: VendorId,
137 pub writer_proxies: Vec<WriterProxy>,
139 pub max_samples_per_proxy: usize,
141 pub heartbeat_response_delay: Duration,
143 pub assembler_caps: AssemblerCaps,
145}
146
147#[derive(Debug, Clone, PartialEq, Eq)]
149pub struct DeliveredSample {
150 pub writer_guid: Guid,
153 pub sequence_number: SequenceNumber,
155 pub payload: alloc::sync::Arc<[u8]>,
158 pub kind: ChangeKind,
164 pub key_hash: Option<[u8; 16]>,
171 pub source_timestamp: Option<crate::header_extension::HeTimestamp>,
176}
177
178impl ReliableReader {
179 #[must_use]
184 pub fn new(cfg: ReliableReaderConfig) -> Self {
185 assert!(
186 cfg.assembler_caps.max_pending_sns > 0,
187 "assembler_caps.max_pending_sns must be > 0; use a Best-Effort reader \
188 or increase the cap to actually accept fragmented samples"
189 );
190 let proxies = cfg
191 .writer_proxies
192 .into_iter()
193 .map(|p| WriterProxyState::new(p, cfg.max_samples_per_proxy, cfg.assembler_caps))
194 .collect();
195 Self {
196 guid: cfg.guid,
197 vendor_id: cfg.vendor_id,
198 writer_proxies: proxies,
199 heartbeat_response_delay: cfg.heartbeat_response_delay,
200 acknack_count: 0,
201 nackfrag_count: 0,
202 duplicate_frag_count: 0,
203 max_samples_per_proxy: cfg.max_samples_per_proxy,
204 assembler_caps: cfg.assembler_caps,
205 unknown_src_count: 0,
206 best_effort: false,
207 }
208 }
209
210 pub fn set_best_effort(&mut self, best_effort: bool) {
214 self.best_effort = best_effort;
215 }
216
217 #[must_use]
219 pub fn guid(&self) -> Guid {
220 self.guid
221 }
222
223 #[must_use]
225 pub fn writer_proxies(&self) -> &[WriterProxyState] {
226 &self.writer_proxies
227 }
228
229 #[must_use]
231 pub fn writer_proxy_count(&self) -> usize {
232 self.writer_proxies.len()
233 }
234
235 #[must_use]
237 pub fn acknack_count(&self) -> i32 {
238 self.acknack_count
239 }
240
241 #[must_use]
243 pub fn nackfrag_count(&self) -> i32 {
244 self.nackfrag_count
245 }
246
247 #[must_use]
250 pub fn pending_fragment_count(&self) -> usize {
251 self.writer_proxies.iter().map(|s| s.assembler.len()).sum()
252 }
253
254 #[must_use]
257 pub fn dropped_fragment_count(&self) -> u64 {
258 self.writer_proxies
259 .iter()
260 .map(|s| s.assembler.drop_count())
261 .sum()
262 }
263
264 #[must_use]
267 pub fn duplicate_fragment_count(&self) -> u64 {
268 self.duplicate_frag_count
269 }
270
271 #[must_use]
274 pub fn unknown_src_count(&self) -> u64 {
275 self.unknown_src_count
276 }
277
278 pub fn add_writer_proxy(&mut self, proxy: WriterProxy) {
293 let guid = proxy.remote_writer_guid;
294 if let Some(idx) = self
295 .writer_proxies
296 .iter()
297 .position(|s| s.proxy.remote_writer_guid == guid)
298 {
299 self.writer_proxies[idx]
302 .proxy
303 .refresh_locators(proxy.unicast_locators, proxy.multicast_locators);
304 self.writer_proxies[idx]
305 .pending_acknack_since
306 .get_or_insert(Duration::ZERO);
307 } else {
308 let mut state =
309 WriterProxyState::new(proxy, self.max_samples_per_proxy, self.assembler_caps);
310 state.pending_acknack_since = Some(Duration::ZERO);
313 self.writer_proxies.push(state);
314 }
315 }
316
317 pub fn remove_writer_proxy(&mut self, guid: Guid) -> Option<WriterProxy> {
319 let idx = self
320 .writer_proxies
321 .iter()
322 .position(|s| s.proxy.remote_writer_guid == guid)?;
323 Some(self.writer_proxies.remove(idx).proxy)
324 }
325
326 pub fn reset_diagnostics(&mut self) {
328 self.acknack_count = 0;
329 self.nackfrag_count = 0;
330 self.duplicate_frag_count = 0;
331 self.unknown_src_count = 0;
332 for s in &mut self.writer_proxies {
333 s.assembler.reset_diagnostics();
334 }
335 }
336
337 pub fn handle_data(
347 &mut self,
348 source_prefix: GuidPrefix,
349 data: &DataSubmessage,
350 source_timestamp: Option<crate::header_extension::HeTimestamp>,
351 ) -> Vec<DeliveredSample> {
352 let Some(idx) = self.proxy_index_by_writer(Guid::new(source_prefix, data.writer_id)) else {
353 self.unknown_src_count = self.unknown_src_count.saturating_add(1);
354 return Vec::new();
355 };
356 let state = &mut self.writer_proxies[idx];
357 let sn = data.writer_sn;
358 if state.proxy.is_known(sn) || sn <= state.delivered_up_to {
359 return Vec::new();
360 }
361 state.proxy.received_change_set(sn);
362 let kind = Self::classify_change_kind(data);
363 if data.key_flag && kind == ChangeKind::Alive {
373 state.proxy.irrelevant_change_set(sn);
378 if sn.0 == state.delivered_up_to.0 + 1 {
379 state.delivered_up_to = sn;
380 }
381 return Self::collect_in_order_for(state, self.best_effort);
382 }
383 let key_hash = data
384 .inline_qos
385 .as_ref()
386 .and_then(crate::inline_qos::find_key_hash);
387 let _ = state.received_cache.insert(CacheChange {
391 sequence_number: sn,
392 payload: alloc::sync::Arc::clone(&data.serialized_payload),
393 kind,
394 key_hash,
395 source_timestamp,
397 });
398 Self::collect_in_order_for(state, self.best_effort)
399 }
400
401 fn classify_change_kind(data: &DataSubmessage) -> ChangeKind {
405 if !data.key_flag {
406 return ChangeKind::Alive;
407 }
408 let Some(pl) = data.inline_qos.as_ref() else {
409 return ChangeKind::Alive;
410 };
411 let Some(bits) = crate::inline_qos::find_status_info(pl) else {
412 return ChangeKind::Alive;
413 };
414 let disposed = bits & crate::inline_qos::status_info::DISPOSED != 0;
415 let unregistered = bits & crate::inline_qos::status_info::UNREGISTERED != 0;
416 match (disposed, unregistered) {
417 (true, true) => ChangeKind::NotAliveDisposedUnregistered,
418 (true, false) => ChangeKind::NotAliveDisposed,
419 (false, true) => ChangeKind::NotAliveUnregistered,
420 (false, false) => ChangeKind::Alive,
421 }
422 }
423
424 pub fn handle_data_frag(
427 &mut self,
428 source_prefix: GuidPrefix,
429 df: &DataFragSubmessage,
430 now: Duration,
431 source_timestamp: Option<crate::header_extension::HeTimestamp>,
432 ) -> Vec<DeliveredSample> {
433 let Some(idx) = self.proxy_index_by_writer(Guid::new(source_prefix, df.writer_id)) else {
434 self.unknown_src_count = self.unknown_src_count.saturating_add(1);
435 return Vec::new();
436 };
437 let state = &mut self.writer_proxies[idx];
438 let sn = df.writer_sn;
439 if state.proxy.is_known(sn) || sn <= state.delivered_up_to {
440 self.duplicate_frag_count = self.duplicate_frag_count.saturating_add(1);
441 return Vec::new();
442 }
443 let result = if let Some(completed) = state.assembler.insert(df) {
444 state.proxy.received_change_set(sn);
445 let _ = state.received_cache.insert(
446 CacheChange::alive(sn, completed.payload).with_source_timestamp(source_timestamp),
447 );
448 Self::collect_in_order_for(state, self.best_effort)
449 } else {
450 Vec::new()
451 };
452 if state.assembler.has_gaps() {
453 state.pending_acknack_since.get_or_insert(now);
454 }
455 result
456 }
457
458 pub fn handle_heartbeat(
460 &mut self,
461 source_prefix: GuidPrefix,
462 hb: &HeartbeatSubmessage,
463 now: Duration,
464 ) -> Vec<DeliveredSample> {
465 let Some(idx) = self.proxy_index_by_writer(Guid::new(source_prefix, hb.writer_id)) else {
466 self.unknown_src_count = self.unknown_src_count.saturating_add(1);
467 return Vec::new();
468 };
469 let state = &mut self.writer_proxies[idx];
470 if hb.liveliness_flag {
471 return Vec::new();
472 }
473 state.proxy.update_from_heartbeat(hb.first_sn, hb.last_sn);
474 let has_missing = state.proxy.has_missing_changes();
475 let has_frag_gaps = state.assembler.has_gaps();
476 if !hb.final_flag || has_missing || has_frag_gaps {
477 state.pending_acknack_since.get_or_insert(now);
478 }
479 Self::collect_in_order_for(state, self.best_effort)
485 }
486
487 pub fn handle_gap(
489 &mut self,
490 source_prefix: GuidPrefix,
491 gap: &GapSubmessage,
492 ) -> Vec<DeliveredSample> {
493 let Some(idx) = self.proxy_index_by_writer(Guid::new(source_prefix, gap.writer_id)) else {
494 self.unknown_src_count = self.unknown_src_count.saturating_add(1);
495 return Vec::new();
496 };
497 let state = &mut self.writer_proxies[idx];
498 let mut sn = gap.gap_start;
499 while sn < gap.gap_list.bitmap_base {
500 state.proxy.irrelevant_change_set(sn);
501 state.assembler.discard(sn);
502 sn = SequenceNumber(sn.0 + 1);
503 }
504 for sn in gap.gap_list.iter_set() {
505 state.proxy.irrelevant_change_set(sn);
506 state.assembler.discard(sn);
507 }
508 Self::collect_in_order_for(state, self.best_effort)
509 }
510
511 pub fn tick(&mut self, now: Duration) -> Result<Vec<Vec<u8>>, WireError> {
518 Ok(self
519 .tick_outbound(now)?
520 .into_iter()
521 .map(|d| d.bytes)
522 .collect())
523 }
524
525 pub fn tick_outbound(&mut self, now: Duration) -> Result<Vec<OutboundDatagram>, WireError> {
532 let mut out = Vec::new();
533 for idx in 0..self.writer_proxies.len() {
534 let Some(since) = self.writer_proxies[idx].pending_acknack_since else {
535 continue;
536 };
537 if now.saturating_sub(since) < self.heartbeat_response_delay {
538 continue;
539 }
540 self.writer_proxies[idx].pending_acknack_since = None;
541 let targets = Rc::new(self.writer_proxies[idx].proxy.unicast_locators.clone());
542
543 let incomplete_sns: Vec<SequenceNumber> = self.writer_proxies[idx]
544 .assembler
545 .incomplete_sns()
546 .collect();
547 for sn in incomplete_sns {
548 let bytes = self.build_nackfrag_datagram(idx, sn)?;
549 out.push(OutboundDatagram {
550 bytes,
551 targets: Rc::clone(&targets),
552 });
553 }
554 let bytes = self.build_acknack_datagram(idx)?;
555 out.push(OutboundDatagram { bytes, targets });
556 }
557 Ok(out)
558 }
559
560 fn proxy_index_by_writer(&self, guid: Guid) -> Option<usize> {
570 self.writer_proxies
571 .iter()
572 .position(|s| s.proxy.remote_writer_guid == guid)
573 }
574
575 fn collect_in_order_for(
576 state: &mut WriterProxyState,
577 best_effort: bool,
578 ) -> Vec<DeliveredSample> {
579 let mut out = Vec::with_capacity(2);
583 loop {
584 let next = SequenceNumber(state.delivered_up_to.0 + 1);
585 if let Some(change) = state.received_cache.get(next) {
586 out.push(DeliveredSample {
587 writer_guid: state.proxy.remote_writer_guid,
588 sequence_number: change.sequence_number,
589 payload: change.payload.clone(),
590 kind: change.kind,
591 key_hash: change.key_hash,
592 source_timestamp: change.source_timestamp,
593 });
594 state.delivered_up_to = next;
595 state.received_cache.remove_up_to(next);
596 } else if state.proxy.is_known(next) && state.proxy.last_available_sn() >= next {
597 state.delivered_up_to = next;
598 } else if next < state.proxy.first_available_sn() {
599 state.delivered_up_to = next;
605 } else if best_effort {
606 match state.received_cache.min_sn() {
615 Some(low) if low.0 > next.0 => {
616 state.delivered_up_to = SequenceNumber(low.0 - 1);
617 }
618 _ => break,
619 }
620 } else {
621 break;
622 }
623 }
624 out
625 }
626
627 fn build_nackfrag_datagram(
628 &mut self,
629 proxy_idx: usize,
630 sn: SequenceNumber,
631 ) -> Result<Vec<u8>, WireError> {
632 let missing = self.writer_proxies[proxy_idx]
633 .assembler
634 .missing_fragments(sn);
635 self.nackfrag_count = self.nackfrag_count.wrapping_add(1);
636 let writer_guid = self.writer_proxies[proxy_idx].proxy.remote_writer_guid;
637 let nf = NackFragSubmessage {
638 reader_id: self.guid.entity_id,
639 writer_id: writer_guid.entity_id,
640 writer_sn: sn,
641 fragment_number_state: missing,
642 count: self.nackfrag_count,
643 };
644 let (body, mut flags) = nf.write_body(true);
645 flags |= FLAG_E_LITTLE_ENDIAN;
646 self.wrap_to_writer(writer_guid.prefix, SubmessageId::NackFrag, flags, &body)
647 }
648
649 fn build_acknack_datagram(&mut self, proxy_idx: usize) -> Result<Vec<u8>, WireError> {
650 let state = &self.writer_proxies[proxy_idx];
651 let base = state.proxy.acknack_base();
652 let missing = state.proxy.missing_changes(256);
653 let snset = SequenceNumberSet::from_missing(base, &missing);
654 self.acknack_count = self.acknack_count.wrapping_add(1);
655 let final_flag = missing.is_empty() && state.proxy.last_available_sn().0 >= 1;
662 let writer_guid = state.proxy.remote_writer_guid;
663 let ack = AckNackSubmessage {
664 reader_id: self.guid.entity_id,
665 writer_id: writer_guid.entity_id,
666 reader_sn_state: snset,
667 count: self.acknack_count,
668 final_flag,
669 };
670 let (body, mut flags) = ack.write_body(true);
671 flags |= FLAG_E_LITTLE_ENDIAN;
672 self.wrap_to_writer(writer_guid.prefix, SubmessageId::AckNack, flags, &body)
673 }
674
675 fn wrap_to_writer(
680 &self,
681 writer_prefix: crate::wire_types::GuidPrefix,
682 id: SubmessageId,
683 flags: u8,
684 body: &[u8],
685 ) -> Result<Vec<u8>, WireError> {
686 let header = RtpsHeader::new(self.vendor_id, self.guid.prefix);
687 let mut out = Vec::new();
688 out.extend_from_slice(&header.to_bytes());
689
690 let info_dst_header = SubmessageHeader {
692 submessage_id: SubmessageId::InfoDst,
693 flags: FLAG_E_LITTLE_ENDIAN,
694 octets_to_next_header: 12,
695 };
696 out.extend_from_slice(&info_dst_header.to_bytes());
697 out.extend_from_slice(&writer_prefix.to_bytes());
698
699 let body_len = u16::try_from(body.len()).map_err(|_| WireError::ValueOutOfRange {
701 message: "submessage body exceeds u16::MAX",
702 })?;
703 let sh = SubmessageHeader {
704 submessage_id: id,
705 flags,
706 octets_to_next_header: body_len,
707 };
708 out.extend_from_slice(&sh.to_bytes());
709 out.extend_from_slice(body);
710 Ok(out)
711 }
712}
713
714#[cfg(test)]
715#[allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
716mod tests {
717 use super::*;
718 use crate::datagram::{ParsedSubmessage, decode_datagram};
719 use crate::wire_types::{EntityId, GuidPrefix, Locator};
720
721 fn single_writer_guid() -> Guid {
722 Guid::new(
723 GuidPrefix::from_bytes([1; 12]),
724 EntityId::user_writer_with_key([0x10, 0x20, 0x30]),
725 )
726 }
727
728 fn make_reader(max_samples: usize) -> ReliableReader {
729 let reader_guid = Guid::new(
730 GuidPrefix::from_bytes([2; 12]),
731 EntityId::user_reader_with_key([0xA0, 0xB0, 0xC0]),
732 );
733 let writer_proxy = WriterProxy::new(
734 single_writer_guid(),
735 alloc::vec![Locator::udp_v4([127, 0, 0, 1], 7420)],
736 alloc::vec![],
737 true,
738 );
739 ReliableReader::new(ReliableReaderConfig {
740 guid: reader_guid,
741 vendor_id: VendorId::ZERODDS,
742 writer_proxies: alloc::vec![writer_proxy],
743 max_samples_per_proxy: max_samples,
744 heartbeat_response_delay: Duration::from_millis(200),
745 assembler_caps: AssemblerCaps::default(),
746 })
747 }
748
749 fn sn(n: i64) -> SequenceNumber {
750 SequenceNumber(n)
751 }
752
753 fn p1() -> GuidPrefix {
755 single_writer_guid().prefix
756 }
757
758 fn p2() -> GuidPrefix {
760 second_writer_guid().prefix
761 }
762
763 fn data(wid: EntityId, rid: EntityId, n: i64, byte: u8) -> DataSubmessage {
764 DataSubmessage {
765 extra_flags: 0,
766 reader_id: rid,
767 writer_id: wid,
768 writer_sn: sn(n),
769 inline_qos: None,
770 key_flag: false,
771 non_standard_flag: false,
772 serialized_payload: alloc::sync::Arc::from(alloc::vec![byte]),
773 }
774 }
775
776 fn heartbeat(
777 wid: EntityId,
778 rid: EntityId,
779 first: i64,
780 last: i64,
781 count: i32,
782 final_flag: bool,
783 ) -> HeartbeatSubmessage {
784 HeartbeatSubmessage {
785 reader_id: rid,
786 writer_id: wid,
787 first_sn: sn(first),
788 last_sn: sn(last),
789 count,
790 final_flag,
791 liveliness_flag: false,
792 group_info: None,
793 }
794 }
795
796 fn first_state(r: &ReliableReader) -> &WriterProxyState {
797 &r.writer_proxies()[0]
798 }
799
800 #[test]
801 fn re_adding_known_writer_preserves_reliability_state() {
802 let mut r = make_reader(10);
808 let w_eid = single_writer_guid().entity_id;
809 let r_eid = r.guid().entity_id;
810 r.handle_heartbeat(
812 p1(),
813 &heartbeat(w_eid, r_eid, 1, 1, 1, false),
814 Duration::ZERO,
815 );
816 assert!(first_state(&r).proxy.has_missing_changes());
817 assert_eq!(first_state(&r).proxy.last_available_sn(), sn(1));
818 r.add_writer_proxy(WriterProxy::new(
820 single_writer_guid(),
821 alloc::vec![Locator::udp_v4([127, 0, 0, 1], 7420)],
822 alloc::vec![],
823 true,
824 ));
825 assert!(
827 first_state(&r).proxy.has_missing_changes(),
828 "re-add must not reset last_available/missing"
829 );
830 assert_eq!(first_state(&r).proxy.last_available_sn(), sn(1));
831 }
832
833 #[test]
834 fn in_order_data_delivered_immediately() {
835 let mut r = make_reader(10);
836 let w_eid = single_writer_guid().entity_id;
837 let r_eid = r.guid().entity_id;
838 let delivered = r.handle_data(p1(), &data(w_eid, r_eid, 1, 0xAA), None);
839 assert_eq!(delivered.len(), 1);
840 assert_eq!(delivered[0].payload.as_ref(), &[0xAA][..]);
841 assert_eq!(delivered[0].writer_guid, single_writer_guid());
842 assert_eq!(first_state(&r).delivered_up_to, sn(1));
843 }
844
845 #[test]
846 fn out_of_order_data_buffered_until_gap_filled() {
847 let mut r = make_reader(10);
848 let w = single_writer_guid().entity_id;
849 let rd = r.guid().entity_id;
850 assert!(r.handle_data(p1(), &data(w, rd, 2, 0x22), None).is_empty());
851 assert!(r.handle_data(p1(), &data(w, rd, 3, 0x33), None).is_empty());
852 let out = r.handle_data(p1(), &data(w, rd, 1, 0x11), None);
853 assert_eq!(
854 out.iter().map(|s| s.sequence_number).collect::<Vec<_>>(),
855 alloc::vec![sn(1), sn(2), sn(3)]
856 );
857 assert_eq!(first_state(&r).delivered_up_to, sn(3));
858 }
859
860 #[test]
865 fn best_effort_reader_skips_leading_gap() {
866 let mut r = make_reader(10);
867 r.set_best_effort(true);
868 let w = single_writer_guid().entity_id;
869 let rd = r.guid().entity_id;
870 let out = r.handle_data(p1(), &data(w, rd, 2, 0x22), None);
874 assert_eq!(
875 out.iter().map(|s| s.sequence_number).collect::<Vec<_>>(),
876 alloc::vec![sn(2)],
877 "best-effort reader must deliver SN 2 despite the missing SN 1"
878 );
879 assert_eq!(first_state(&r).delivered_up_to, sn(2));
880 let out3 = r.handle_data(p1(), &data(w, rd, 3, 0x33), None);
882 assert_eq!(out3.len(), 1);
883 assert_eq!(out3[0].sequence_number, sn(3));
884 }
885
886 #[test]
887 fn duplicate_data_is_rejected() {
888 let mut r = make_reader(10);
889 let w = single_writer_guid().entity_id;
890 let rd = r.guid().entity_id;
891 r.handle_data(p1(), &data(w, rd, 1, 0xAA), None);
892 let second = r.handle_data(p1(), &data(w, rd, 1, 0xAA), None);
893 assert!(second.is_empty());
894 }
895
896 #[test]
897 fn mismatched_writer_id_is_counted() {
898 let mut r = make_reader(10);
899 let rd = r.guid().entity_id;
900 let foreign = EntityId::user_writer_with_key([0xFF, 0xFF, 0xFF]);
901 assert!(
902 r.handle_data(p1(), &data(foreign, rd, 1, 0xAA), None)
903 .is_empty()
904 );
905 assert_eq!(r.unknown_src_count(), 1);
906 }
907
908 #[test]
911 fn alive_data_yields_alive_changekind() {
912 let mut r = make_reader(10);
913 let w = single_writer_guid().entity_id;
914 let rd = r.guid().entity_id;
915 let delivered = r.handle_data(p1(), &data(w, rd, 1, 0xAA), None);
916 assert_eq!(delivered.len(), 1);
917 assert_eq!(delivered[0].kind, ChangeKind::Alive);
918 }
919
920 fn lifecycle_data(
921 wid: EntityId,
922 rid: EntityId,
923 n: i64,
924 key_hash: [u8; 16],
925 status_bits: u32,
926 ) -> DataSubmessage {
927 DataSubmessage {
928 extra_flags: 0,
929 reader_id: rid,
930 writer_id: wid,
931 writer_sn: sn(n),
932 inline_qos: Some(crate::inline_qos::lifecycle_inline_qos(
933 key_hash,
934 status_bits,
935 )),
936 key_flag: true,
937 non_standard_flag: false,
938 serialized_payload: alloc::sync::Arc::from(alloc::vec![0u8; 0]),
939 }
940 }
941
942 #[test]
943 fn dispose_data_yields_not_alive_disposed() {
944 let mut r = make_reader(10);
945 let w = single_writer_guid().entity_id;
946 let rd = r.guid().entity_id;
947 let delivered = r.handle_data(
948 p1(),
949 &lifecycle_data(
950 w,
951 rd,
952 1,
953 [0xAB; 16],
954 crate::inline_qos::status_info::DISPOSED,
955 ),
956 None,
957 );
958 assert_eq!(delivered.len(), 1);
959 assert_eq!(delivered[0].kind, ChangeKind::NotAliveDisposed);
960 }
961
962 #[test]
963 fn unregister_data_yields_not_alive_unregistered() {
964 let mut r = make_reader(10);
965 let w = single_writer_guid().entity_id;
966 let rd = r.guid().entity_id;
967 let delivered = r.handle_data(
968 p1(),
969 &lifecycle_data(
970 w,
971 rd,
972 1,
973 [0xCD; 16],
974 crate::inline_qos::status_info::UNREGISTERED,
975 ),
976 None,
977 );
978 assert_eq!(delivered.len(), 1);
979 assert_eq!(delivered[0].kind, ChangeKind::NotAliveUnregistered);
980 }
981
982 #[test]
983 fn dispose_and_unregister_combined() {
984 let mut r = make_reader(10);
985 let w = single_writer_guid().entity_id;
986 let rd = r.guid().entity_id;
987 let bits =
988 crate::inline_qos::status_info::DISPOSED | crate::inline_qos::status_info::UNREGISTERED;
989 let delivered = r.handle_data(p1(), &lifecycle_data(w, rd, 1, [0xEF; 16], bits), None);
990 assert_eq!(delivered.len(), 1);
991 assert_eq!(delivered[0].kind, ChangeKind::NotAliveDisposedUnregistered);
992 }
993
994 #[test]
995 fn key_only_alive_registration_is_acked_but_not_delivered() {
996 let mut r = make_reader(10);
1004 let w = single_writer_guid().entity_id;
1005 let rd = r.guid().entity_id;
1006 let mut d = data(w, rd, 1, 0xAA);
1007 d.key_flag = true;
1008 let delivered = r.handle_data(p1(), &d, None);
1009 assert!(
1010 delivered.is_empty(),
1011 "key-only ALIVE registration must not be delivered"
1012 );
1013 let d2 = data(w, rd, 2, 0xBB);
1016 let delivered2 = r.handle_data(p1(), &d2, None);
1017 assert_eq!(delivered2.len(), 1);
1018 assert_eq!(delivered2[0].sequence_number, SequenceNumber(2));
1019 }
1020
1021 #[test]
1022 fn heartbeat_with_missing_triggers_acknack_after_delay() {
1023 let mut r = make_reader(10);
1024 let w = single_writer_guid().entity_id;
1025 let rd = r.guid().entity_id;
1026 r.handle_heartbeat(p1(), &heartbeat(w, rd, 1, 3, 1, false), Duration::ZERO);
1027 assert!(r.tick(Duration::from_millis(100)).unwrap().is_empty());
1028 let out = r.tick(Duration::from_millis(250)).unwrap();
1029 assert_eq!(out.len(), 1);
1030 }
1031
1032 #[test]
1033 fn heartbeat_without_missing_and_final_schedules_no_acknack() {
1034 let mut r = make_reader(10);
1035 let w = single_writer_guid().entity_id;
1036 let rd = r.guid().entity_id;
1037 r.handle_data(p1(), &data(w, rd, 1, 0xAA), None);
1038 r.handle_heartbeat(p1(), &heartbeat(w, rd, 1, 1, 1, true), Duration::ZERO);
1039 assert!(r.tick(Duration::from_secs(10)).unwrap().is_empty());
1040 }
1041
1042 fn second_writer_guid() -> Guid {
1045 Guid::new(
1046 GuidPrefix::from_bytes([3; 12]),
1047 EntityId::user_writer_with_key([0x40, 0x50, 0x60]),
1048 )
1049 }
1050
1051 fn add_second_writer(r: &mut ReliableReader) {
1052 r.add_writer_proxy(WriterProxy::new(
1053 second_writer_guid(),
1054 alloc::vec![Locator::udp_v4([127, 0, 0, 2], 7420)],
1055 alloc::vec![],
1056 true,
1057 ));
1058 }
1059
1060 #[test]
1061 fn add_writer_proxy_increases_count() {
1062 let mut r = make_reader(10);
1063 add_second_writer(&mut r);
1064 assert_eq!(r.writer_proxy_count(), 2);
1065 }
1066
1067 #[test]
1068 fn two_writers_with_overlapping_sn_spaces_both_delivered() {
1069 let mut r = make_reader(10);
1072 add_second_writer(&mut r);
1073 let w1 = single_writer_guid().entity_id;
1074 let w2 = second_writer_guid().entity_id;
1075 let rd = r.guid().entity_id;
1076
1077 let d1 = r.handle_data(p1(), &data(w1, rd, 1, 0xAA), None);
1078 let d2 = r.handle_data(p2(), &data(w2, rd, 1, 0xBB), None);
1079
1080 assert_eq!(d1.len(), 1);
1081 assert_eq!(d1[0].payload.as_ref(), &[0xAA][..]);
1082 assert_eq!(d1[0].writer_guid, single_writer_guid());
1083 assert_eq!(d2.len(), 1);
1084 assert_eq!(d2[0].payload.as_ref(), &[0xBB][..]);
1085 assert_eq!(d2[0].writer_guid, second_writer_guid());
1086
1087 assert_eq!(r.writer_proxies()[0].delivered_up_to, sn(1));
1088 assert_eq!(r.writer_proxies()[1].delivered_up_to, sn(1));
1089 }
1090
1091 #[test]
1092 fn same_entity_id_different_prefix_not_confused() {
1093 let mut r = make_reader(10);
1099 let eid = single_writer_guid().entity_id; let prefix_b = GuidPrefix::from_bytes([9; 12]);
1101 let guid_b = Guid::new(prefix_b, eid);
1102 r.add_writer_proxy(WriterProxy::new(
1103 guid_b,
1104 alloc::vec![Locator::udp_v4([127, 0, 0, 9], 7420)],
1105 alloc::vec![],
1106 true,
1107 ));
1108 assert_eq!(r.writer_proxy_count(), 2);
1109 let rd = r.guid().entity_id;
1110 let d = r.handle_data(prefix_b, &data(eid, rd, 1, 0xBB), None);
1112 assert_eq!(d.len(), 1);
1113 assert_eq!(d[0].writer_guid, guid_b, "sample misattributed");
1114 assert_eq!(r.writer_proxies()[0].delivered_up_to, sn(0));
1116 assert_eq!(r.writer_proxies()[1].delivered_up_to, sn(1));
1117 }
1118
1119 #[test]
1120 fn remove_writer_proxy_drops_its_state() {
1121 let mut r = make_reader(10);
1122 add_second_writer(&mut r);
1123 let removed = r.remove_writer_proxy(single_writer_guid());
1124 assert!(removed.is_some());
1125 assert_eq!(r.writer_proxy_count(), 1);
1126 assert_eq!(
1127 r.writer_proxies()[0].proxy.remote_writer_guid,
1128 second_writer_guid()
1129 );
1130 }
1131
1132 #[test]
1133 fn tick_emits_one_acknack_per_writer_with_missing() {
1134 let mut r = make_reader(10);
1135 add_second_writer(&mut r);
1136 let rd = r.guid().entity_id;
1137 r.handle_heartbeat(
1139 p1(),
1140 &heartbeat(single_writer_guid().entity_id, rd, 1, 3, 1, false),
1141 Duration::ZERO,
1142 );
1143 r.handle_heartbeat(
1144 p2(),
1145 &heartbeat(second_writer_guid().entity_id, rd, 1, 5, 1, false),
1146 Duration::ZERO,
1147 );
1148 let out = r.tick(Duration::from_millis(250)).unwrap();
1149 assert_eq!(out.len(), 2);
1151 }
1152
1153 #[test]
1160 fn pre_emptive_acknack_emitted_after_add_writer_proxy() {
1161 let reader_guid = Guid::new(
1162 GuidPrefix::from_bytes([2; 12]),
1163 EntityId::user_reader_with_key([0xA0, 0xB0, 0xC0]),
1164 );
1165 let mut r = ReliableReader::new(ReliableReaderConfig {
1166 guid: reader_guid,
1167 vendor_id: VendorId::ZERODDS,
1168 writer_proxies: alloc::vec![],
1169 max_samples_per_proxy: 10,
1170 heartbeat_response_delay: Duration::from_millis(200),
1171 assembler_caps: AssemblerCaps::default(),
1172 });
1173 r.add_writer_proxy(WriterProxy::new(
1174 single_writer_guid(),
1175 alloc::vec![Locator::udp_v4([127, 0, 0, 1], 7420)],
1176 alloc::vec![],
1177 true,
1178 ));
1179 let out = r.tick(Duration::from_millis(250)).unwrap();
1182 assert_eq!(out.len(), 1, "exactly one Pre-Emptive ACKNACK expected");
1183 let parsed = decode_datagram(&out[0]).unwrap();
1184 let ack = parsed
1185 .submessages
1186 .iter()
1187 .find_map(|s| {
1188 if let ParsedSubmessage::AckNack(a) = s {
1189 Some(a)
1190 } else {
1191 None
1192 }
1193 })
1194 .expect("ACKNACK in datagram");
1195 assert_eq!(ack.reader_sn_state.bitmap_base, sn(1));
1196 assert_eq!(ack.reader_sn_state.num_bits, 0);
1197 assert!(
1198 !ack.final_flag,
1199 "Pre-Emptive ACKNACK must be non-final (force HB-response)"
1200 );
1201 }
1202
1203 #[test]
1206 fn no_pre_emptive_acknack_without_proxy() {
1207 let reader_guid = Guid::new(
1208 GuidPrefix::from_bytes([2; 12]),
1209 EntityId::user_reader_with_key([0xA0, 0xB0, 0xC0]),
1210 );
1211 let mut r = ReliableReader::new(ReliableReaderConfig {
1212 guid: reader_guid,
1213 vendor_id: VendorId::ZERODDS,
1214 writer_proxies: alloc::vec![],
1215 max_samples_per_proxy: 10,
1216 heartbeat_response_delay: Duration::from_millis(200),
1217 assembler_caps: AssemblerCaps::default(),
1218 });
1219 assert!(r.tick(Duration::from_secs(10)).unwrap().is_empty());
1221 }
1222
1223 #[test]
1228 fn initial_proxy_from_config_does_not_send_pre_emptive() {
1229 let mut r = make_reader(10);
1231 assert!(
1233 r.tick(Duration::from_secs(10)).unwrap().is_empty(),
1234 "initial proxy from config must not emit Pre-Emptive"
1235 );
1236 }
1237
1238 #[test]
1239 fn pre_emptive_acknack_carries_info_dst() {
1240 let reader_guid = Guid::new(
1244 GuidPrefix::from_bytes([2; 12]),
1245 EntityId::user_reader_with_key([0xA0, 0xB0, 0xC0]),
1246 );
1247 let mut r = ReliableReader::new(ReliableReaderConfig {
1248 guid: reader_guid,
1249 vendor_id: VendorId::ZERODDS,
1250 writer_proxies: alloc::vec![],
1251 max_samples_per_proxy: 10,
1252 heartbeat_response_delay: Duration::from_millis(200),
1253 assembler_caps: AssemblerCaps::default(),
1254 });
1255 r.add_writer_proxy(WriterProxy::new(
1256 single_writer_guid(),
1257 alloc::vec![Locator::udp_v4([127, 0, 0, 1], 7420)],
1258 alloc::vec![],
1259 true,
1260 ));
1261 let out = r.tick(Duration::from_millis(250)).unwrap();
1262 assert_eq!(out.len(), 1);
1263 let parsed = decode_datagram(&out[0]).unwrap();
1264 assert!(parsed.submessages.len() >= 2, "INFO_DST + ACKNACK");
1267 match &parsed.submessages[0] {
1268 ParsedSubmessage::Unknown { id, .. } => assert_eq!(*id, 0x0E),
1269 other => panic!("expected INFO_DST first, got {other:?}"),
1270 }
1271 }
1272
1273 #[test]
1274 fn unknown_writer_id_in_heartbeat_counts_not_crashes() {
1275 let mut r = make_reader(10);
1276 let rd = r.guid().entity_id;
1277 let foreign = EntityId::user_writer_with_key([0xFF, 0xFF, 0xFF]);
1278 r.handle_heartbeat(
1279 p1(),
1280 &heartbeat(foreign, rd, 1, 3, 1, false),
1281 Duration::ZERO,
1282 );
1283 assert_eq!(r.unknown_src_count(), 1);
1284 assert!(r.tick(Duration::from_secs(1)).unwrap().is_empty());
1285 }
1286}