1extern crate alloc;
29use alloc::sync::Arc;
30use alloc::vec::Vec;
31
32use crate::error::WireError;
33use crate::submessage_header::FLAG_E_LITTLE_ENDIAN;
34use crate::wire_types::{EntityId, FragmentNumber, SequenceNumber};
35
36pub const RTPS_BITMAP_MAX_BITS: u32 = 256;
42
43#[derive(Debug, Clone, PartialEq, Eq)]
55pub struct SequenceNumberSet {
56 pub bitmap_base: SequenceNumber,
58 pub num_bits: u32,
60 pub bitmap: Vec<u32>,
62}
63
64impl SequenceNumberSet {
65 #[must_use]
67 pub fn wire_size(num_bits: u32) -> usize {
68 let words = (num_bits as usize).div_ceil(32);
69 8 + 4 + words * 4
70 }
71
72 #[must_use]
79 pub fn from_missing(base: SequenceNumber, missing: &[SequenceNumber]) -> Self {
80 let Some(last) = missing.last().copied() else {
81 return Self {
82 bitmap_base: base,
83 num_bits: 0,
84 bitmap: Vec::new(),
85 };
86 };
87 if last < base {
88 return Self {
89 bitmap_base: base,
90 num_bits: 0,
91 bitmap: Vec::new(),
92 };
93 }
94 let num_bits = u32::try_from(last.0 - base.0 + 1).unwrap_or(u32::MAX);
95 let num_words = (num_bits as usize).div_ceil(32);
96 let mut bitmap = alloc::vec![0u32; num_words];
97 for sn in missing {
98 if *sn < base {
99 continue;
100 }
101 let offset = (sn.0 - base.0) as usize;
102 let word_idx = offset / 32;
103 let bit = 31 - (offset % 32);
104 if word_idx < bitmap.len() {
105 bitmap[word_idx] |= 1u32 << bit;
106 }
107 }
108 Self {
109 bitmap_base: base,
110 num_bits,
111 bitmap,
112 }
113 }
114
115 pub fn iter_set(&self) -> impl Iterator<Item = SequenceNumber> + '_ {
117 (0..self.num_bits).filter_map(move |i| {
118 let word_idx = (i / 32) as usize;
119 let bit = 31 - (i as usize % 32);
120 if word_idx < self.bitmap.len() && (self.bitmap[word_idx] >> bit) & 1 == 1 {
121 Some(SequenceNumber(self.bitmap_base.0 + i64::from(i)))
122 } else {
123 None
124 }
125 })
126 }
127
128 #[must_use]
130 pub fn encoded_size(&self) -> usize {
131 Self::wire_size(self.num_bits)
132 }
133
134 pub fn write_to(&self, out: &mut Vec<u8>, little_endian: bool) {
136 if little_endian {
137 out.extend_from_slice(&self.bitmap_base.to_bytes_le());
138 out.extend_from_slice(&self.num_bits.to_le_bytes());
139 for w in &self.bitmap {
140 out.extend_from_slice(&w.to_le_bytes());
141 }
142 } else {
143 out.extend_from_slice(&self.bitmap_base.to_bytes_be());
144 out.extend_from_slice(&self.num_bits.to_be_bytes());
145 for w in &self.bitmap {
146 out.extend_from_slice(&w.to_be_bytes());
147 }
148 }
149 }
150
151 pub fn read_from(
156 bytes: &[u8],
157 offset: usize,
158 little_endian: bool,
159 ) -> Result<(Self, usize), WireError> {
160 let mut pos = offset;
161 if bytes.len() < pos + 8 {
162 return Err(WireError::UnexpectedEof {
163 needed: 8,
164 offset: pos,
165 });
166 }
167 let mut sn_bytes = [0u8; 8];
168 sn_bytes.copy_from_slice(&bytes[pos..pos + 8]);
169 let bitmap_base = if little_endian {
170 SequenceNumber::from_bytes_le(sn_bytes)
171 } else {
172 SequenceNumber::from_bytes_be(sn_bytes)
173 };
174 pos += 8;
175 if bytes.len() < pos + 4 {
176 return Err(WireError::UnexpectedEof {
177 needed: 4,
178 offset: pos,
179 });
180 }
181 let mut num_bytes = [0u8; 4];
182 num_bytes.copy_from_slice(&bytes[pos..pos + 4]);
183 let num_bits = if little_endian {
184 u32::from_le_bytes(num_bytes)
185 } else {
186 u32::from_be_bytes(num_bytes)
187 };
188 pos += 4;
189 if num_bits > RTPS_BITMAP_MAX_BITS {
190 return Err(WireError::ValueOutOfRange {
191 message: "SequenceNumberSet.numBits exceeds RTPS_BITMAP_MAX_BITS (256)",
192 });
193 }
194 let words = (num_bits as usize).div_ceil(32);
195 let bitmap_bytes = words * 4;
196 if bytes.len() < pos + bitmap_bytes {
197 return Err(WireError::UnexpectedEof {
198 needed: bitmap_bytes,
199 offset: pos,
200 });
201 }
202 let mut bitmap = Vec::with_capacity(words);
203 for _ in 0..words {
204 let mut w = [0u8; 4];
205 w.copy_from_slice(&bytes[pos..pos + 4]);
206 bitmap.push(if little_endian {
207 u32::from_le_bytes(w)
208 } else {
209 u32::from_be_bytes(w)
210 });
211 pos += 4;
212 }
213 Ok((
214 Self {
215 bitmap_base,
216 num_bits,
217 bitmap,
218 },
219 pos,
220 ))
221 }
222}
223
224#[derive(Debug, Clone, PartialEq, Eq)]
237pub struct FragmentNumberSet {
238 pub bitmap_base: FragmentNumber,
240 pub num_bits: u32,
242 pub bitmap: Vec<u32>,
244}
245
246impl FragmentNumberSet {
247 #[must_use]
249 pub fn wire_size(num_bits: u32) -> usize {
250 let words = (num_bits as usize).div_ceil(32);
251 4 + 4 + words * 4
252 }
253
254 #[must_use]
257 pub fn from_missing(base: FragmentNumber, missing: &[FragmentNumber]) -> Self {
258 let Some(last) = missing.last().copied() else {
259 return Self {
260 bitmap_base: base,
261 num_bits: 0,
262 bitmap: Vec::new(),
263 };
264 };
265 if last < base {
266 return Self {
267 bitmap_base: base,
268 num_bits: 0,
269 bitmap: Vec::new(),
270 };
271 }
272 let num_bits = last.0.saturating_sub(base.0).saturating_add(1).min(256);
279 let num_words = (num_bits as usize).div_ceil(32);
280 let mut bitmap = alloc::vec![0u32; num_words];
281 for fnum in missing {
282 if *fnum < base {
283 continue;
284 }
285 let offset = (fnum.0 - base.0) as usize;
286 if offset >= num_bits as usize {
288 continue;
289 }
290 let word_idx = offset / 32;
291 let bit = 31 - (offset % 32);
292 if word_idx < bitmap.len() {
293 bitmap[word_idx] |= 1u32 << bit;
294 }
295 }
296 Self {
297 bitmap_base: base,
298 num_bits,
299 bitmap,
300 }
301 }
302
303 pub fn iter_set(&self) -> impl Iterator<Item = FragmentNumber> + '_ {
305 (0..self.num_bits).filter_map(move |i| {
306 let word_idx = (i / 32) as usize;
307 let bit = 31 - (i as usize % 32);
308 if word_idx < self.bitmap.len() && (self.bitmap[word_idx] >> bit) & 1 == 1 {
309 Some(FragmentNumber(self.bitmap_base.0.wrapping_add(i)))
310 } else {
311 None
312 }
313 })
314 }
315
316 #[must_use]
318 pub fn encoded_size(&self) -> usize {
319 Self::wire_size(self.num_bits)
320 }
321
322 pub fn write_to(&self, out: &mut Vec<u8>, little_endian: bool) {
324 if little_endian {
325 out.extend_from_slice(&self.bitmap_base.to_bytes_le());
326 out.extend_from_slice(&self.num_bits.to_le_bytes());
327 for w in &self.bitmap {
328 out.extend_from_slice(&w.to_le_bytes());
329 }
330 } else {
331 out.extend_from_slice(&self.bitmap_base.to_bytes_be());
332 out.extend_from_slice(&self.num_bits.to_be_bytes());
333 for w in &self.bitmap {
334 out.extend_from_slice(&w.to_be_bytes());
335 }
336 }
337 }
338
339 pub fn read_from(
344 bytes: &[u8],
345 offset: usize,
346 little_endian: bool,
347 ) -> Result<(Self, usize), WireError> {
348 let mut pos = offset;
349 if bytes.len() < pos + 4 {
350 return Err(WireError::UnexpectedEof {
351 needed: 4,
352 offset: pos,
353 });
354 }
355 let mut bb = [0u8; 4];
356 bb.copy_from_slice(&bytes[pos..pos + 4]);
357 let bitmap_base = if little_endian {
358 FragmentNumber::from_bytes_le(bb)
359 } else {
360 FragmentNumber::from_bytes_be(bb)
361 };
362 pos += 4;
363 if bytes.len() < pos + 4 {
364 return Err(WireError::UnexpectedEof {
365 needed: 4,
366 offset: pos,
367 });
368 }
369 let mut nb = [0u8; 4];
370 nb.copy_from_slice(&bytes[pos..pos + 4]);
371 let num_bits = if little_endian {
372 u32::from_le_bytes(nb)
373 } else {
374 u32::from_be_bytes(nb)
375 };
376 pos += 4;
377 if num_bits > RTPS_BITMAP_MAX_BITS {
378 return Err(WireError::ValueOutOfRange {
379 message: "FragmentNumberSet.numBits exceeds RTPS_BITMAP_MAX_BITS (256)",
380 });
381 }
382 let words = (num_bits as usize).div_ceil(32);
383 let need = words * 4;
384 if bytes.len() < pos + need {
385 return Err(WireError::UnexpectedEof {
386 needed: need,
387 offset: pos,
388 });
389 }
390 let mut bitmap = Vec::with_capacity(words);
391 for _ in 0..words {
392 let mut w = [0u8; 4];
393 w.copy_from_slice(&bytes[pos..pos + 4]);
394 bitmap.push(if little_endian {
395 u32::from_le_bytes(w)
396 } else {
397 u32::from_be_bytes(w)
398 });
399 pos += 4;
400 }
401 Ok((
402 Self {
403 bitmap_base,
404 num_bits,
405 bitmap,
406 },
407 pos,
408 ))
409 }
410}
411
412pub const DATA_FLAG_INLINE_QOS: u8 = 0x02;
418pub const DATA_FLAG_DATA: u8 = 0x04;
420pub const DATA_FLAG_KEY: u8 = 0x08;
422pub const DATA_FLAG_NON_STANDARD: u8 = 0x10;
424
425#[derive(Debug, Clone, PartialEq, Eq)]
433pub struct DataSubmessage {
434 pub extra_flags: u16,
436 pub reader_id: EntityId,
438 pub writer_id: EntityId,
440 pub writer_sn: SequenceNumber,
442 pub inline_qos: Option<crate::parameter_list::ParameterList>,
446 pub key_flag: bool,
452 pub non_standard_flag: bool,
457 pub serialized_payload: Arc<[u8]>,
459}
460
461impl DataSubmessage {
462 #[must_use]
470 pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
471 let inline_qos_buf = self
477 .inline_qos
478 .as_ref()
479 .map(|pl| pl.to_bytes(little_endian));
480 let inline_qos_len = inline_qos_buf.as_ref().map_or(0, |v| v.len());
481 let mut out = Vec::with_capacity(20 + inline_qos_len + self.serialized_payload.len());
484 let extra = if little_endian {
486 self.extra_flags.to_le_bytes()
487 } else {
488 self.extra_flags.to_be_bytes()
489 };
490 out.extend_from_slice(&extra);
491 let octets_to_inline_qos: u16 = 16;
495 let oti = if little_endian {
496 octets_to_inline_qos.to_le_bytes()
497 } else {
498 octets_to_inline_qos.to_be_bytes()
499 };
500 out.extend_from_slice(&oti);
501 out.extend_from_slice(&self.reader_id.to_bytes());
503 out.extend_from_slice(&self.writer_id.to_bytes());
504 out.extend_from_slice(&if little_endian {
506 self.writer_sn.to_bytes_le()
507 } else {
508 self.writer_sn.to_bytes_be()
509 });
510 if let Some(qos_bytes) = inline_qos_buf {
512 out.extend_from_slice(&qos_bytes);
513 }
514 out.extend_from_slice(&self.serialized_payload);
516
517 let mut flags = 0u8;
518 if little_endian {
519 flags |= FLAG_E_LITTLE_ENDIAN;
520 }
521 flags |= DATA_FLAG_DATA;
522 if self.key_flag {
523 flags |= DATA_FLAG_KEY;
524 }
525 if self.non_standard_flag {
526 flags |= DATA_FLAG_NON_STANDARD;
527 }
528 if self.inline_qos.is_some() {
529 flags |= DATA_FLAG_INLINE_QOS;
530 }
531 (out, flags)
532 }
533
534 pub fn read_body(body: &[u8], little_endian: bool) -> Result<Self, WireError> {
540 Self::read_body_with_flags(body, little_endian, 0)
541 }
542
543 pub fn read_body_with_flags(
551 body: &[u8],
552 little_endian: bool,
553 flags: u8,
554 ) -> Result<Self, WireError> {
555 if body.len() < 4 + 4 + 4 + 8 {
556 return Err(WireError::UnexpectedEof {
557 needed: 20,
558 offset: 0,
559 });
560 }
561 let mut pos = 0usize;
562 let mut ef = [0u8; 2];
563 ef.copy_from_slice(&body[pos..pos + 2]);
564 let extra_flags = if little_endian {
565 u16::from_le_bytes(ef)
566 } else {
567 u16::from_be_bytes(ef)
568 };
569 pos += 2;
570 pos += 2;
573 let mut rid = [0u8; 4];
574 rid.copy_from_slice(&body[pos..pos + 4]);
575 let reader_id = EntityId::from_bytes(rid);
576 pos += 4;
577 let mut wid = [0u8; 4];
578 wid.copy_from_slice(&body[pos..pos + 4]);
579 let writer_id = EntityId::from_bytes(wid);
580 pos += 4;
581 let mut sn = [0u8; 8];
582 sn.copy_from_slice(&body[pos..pos + 8]);
583 let writer_sn = if little_endian {
584 SequenceNumber::from_bytes_le(sn)
585 } else {
586 SequenceNumber::from_bytes_be(sn)
587 };
588 pos += 8;
589
590 let inline_qos = if flags & DATA_FLAG_INLINE_QOS != 0 {
592 let pl = crate::parameter_list::ParameterList::from_bytes(&body[pos..], little_endian)?;
598 let consumed = pl.to_bytes(little_endian).len();
601 pos += consumed;
602 Some(pl)
603 } else {
604 None
605 };
606
607 let serialized_payload: Arc<[u8]> = Arc::from(&body[pos..]);
609 let key_flag = (flags & DATA_FLAG_KEY) != 0;
610 let non_standard_flag = (flags & DATA_FLAG_NON_STANDARD) != 0;
611 Ok(Self {
612 extra_flags,
613 reader_id,
614 writer_id,
615 writer_sn,
616 inline_qos,
617 key_flag,
618 non_standard_flag,
619 serialized_payload,
620 })
621 }
622}
623
624pub const HEARTBEAT_FLAG_FINAL: u8 = 0x02;
630pub const HEARTBEAT_FLAG_LIVELINESS: u8 = 0x04;
632pub const HEARTBEAT_FLAG_GROUP_INFO: u8 = 0x08;
637
638#[derive(Debug, Clone, PartialEq, Eq)]
646pub struct HeartbeatGroupInfo {
647 pub current_gsn: SequenceNumber,
649 pub first_gsn: SequenceNumber,
651 pub last_gsn: SequenceNumber,
654 pub writer_set: Vec<crate::wire_types::GuidPrefix>,
656}
657
658#[derive(Debug, Clone, PartialEq, Eq)]
665pub struct HeartbeatSubmessage {
666 pub reader_id: EntityId,
668 pub writer_id: EntityId,
670 pub first_sn: SequenceNumber,
672 pub last_sn: SequenceNumber,
674 pub count: i32,
676 pub final_flag: bool,
678 pub liveliness_flag: bool,
680 pub group_info: Option<HeartbeatGroupInfo>,
682}
683
684impl HeartbeatSubmessage {
685 pub const WIRE_SIZE: usize = 28;
688
689 #[must_use]
692 pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
693 let mut out = Vec::with_capacity(Self::WIRE_SIZE);
694 out.extend_from_slice(&self.reader_id.to_bytes());
695 out.extend_from_slice(&self.writer_id.to_bytes());
696 out.extend_from_slice(&if little_endian {
697 self.first_sn.to_bytes_le()
698 } else {
699 self.first_sn.to_bytes_be()
700 });
701 out.extend_from_slice(&if little_endian {
702 self.last_sn.to_bytes_le()
703 } else {
704 self.last_sn.to_bytes_be()
705 });
706 out.extend_from_slice(&if little_endian {
707 self.count.to_le_bytes()
708 } else {
709 self.count.to_be_bytes()
710 });
711 let mut flags = 0u8;
712 if little_endian {
713 flags |= FLAG_E_LITTLE_ENDIAN;
714 }
715 if self.final_flag {
716 flags |= HEARTBEAT_FLAG_FINAL;
717 }
718 if self.liveliness_flag {
719 flags |= HEARTBEAT_FLAG_LIVELINESS;
720 }
721 if let Some(gi) = &self.group_info {
722 flags |= HEARTBEAT_FLAG_GROUP_INFO;
723 for sn in [gi.current_gsn, gi.first_gsn, gi.last_gsn] {
724 out.extend_from_slice(&if little_endian {
725 sn.to_bytes_le()
726 } else {
727 sn.to_bytes_be()
728 });
729 }
730 let len = u32::try_from(gi.writer_set.len()).unwrap_or(u32::MAX);
731 out.extend_from_slice(&if little_endian {
732 len.to_le_bytes()
733 } else {
734 len.to_be_bytes()
735 });
736 for prefix in &gi.writer_set {
737 out.extend_from_slice(&prefix.to_bytes());
738 }
739 }
740 (out, flags)
741 }
742
743 pub fn read_body(
749 body: &[u8],
750 little_endian: bool,
751 final_flag: bool,
752 liveliness_flag: bool,
753 group_info_flag: bool,
754 ) -> Result<Self, WireError> {
755 if body.len() < Self::WIRE_SIZE {
756 return Err(WireError::UnexpectedEof {
757 needed: Self::WIRE_SIZE,
758 offset: 0,
759 });
760 }
761 let mut pos = 0usize;
762 let mut rid = [0u8; 4];
763 rid.copy_from_slice(&body[pos..pos + 4]);
764 let reader_id = EntityId::from_bytes(rid);
765 pos += 4;
766 let mut wid = [0u8; 4];
767 wid.copy_from_slice(&body[pos..pos + 4]);
768 let writer_id = EntityId::from_bytes(wid);
769 pos += 4;
770 let mut sn = [0u8; 8];
771 sn.copy_from_slice(&body[pos..pos + 8]);
772 let first_sn = if little_endian {
773 SequenceNumber::from_bytes_le(sn)
774 } else {
775 SequenceNumber::from_bytes_be(sn)
776 };
777 pos += 8;
778 sn.copy_from_slice(&body[pos..pos + 8]);
779 let last_sn = if little_endian {
780 SequenceNumber::from_bytes_le(sn)
781 } else {
782 SequenceNumber::from_bytes_be(sn)
783 };
784 pos += 8;
785 let mut cnt = [0u8; 4];
786 cnt.copy_from_slice(&body[pos..pos + 4]);
787 let count = if little_endian {
788 i32::from_le_bytes(cnt)
789 } else {
790 i32::from_be_bytes(cnt)
791 };
792 pos += 4;
793 let group_info = if group_info_flag {
794 if body.len() < pos + 28 {
796 return Err(WireError::UnexpectedEof {
797 needed: 28,
798 offset: pos,
799 });
800 }
801 let mut s = [0u8; 8];
802 s.copy_from_slice(&body[pos..pos + 8]);
803 let current_gsn = if little_endian {
804 SequenceNumber::from_bytes_le(s)
805 } else {
806 SequenceNumber::from_bytes_be(s)
807 };
808 pos += 8;
809 s.copy_from_slice(&body[pos..pos + 8]);
810 let first_gsn = if little_endian {
811 SequenceNumber::from_bytes_le(s)
812 } else {
813 SequenceNumber::from_bytes_be(s)
814 };
815 pos += 8;
816 s.copy_from_slice(&body[pos..pos + 8]);
817 let last_gsn = if little_endian {
818 SequenceNumber::from_bytes_le(s)
819 } else {
820 SequenceNumber::from_bytes_be(s)
821 };
822 pos += 8;
823 let mut len_bytes = [0u8; 4];
824 len_bytes.copy_from_slice(&body[pos..pos + 4]);
825 let len = if little_endian {
826 u32::from_le_bytes(len_bytes)
827 } else {
828 u32::from_be_bytes(len_bytes)
829 } as usize;
830 pos += 4;
831 let remaining = body.len().saturating_sub(pos);
834 if len.saturating_mul(12) > remaining {
835 return Err(WireError::ValueOutOfRange {
836 message: "HEARTBEAT.groupInfo.writerSet length exceeds body",
837 });
838 }
839 let mut writer_set = Vec::with_capacity(len);
840 for _ in 0..len {
841 let mut p = [0u8; 12];
842 p.copy_from_slice(&body[pos..pos + 12]);
843 writer_set.push(crate::wire_types::GuidPrefix::from_bytes(p));
844 pos += 12;
845 }
846 Some(HeartbeatGroupInfo {
847 current_gsn,
848 first_gsn,
849 last_gsn,
850 writer_set,
851 })
852 } else {
853 None
854 };
855 Ok(Self {
856 reader_id,
857 writer_id,
858 first_sn,
859 last_sn,
860 count,
861 final_flag,
862 liveliness_flag,
863 group_info,
864 })
865 }
866}
867
868pub const ACKNACK_FLAG_FINAL: u8 = 0x02;
874
875#[derive(Debug, Clone, PartialEq, Eq)]
881pub struct AckNackSubmessage {
882 pub reader_id: EntityId,
884 pub writer_id: EntityId,
886 pub reader_sn_state: SequenceNumberSet,
888 pub count: i32,
890 pub final_flag: bool,
892}
893
894impl AckNackSubmessage {
895 #[must_use]
897 pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
898 let snset_words = self.reader_sn_state.bitmap.len();
903 let mut out = Vec::with_capacity(4 + 4 + 12 + snset_words * 4 + 4);
904 out.extend_from_slice(&self.reader_id.to_bytes());
905 out.extend_from_slice(&self.writer_id.to_bytes());
906 self.reader_sn_state.write_to(&mut out, little_endian);
907 out.extend_from_slice(&if little_endian {
908 self.count.to_le_bytes()
909 } else {
910 self.count.to_be_bytes()
911 });
912 let mut flags = 0u8;
913 if little_endian {
914 flags |= FLAG_E_LITTLE_ENDIAN;
915 }
916 if self.final_flag {
917 flags |= ACKNACK_FLAG_FINAL;
918 }
919 (out, flags)
920 }
921
922 pub fn read_body(
928 body: &[u8],
929 little_endian: bool,
930 final_flag: bool,
931 ) -> Result<Self, WireError> {
932 if body.len() < 8 {
933 return Err(WireError::UnexpectedEof {
934 needed: 8,
935 offset: 0,
936 });
937 }
938 let mut pos = 0usize;
939 let mut rid = [0u8; 4];
940 rid.copy_from_slice(&body[pos..pos + 4]);
941 let reader_id = EntityId::from_bytes(rid);
942 pos += 4;
943 let mut wid = [0u8; 4];
944 wid.copy_from_slice(&body[pos..pos + 4]);
945 let writer_id = EntityId::from_bytes(wid);
946 pos += 4;
947 let (reader_sn_state, new_pos) = SequenceNumberSet::read_from(body, pos, little_endian)?;
948 pos = new_pos;
949 if body.len() < pos + 4 {
950 return Err(WireError::UnexpectedEof {
951 needed: 4,
952 offset: pos,
953 });
954 }
955 let mut cnt = [0u8; 4];
956 cnt.copy_from_slice(&body[pos..pos + 4]);
957 let count = if little_endian {
958 i32::from_le_bytes(cnt)
959 } else {
960 i32::from_be_bytes(cnt)
961 };
962 Ok(Self {
963 reader_id,
964 writer_id,
965 reader_sn_state,
966 count,
967 final_flag,
968 })
969 }
970}
971
972pub const GAP_FLAG_GROUP_INFO: u8 = 0x04;
979
980pub const GAP_FLAG_FILTERED_COUNT: u8 = 0x08;
986
987#[derive(Debug, Clone, Copy, PartialEq, Eq)]
989pub struct GapGroupInfo {
990 pub gap_start_gsn: SequenceNumber,
992 pub gap_end_gsn: SequenceNumber,
994}
995
996#[derive(Debug, Clone, PartialEq, Eq)]
1001pub struct GapSubmessage {
1002 pub reader_id: EntityId,
1004 pub writer_id: EntityId,
1006 pub gap_start: SequenceNumber,
1008 pub gap_list: SequenceNumberSet,
1010 pub group_info: Option<GapGroupInfo>,
1012 pub filtered_count: Option<u32>,
1017}
1018
1019impl GapSubmessage {
1020 #[must_use]
1022 pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
1023 let snset_words = self.gap_list.bitmap.len();
1026 let extra =
1027 self.group_info.as_ref().map_or(0, |_| 24) + self.filtered_count.map_or(0, |_| 4);
1028 let mut out = Vec::with_capacity(4 + 4 + 8 + 12 + snset_words * 4 + extra);
1029 out.extend_from_slice(&self.reader_id.to_bytes());
1030 out.extend_from_slice(&self.writer_id.to_bytes());
1031 out.extend_from_slice(&if little_endian {
1032 self.gap_start.to_bytes_le()
1033 } else {
1034 self.gap_start.to_bytes_be()
1035 });
1036 self.gap_list.write_to(&mut out, little_endian);
1037 let mut flags = 0u8;
1038 if little_endian {
1039 flags |= FLAG_E_LITTLE_ENDIAN;
1040 }
1041 if let Some(gi) = self.group_info {
1042 flags |= GAP_FLAG_GROUP_INFO;
1043 out.extend_from_slice(&if little_endian {
1044 gi.gap_start_gsn.to_bytes_le()
1045 } else {
1046 gi.gap_start_gsn.to_bytes_be()
1047 });
1048 out.extend_from_slice(&if little_endian {
1049 gi.gap_end_gsn.to_bytes_le()
1050 } else {
1051 gi.gap_end_gsn.to_bytes_be()
1052 });
1053 }
1054 if let Some(fc) = self.filtered_count {
1055 flags |= GAP_FLAG_FILTERED_COUNT;
1056 out.extend_from_slice(&if little_endian {
1057 fc.to_le_bytes()
1058 } else {
1059 fc.to_be_bytes()
1060 });
1061 }
1062 (out, flags)
1063 }
1064
1065 pub fn read_body(
1071 body: &[u8],
1072 little_endian: bool,
1073 group_info_flag: bool,
1074 filtered_count_flag: bool,
1075 ) -> Result<Self, WireError> {
1076 if body.len() < 4 + 4 + 8 {
1077 return Err(WireError::UnexpectedEof {
1078 needed: 16,
1079 offset: 0,
1080 });
1081 }
1082 let mut pos = 0usize;
1083 let mut rid = [0u8; 4];
1084 rid.copy_from_slice(&body[pos..pos + 4]);
1085 let reader_id = EntityId::from_bytes(rid);
1086 pos += 4;
1087 let mut wid = [0u8; 4];
1088 wid.copy_from_slice(&body[pos..pos + 4]);
1089 let writer_id = EntityId::from_bytes(wid);
1090 pos += 4;
1091 let mut sn = [0u8; 8];
1092 sn.copy_from_slice(&body[pos..pos + 8]);
1093 let gap_start = if little_endian {
1094 SequenceNumber::from_bytes_le(sn)
1095 } else {
1096 SequenceNumber::from_bytes_be(sn)
1097 };
1098 pos += 8;
1099 let (gap_list, new_pos) = SequenceNumberSet::read_from(body, pos, little_endian)?;
1100 pos = new_pos;
1101 let group_info = if group_info_flag {
1102 if body.len() < pos + 16 {
1103 return Err(WireError::UnexpectedEof {
1104 needed: 16,
1105 offset: pos,
1106 });
1107 }
1108 let mut s = [0u8; 8];
1109 s.copy_from_slice(&body[pos..pos + 8]);
1110 let gap_start_gsn = if little_endian {
1111 SequenceNumber::from_bytes_le(s)
1112 } else {
1113 SequenceNumber::from_bytes_be(s)
1114 };
1115 pos += 8;
1116 s.copy_from_slice(&body[pos..pos + 8]);
1117 let gap_end_gsn = if little_endian {
1118 SequenceNumber::from_bytes_le(s)
1119 } else {
1120 SequenceNumber::from_bytes_be(s)
1121 };
1122 pos += 8;
1123 Some(GapGroupInfo {
1124 gap_start_gsn,
1125 gap_end_gsn,
1126 })
1127 } else {
1128 None
1129 };
1130 let filtered_count = if filtered_count_flag {
1131 if body.len() < pos + 4 {
1132 return Err(WireError::UnexpectedEof {
1133 needed: 4,
1134 offset: pos,
1135 });
1136 }
1137 let mut c = [0u8; 4];
1138 c.copy_from_slice(&body[pos..pos + 4]);
1139 let fc = if little_endian {
1140 u32::from_le_bytes(c)
1141 } else {
1142 u32::from_be_bytes(c)
1143 };
1144 Some(fc)
1145 } else {
1146 None
1147 };
1148 Ok(Self {
1149 reader_id,
1150 writer_id,
1151 gap_start,
1152 gap_list,
1153 group_info,
1154 filtered_count,
1155 })
1156 }
1157}
1158
1159pub const DATA_FRAG_FLAG_INLINE_QOS: u8 = 0x02;
1165pub const DATA_FRAG_FLAG_HASH_KEY: u8 = 0x04;
1167pub const DATA_FRAG_FLAG_KEY: u8 = 0x08;
1169pub const DATA_FRAG_FLAG_NON_STANDARD: u8 = 0x10;
1171
1172#[derive(Debug, Clone, PartialEq, Eq)]
1177pub struct DataFragSubmessage {
1178 pub extra_flags: u16,
1181 pub reader_id: EntityId,
1183 pub writer_id: EntityId,
1185 pub writer_sn: SequenceNumber,
1187 pub fragment_starting_num: FragmentNumber,
1189 pub fragments_in_submessage: u16,
1191 pub fragment_size: u16,
1193 pub sample_size: u32,
1195 pub serialized_payload: Arc<[u8]>,
1198 pub inline_qos_flag: bool,
1200 pub hash_key_flag: bool,
1202 pub key_flag: bool,
1204 pub non_standard_flag: bool,
1206}
1207
1208impl DataFragSubmessage {
1209 pub const HEADER_WIRE_SIZE: usize = 32;
1213
1214 pub const OCTETS_TO_INLINE_QOS: u16 = 28;
1218
1219 #[must_use]
1221 pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
1222 let mut out = Vec::with_capacity(Self::HEADER_WIRE_SIZE + self.serialized_payload.len());
1223 if little_endian {
1224 out.extend_from_slice(&self.extra_flags.to_le_bytes());
1225 out.extend_from_slice(&Self::OCTETS_TO_INLINE_QOS.to_le_bytes());
1226 } else {
1227 out.extend_from_slice(&self.extra_flags.to_be_bytes());
1228 out.extend_from_slice(&Self::OCTETS_TO_INLINE_QOS.to_be_bytes());
1229 }
1230 out.extend_from_slice(&self.reader_id.to_bytes());
1231 out.extend_from_slice(&self.writer_id.to_bytes());
1232 out.extend_from_slice(&if little_endian {
1233 self.writer_sn.to_bytes_le()
1234 } else {
1235 self.writer_sn.to_bytes_be()
1236 });
1237 out.extend_from_slice(&if little_endian {
1238 self.fragment_starting_num.to_bytes_le()
1239 } else {
1240 self.fragment_starting_num.to_bytes_be()
1241 });
1242 if little_endian {
1243 out.extend_from_slice(&self.fragments_in_submessage.to_le_bytes());
1244 out.extend_from_slice(&self.fragment_size.to_le_bytes());
1245 out.extend_from_slice(&self.sample_size.to_le_bytes());
1246 } else {
1247 out.extend_from_slice(&self.fragments_in_submessage.to_be_bytes());
1248 out.extend_from_slice(&self.fragment_size.to_be_bytes());
1249 out.extend_from_slice(&self.sample_size.to_be_bytes());
1250 }
1251 out.extend_from_slice(&self.serialized_payload);
1252 let mut flags = 0u8;
1253 if little_endian {
1254 flags |= FLAG_E_LITTLE_ENDIAN;
1255 }
1256 if self.inline_qos_flag {
1257 flags |= DATA_FRAG_FLAG_INLINE_QOS;
1258 }
1259 if self.hash_key_flag {
1260 flags |= DATA_FRAG_FLAG_HASH_KEY;
1261 }
1262 if self.key_flag {
1263 flags |= DATA_FRAG_FLAG_KEY;
1264 }
1265 if self.non_standard_flag {
1266 flags |= DATA_FRAG_FLAG_NON_STANDARD;
1267 }
1268 (out, flags)
1269 }
1270
1271 pub fn read_body(
1277 body: &[u8],
1278 little_endian: bool,
1279 inline_qos_flag: bool,
1280 hash_key_flag: bool,
1281 key_flag: bool,
1282 non_standard_flag: bool,
1283 ) -> Result<Self, WireError> {
1284 if body.len() < Self::HEADER_WIRE_SIZE {
1285 return Err(WireError::UnexpectedEof {
1286 needed: Self::HEADER_WIRE_SIZE,
1287 offset: 0,
1288 });
1289 }
1290 let mut pos = 0usize;
1291 let mut ef = [0u8; 2];
1292 ef.copy_from_slice(&body[pos..pos + 2]);
1293 let extra_flags = if little_endian {
1294 u16::from_le_bytes(ef)
1295 } else {
1296 u16::from_be_bytes(ef)
1297 };
1298 pos += 2;
1299 let mut otq = [0u8; 2];
1307 otq.copy_from_slice(&body[pos..pos + 2]);
1308 let octets_to_inline_qos = if little_endian {
1309 u16::from_le_bytes(otq)
1310 } else {
1311 u16::from_be_bytes(otq)
1312 };
1313 pos += 2;
1314 if !inline_qos_flag && octets_to_inline_qos != Self::OCTETS_TO_INLINE_QOS {
1315 return Err(WireError::ValueOutOfRange {
1316 message: "DATA_FRAG.octetsToInlineQos must equal 28 when Q=false",
1317 });
1318 }
1319 let mut rid = [0u8; 4];
1320 rid.copy_from_slice(&body[pos..pos + 4]);
1321 let reader_id = EntityId::from_bytes(rid);
1322 pos += 4;
1323 let mut wid = [0u8; 4];
1324 wid.copy_from_slice(&body[pos..pos + 4]);
1325 let writer_id = EntityId::from_bytes(wid);
1326 pos += 4;
1327 let mut sn = [0u8; 8];
1328 sn.copy_from_slice(&body[pos..pos + 8]);
1329 let writer_sn = if little_endian {
1330 SequenceNumber::from_bytes_le(sn)
1331 } else {
1332 SequenceNumber::from_bytes_be(sn)
1333 };
1334 pos += 8;
1335 let mut fsn = [0u8; 4];
1336 fsn.copy_from_slice(&body[pos..pos + 4]);
1337 let fragment_starting_num = if little_endian {
1338 FragmentNumber::from_bytes_le(fsn)
1339 } else {
1340 FragmentNumber::from_bytes_be(fsn)
1341 };
1342 pos += 4;
1343 let mut fis = [0u8; 2];
1344 fis.copy_from_slice(&body[pos..pos + 2]);
1345 let fragments_in_submessage = if little_endian {
1346 u16::from_le_bytes(fis)
1347 } else {
1348 u16::from_be_bytes(fis)
1349 };
1350 pos += 2;
1351 let mut fs = [0u8; 2];
1352 fs.copy_from_slice(&body[pos..pos + 2]);
1353 let fragment_size = if little_endian {
1354 u16::from_le_bytes(fs)
1355 } else {
1356 u16::from_be_bytes(fs)
1357 };
1358 pos += 2;
1359 let mut ss = [0u8; 4];
1360 ss.copy_from_slice(&body[pos..pos + 4]);
1361 let sample_size = if little_endian {
1362 u32::from_le_bytes(ss)
1363 } else {
1364 u32::from_be_bytes(ss)
1365 };
1366 pos += 4;
1367 if inline_qos_flag {
1371 return Err(WireError::UnsupportedFeature {
1372 what: "DATA_FRAG with inline_qos",
1373 });
1374 }
1375 let serialized_payload: Arc<[u8]> = Arc::from(&body[pos..]);
1376 Ok(Self {
1377 extra_flags,
1378 reader_id,
1379 writer_id,
1380 writer_sn,
1381 fragment_starting_num,
1382 fragments_in_submessage,
1383 fragment_size,
1384 sample_size,
1385 serialized_payload,
1386 inline_qos_flag,
1387 hash_key_flag,
1388 key_flag,
1389 non_standard_flag,
1390 })
1391 }
1392}
1393
1394#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1410pub struct InfoSourceSubmessage {
1411 pub unused: u32,
1414 pub protocol_version: crate::wire_types::ProtocolVersion,
1416 pub vendor_id: crate::wire_types::VendorId,
1418 pub guid_prefix: crate::wire_types::GuidPrefix,
1420}
1421
1422impl InfoSourceSubmessage {
1423 pub const WIRE_SIZE: usize = 20;
1425
1426 #[must_use]
1428 pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
1429 let mut out = Vec::with_capacity(Self::WIRE_SIZE);
1430 out.extend_from_slice(&if little_endian {
1431 self.unused.to_le_bytes()
1432 } else {
1433 self.unused.to_be_bytes()
1434 });
1435 out.extend_from_slice(&self.protocol_version.to_bytes());
1436 out.extend_from_slice(&self.vendor_id.to_bytes());
1437 out.extend_from_slice(&self.guid_prefix.to_bytes());
1438 let mut flags = 0u8;
1439 if little_endian {
1440 flags |= FLAG_E_LITTLE_ENDIAN;
1441 }
1442 (out, flags)
1443 }
1444
1445 pub fn read_body(body: &[u8], little_endian: bool) -> Result<Self, WireError> {
1450 if body.len() < Self::WIRE_SIZE {
1451 return Err(WireError::UnexpectedEof {
1452 needed: Self::WIRE_SIZE,
1453 offset: 0,
1454 });
1455 }
1456 let mut pos = 0usize;
1457 let mut u = [0u8; 4];
1458 u.copy_from_slice(&body[pos..pos + 4]);
1459 let unused = if little_endian {
1460 u32::from_le_bytes(u)
1461 } else {
1462 u32::from_be_bytes(u)
1463 };
1464 pos += 4;
1465 let mut pv = [0u8; 2];
1466 pv.copy_from_slice(&body[pos..pos + 2]);
1467 let protocol_version = crate::wire_types::ProtocolVersion::from_bytes(pv);
1468 pos += 2;
1469 let mut vid = [0u8; 2];
1470 vid.copy_from_slice(&body[pos..pos + 2]);
1471 let vendor_id = crate::wire_types::VendorId::from_bytes(vid);
1472 pos += 2;
1473 let mut gp = [0u8; 12];
1474 gp.copy_from_slice(&body[pos..pos + 12]);
1475 let guid_prefix = crate::wire_types::GuidPrefix::from_bytes(gp);
1476 Ok(Self {
1477 unused,
1478 protocol_version,
1479 vendor_id,
1480 guid_prefix,
1481 })
1482 }
1483}
1484
1485pub const INFO_TIMESTAMP_FLAG_INVALIDATE: u8 = 0x02;
1492
1493#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1497pub struct InfoTimestampSubmessage {
1498 pub timestamp: crate::header_extension::HeTimestamp,
1501 pub invalidate: bool,
1504}
1505
1506impl InfoTimestampSubmessage {
1507 #[must_use]
1510 pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
1511 let mut flags = 0u8;
1512 if little_endian {
1513 flags |= FLAG_E_LITTLE_ENDIAN;
1514 }
1515 if self.invalidate {
1516 flags |= INFO_TIMESTAMP_FLAG_INVALIDATE;
1517 return (Vec::new(), flags);
1518 }
1519 let mut out = Vec::with_capacity(8);
1520 let s = if little_endian {
1521 self.timestamp.seconds.to_le_bytes()
1522 } else {
1523 self.timestamp.seconds.to_be_bytes()
1524 };
1525 let f = if little_endian {
1526 self.timestamp.fraction.to_le_bytes()
1527 } else {
1528 self.timestamp.fraction.to_be_bytes()
1529 };
1530 out.extend_from_slice(&s);
1531 out.extend_from_slice(&f);
1532 (out, flags)
1533 }
1534
1535 pub fn read_body(
1541 body: &[u8],
1542 little_endian: bool,
1543 invalidate_flag: bool,
1544 ) -> Result<Self, WireError> {
1545 if invalidate_flag {
1546 return Ok(Self {
1547 timestamp: crate::header_extension::HeTimestamp::default(),
1548 invalidate: true,
1549 });
1550 }
1551 if body.len() < 8 {
1552 return Err(WireError::UnexpectedEof {
1553 needed: 8,
1554 offset: 0,
1555 });
1556 }
1557 let mut s = [0u8; 4];
1558 s.copy_from_slice(&body[0..4]);
1559 let mut f = [0u8; 4];
1560 f.copy_from_slice(&body[4..8]);
1561 let seconds = if little_endian {
1562 i32::from_le_bytes(s)
1563 } else {
1564 i32::from_be_bytes(s)
1565 };
1566 let fraction = if little_endian {
1567 u32::from_le_bytes(f)
1568 } else {
1569 u32::from_be_bytes(f)
1570 };
1571 Ok(Self {
1572 timestamp: crate::header_extension::HeTimestamp { seconds, fraction },
1573 invalidate: false,
1574 })
1575 }
1576}
1577
1578pub const INFO_REPLY_FLAG_MULTICAST: u8 = 0x02;
1585
1586#[derive(Debug, Clone, PartialEq, Eq)]
1594pub struct InfoReplySubmessage {
1595 pub unicast_locators: Vec<crate::wire_types::Locator>,
1597 pub multicast_locators: Option<Vec<crate::wire_types::Locator>>,
1599}
1600
1601impl InfoReplySubmessage {
1602 #[must_use]
1604 pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
1605 let mut out = Vec::new();
1606 Self::write_locator_list(&mut out, &self.unicast_locators, little_endian);
1607 let mut flags = 0u8;
1608 if little_endian {
1609 flags |= FLAG_E_LITTLE_ENDIAN;
1610 }
1611 if let Some(mcast) = &self.multicast_locators {
1612 flags |= INFO_REPLY_FLAG_MULTICAST;
1613 Self::write_locator_list(&mut out, mcast, little_endian);
1614 }
1615 (out, flags)
1616 }
1617
1618 fn write_locator_list(
1619 out: &mut Vec<u8>,
1620 list: &[crate::wire_types::Locator],
1621 little_endian: bool,
1622 ) {
1623 let len = u32::try_from(list.len()).unwrap_or(u32::MAX);
1624 out.extend_from_slice(&if little_endian {
1625 len.to_le_bytes()
1626 } else {
1627 len.to_be_bytes()
1628 });
1629 for loc in list {
1630 if little_endian {
1635 out.extend_from_slice(&loc.to_bytes_le());
1636 } else {
1637 out.extend_from_slice(&(loc.kind.as_i32()).to_be_bytes());
1639 out.extend_from_slice(&loc.port.to_be_bytes());
1640 out.extend_from_slice(&loc.address);
1641 }
1642 }
1643 }
1644
1645 pub fn read_body(
1651 body: &[u8],
1652 little_endian: bool,
1653 multicast_flag: bool,
1654 ) -> Result<Self, WireError> {
1655 let mut pos = 0usize;
1656 let unicast_locators = Self::read_locator_list(body, &mut pos, little_endian)?;
1657 let multicast_locators = if multicast_flag {
1658 Some(Self::read_locator_list(body, &mut pos, little_endian)?)
1659 } else {
1660 None
1661 };
1662 Ok(Self {
1663 unicast_locators,
1664 multicast_locators,
1665 })
1666 }
1667
1668 fn read_locator_list(
1669 body: &[u8],
1670 pos: &mut usize,
1671 little_endian: bool,
1672 ) -> Result<Vec<crate::wire_types::Locator>, WireError> {
1673 if body.len() < *pos + 4 {
1674 return Err(WireError::UnexpectedEof {
1675 needed: 4,
1676 offset: *pos,
1677 });
1678 }
1679 let mut len_bytes = [0u8; 4];
1680 len_bytes.copy_from_slice(&body[*pos..*pos + 4]);
1681 let len = if little_endian {
1682 u32::from_le_bytes(len_bytes)
1683 } else {
1684 u32::from_be_bytes(len_bytes)
1685 } as usize;
1686 *pos += 4;
1687 let remaining = body.len().saturating_sub(*pos);
1688 if len.saturating_mul(24) > remaining {
1689 return Err(WireError::ValueOutOfRange {
1690 message: "InfoReply.locatorList length exceeds body",
1691 });
1692 }
1693 let mut out = Vec::with_capacity(len);
1694 for _ in 0..len {
1695 let mut buf = [0u8; 24];
1696 buf.copy_from_slice(&body[*pos..*pos + 24]);
1697 let loc = if little_endian {
1700 crate::wire_types::Locator::from_bytes_le(buf)?
1701 } else {
1702 let mut k = [0u8; 4];
1703 k.copy_from_slice(&buf[0..4]);
1704 let kind_raw = i32::from_be_bytes(k);
1705 let kind = crate::wire_types::LocatorKind::from_i32(kind_raw)?;
1706 let mut p = [0u8; 4];
1707 p.copy_from_slice(&buf[4..8]);
1708 let port = u32::from_be_bytes(p);
1709 let mut address = [0u8; 16];
1710 address.copy_from_slice(&buf[8..24]);
1711 crate::wire_types::Locator {
1712 kind,
1713 port,
1714 address,
1715 }
1716 };
1717 out.push(loc);
1718 *pos += 24;
1719 }
1720 Ok(out)
1721 }
1722}
1723
1724#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1733pub struct HeartbeatFragSubmessage {
1734 pub reader_id: EntityId,
1736 pub writer_id: EntityId,
1738 pub writer_sn: SequenceNumber,
1740 pub last_fragment_num: FragmentNumber,
1742 pub count: i32,
1744}
1745
1746impl HeartbeatFragSubmessage {
1747 pub const WIRE_SIZE: usize = 24;
1749
1750 #[must_use]
1752 pub fn write_body(self, little_endian: bool) -> (Vec<u8>, u8) {
1753 let mut out = Vec::with_capacity(Self::WIRE_SIZE);
1754 out.extend_from_slice(&self.reader_id.to_bytes());
1755 out.extend_from_slice(&self.writer_id.to_bytes());
1756 out.extend_from_slice(&if little_endian {
1757 self.writer_sn.to_bytes_le()
1758 } else {
1759 self.writer_sn.to_bytes_be()
1760 });
1761 out.extend_from_slice(&if little_endian {
1762 self.last_fragment_num.to_bytes_le()
1763 } else {
1764 self.last_fragment_num.to_bytes_be()
1765 });
1766 out.extend_from_slice(&if little_endian {
1767 self.count.to_le_bytes()
1768 } else {
1769 self.count.to_be_bytes()
1770 });
1771 let mut flags = 0u8;
1772 if little_endian {
1773 flags |= FLAG_E_LITTLE_ENDIAN;
1774 }
1775 (out, flags)
1776 }
1777
1778 pub fn read_body(body: &[u8], little_endian: bool) -> Result<Self, WireError> {
1783 if body.len() < Self::WIRE_SIZE {
1784 return Err(WireError::UnexpectedEof {
1785 needed: Self::WIRE_SIZE,
1786 offset: 0,
1787 });
1788 }
1789 let mut pos = 0usize;
1790 let mut rid = [0u8; 4];
1791 rid.copy_from_slice(&body[pos..pos + 4]);
1792 let reader_id = EntityId::from_bytes(rid);
1793 pos += 4;
1794 let mut wid = [0u8; 4];
1795 wid.copy_from_slice(&body[pos..pos + 4]);
1796 let writer_id = EntityId::from_bytes(wid);
1797 pos += 4;
1798 let mut sn = [0u8; 8];
1799 sn.copy_from_slice(&body[pos..pos + 8]);
1800 let writer_sn = if little_endian {
1801 SequenceNumber::from_bytes_le(sn)
1802 } else {
1803 SequenceNumber::from_bytes_be(sn)
1804 };
1805 pos += 8;
1806 let mut lf = [0u8; 4];
1807 lf.copy_from_slice(&body[pos..pos + 4]);
1808 let last_fragment_num = if little_endian {
1809 FragmentNumber::from_bytes_le(lf)
1810 } else {
1811 FragmentNumber::from_bytes_be(lf)
1812 };
1813 pos += 4;
1814 let mut cnt = [0u8; 4];
1815 cnt.copy_from_slice(&body[pos..pos + 4]);
1816 let count = if little_endian {
1817 i32::from_le_bytes(cnt)
1818 } else {
1819 i32::from_be_bytes(cnt)
1820 };
1821 Ok(Self {
1822 reader_id,
1823 writer_id,
1824 writer_sn,
1825 last_fragment_num,
1826 count,
1827 })
1828 }
1829}
1830
1831#[derive(Debug, Clone, PartialEq, Eq)]
1838pub struct NackFragSubmessage {
1839 pub reader_id: EntityId,
1841 pub writer_id: EntityId,
1843 pub writer_sn: SequenceNumber,
1845 pub fragment_number_state: FragmentNumberSet,
1847 pub count: i32,
1849}
1850
1851impl NackFragSubmessage {
1852 #[must_use]
1854 pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
1855 let mut out = Vec::new();
1856 out.extend_from_slice(&self.reader_id.to_bytes());
1857 out.extend_from_slice(&self.writer_id.to_bytes());
1858 out.extend_from_slice(&if little_endian {
1859 self.writer_sn.to_bytes_le()
1860 } else {
1861 self.writer_sn.to_bytes_be()
1862 });
1863 self.fragment_number_state.write_to(&mut out, little_endian);
1864 out.extend_from_slice(&if little_endian {
1865 self.count.to_le_bytes()
1866 } else {
1867 self.count.to_be_bytes()
1868 });
1869 let mut flags = 0u8;
1870 if little_endian {
1871 flags |= FLAG_E_LITTLE_ENDIAN;
1872 }
1873 (out, flags)
1874 }
1875
1876 pub fn read_body(body: &[u8], little_endian: bool) -> Result<Self, WireError> {
1881 if body.len() < 4 + 4 + 8 + 4 + 4 + 4 {
1882 return Err(WireError::UnexpectedEof {
1883 needed: 4 + 4 + 8 + 4 + 4 + 4,
1884 offset: 0,
1885 });
1886 }
1887 let mut pos = 0usize;
1888 let mut rid = [0u8; 4];
1889 rid.copy_from_slice(&body[pos..pos + 4]);
1890 let reader_id = EntityId::from_bytes(rid);
1891 pos += 4;
1892 let mut wid = [0u8; 4];
1893 wid.copy_from_slice(&body[pos..pos + 4]);
1894 let writer_id = EntityId::from_bytes(wid);
1895 pos += 4;
1896 let mut sn = [0u8; 8];
1897 sn.copy_from_slice(&body[pos..pos + 8]);
1898 let writer_sn = if little_endian {
1899 SequenceNumber::from_bytes_le(sn)
1900 } else {
1901 SequenceNumber::from_bytes_be(sn)
1902 };
1903 pos += 8;
1904 let (fragment_number_state, new_pos) =
1905 FragmentNumberSet::read_from(body, pos, little_endian)?;
1906 pos = new_pos;
1907 if body.len() < pos + 4 {
1908 return Err(WireError::UnexpectedEof {
1909 needed: 4,
1910 offset: pos,
1911 });
1912 }
1913 let mut cnt = [0u8; 4];
1914 cnt.copy_from_slice(&body[pos..pos + 4]);
1915 let count = if little_endian {
1916 i32::from_le_bytes(cnt)
1917 } else {
1918 i32::from_be_bytes(cnt)
1919 };
1920 Ok(Self {
1921 reader_id,
1922 writer_id,
1923 writer_sn,
1924 fragment_number_state,
1925 count,
1926 })
1927 }
1928}
1929
1930#[cfg(test)]
1931mod tests {
1932 #![allow(clippy::expect_used, clippy::panic, clippy::unwrap_used)]
1933 use super::*;
1934 use alloc::vec;
1935
1936 fn writer_id() -> EntityId {
1937 EntityId::user_writer_with_key([0x10, 0x20, 0x30])
1938 }
1939 fn reader_id() -> EntityId {
1940 EntityId::user_reader_with_key([0x40, 0x50, 0x60])
1941 }
1942
1943 #[test]
1946 fn snset_wire_size_zero_bits_is_12_bytes() {
1947 assert_eq!(SequenceNumberSet::wire_size(0), 12);
1948 }
1949
1950 #[test]
1951 fn snset_wire_size_32_bits_is_16_bytes() {
1952 assert_eq!(SequenceNumberSet::wire_size(32), 16);
1953 }
1954
1955 #[test]
1956 fn snset_wire_size_33_bits_is_20_bytes() {
1957 assert_eq!(SequenceNumberSet::wire_size(33), 20);
1958 }
1959
1960 #[test]
1961 fn snset_roundtrip_le() {
1962 let s = SequenceNumberSet {
1963 bitmap_base: SequenceNumber(100),
1964 num_bits: 5,
1965 bitmap: vec![0b0000_1010_0000_0000_0000_0000_0000_0000],
1966 };
1967 let mut buf = Vec::new();
1968 s.write_to(&mut buf, true);
1969 let (decoded, end) = SequenceNumberSet::read_from(&buf, 0, true).unwrap();
1970 assert_eq!(decoded, s);
1971 assert_eq!(end, buf.len());
1972 }
1973
1974 #[test]
1975 fn snset_roundtrip_be() {
1976 let s = SequenceNumberSet {
1977 bitmap_base: SequenceNumber(0xDEAD_BEEF),
1978 num_bits: 64,
1979 bitmap: vec![0x1234_5678, 0x9ABC_DEF0],
1980 };
1981 let mut buf = Vec::new();
1982 s.write_to(&mut buf, false);
1983 let (decoded, _) = SequenceNumberSet::read_from(&buf, 0, false).unwrap();
1984 assert_eq!(decoded, s);
1985 }
1986
1987 #[test]
1988 fn snset_decode_rejects_truncated_bitmap() {
1989 let mut buf = Vec::new();
1991 buf.extend_from_slice(&SequenceNumber(0).to_bytes_le());
1992 buf.extend_from_slice(&64_u32.to_le_bytes());
1993 buf.extend_from_slice(&[0u8; 4]); let res = SequenceNumberSet::read_from(&buf, 0, true);
1995 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
1996 }
1997
1998 #[test]
2001 fn data_submessage_roundtrip_le() {
2002 let d = DataSubmessage {
2003 extra_flags: 0,
2004 reader_id: reader_id(),
2005 writer_id: writer_id(),
2006 writer_sn: SequenceNumber(42),
2007 inline_qos: None,
2008 key_flag: false,
2009 non_standard_flag: false,
2010 serialized_payload: Arc::<[u8]>::from([1u8, 2, 3, 4, 5].as_slice()),
2011 };
2012 let (bytes, flags) = d.write_body(true);
2013 assert!(flags & FLAG_E_LITTLE_ENDIAN != 0);
2014 assert!(flags & DATA_FLAG_DATA != 0);
2015 let decoded = DataSubmessage::read_body(&bytes, true).unwrap();
2016 assert_eq!(decoded, d);
2017 }
2018
2019 #[test]
2020 fn data_submessage_roundtrip_be_with_empty_payload() {
2021 let d = DataSubmessage {
2022 extra_flags: 0,
2023 reader_id: reader_id(),
2024 writer_id: writer_id(),
2025 writer_sn: SequenceNumber(0xDEAD_BEEF),
2026 inline_qos: None,
2027 key_flag: false,
2028 non_standard_flag: false,
2029 serialized_payload: Arc::<[u8]>::from([].as_slice()),
2030 };
2031 let (bytes, flags) = d.write_body(false);
2032 assert_eq!(flags & FLAG_E_LITTLE_ENDIAN, 0);
2033 let decoded = DataSubmessage::read_body(&bytes, false).unwrap();
2034 assert_eq!(decoded, d);
2035 }
2036
2037 #[test]
2038 fn data_submessage_key_flag_roundtrip() {
2039 let d = DataSubmessage {
2041 extra_flags: 0,
2042 reader_id: reader_id(),
2043 writer_id: writer_id(),
2044 writer_sn: SequenceNumber(7),
2045 inline_qos: None,
2046 key_flag: true,
2047 non_standard_flag: false,
2048 serialized_payload: Arc::<[u8]>::from([0xAA, 0xBB].as_slice()),
2049 };
2050 let (bytes, flags) = d.write_body(true);
2051 assert!(flags & DATA_FLAG_KEY != 0, "K-Flag must be set");
2052 let decoded = DataSubmessage::read_body_with_flags(&bytes, true, flags).unwrap();
2053 assert!(decoded.key_flag);
2054 assert!(!decoded.non_standard_flag);
2055 assert_eq!(decoded, d);
2056 }
2057
2058 #[test]
2059 fn data_submessage_non_standard_flag_roundtrip() {
2060 let d = DataSubmessage {
2062 extra_flags: 0,
2063 reader_id: reader_id(),
2064 writer_id: writer_id(),
2065 writer_sn: SequenceNumber(8),
2066 inline_qos: None,
2067 key_flag: false,
2068 non_standard_flag: true,
2069 serialized_payload: Arc::<[u8]>::from([0xCC, 0xDD].as_slice()),
2070 };
2071 let (bytes, flags) = d.write_body(true);
2072 assert!(flags & DATA_FLAG_NON_STANDARD != 0, "N-Flag must be set");
2073 let decoded = DataSubmessage::read_body_with_flags(&bytes, true, flags).unwrap();
2074 assert!(!decoded.key_flag);
2075 assert!(decoded.non_standard_flag);
2076 assert_eq!(decoded, d);
2077 }
2078
2079 #[test]
2080 fn data_submessage_all_flags_combined_roundtrip() {
2081 let mut pl = crate::parameter_list::ParameterList::new();
2083 pl.push(crate::parameter_list::Parameter::new(0x0070, vec![1; 4]));
2084 let d = DataSubmessage {
2085 extra_flags: 0xABCD,
2086 reader_id: reader_id(),
2087 writer_id: writer_id(),
2088 writer_sn: SequenceNumber(9),
2089 inline_qos: Some(pl),
2090 key_flag: true,
2091 non_standard_flag: true,
2092 serialized_payload: Arc::<[u8]>::from([0xEE; 8].as_slice()),
2093 };
2094 let (bytes, flags) = d.write_body(true);
2095 assert!(flags & FLAG_E_LITTLE_ENDIAN != 0);
2096 assert!(flags & DATA_FLAG_INLINE_QOS != 0);
2097 assert!(flags & DATA_FLAG_DATA != 0);
2098 assert!(flags & DATA_FLAG_KEY != 0);
2099 assert!(flags & DATA_FLAG_NON_STANDARD != 0);
2100 let decoded = DataSubmessage::read_body_with_flags(&bytes, true, flags).unwrap();
2101 assert_eq!(decoded, d);
2102 }
2103
2104 #[test]
2105 fn data_submessage_octets_to_inline_qos_is_16() {
2106 let d = DataSubmessage {
2107 extra_flags: 0,
2108 reader_id: reader_id(),
2109 writer_id: writer_id(),
2110 writer_sn: SequenceNumber(1),
2111 inline_qos: None,
2112 key_flag: false,
2113 non_standard_flag: false,
2114 serialized_payload: Arc::<[u8]>::from([].as_slice()),
2115 };
2116 let (bytes, _) = d.write_body(true);
2117 assert_eq!(&bytes[2..4], &[16, 0]);
2119 }
2120
2121 #[test]
2122 fn data_submessage_decode_rejects_truncated() {
2123 let res = DataSubmessage::read_body(&[1, 2, 3], true);
2124 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2125 }
2126
2127 #[test]
2130 fn heartbeat_submessage_roundtrip_le() {
2131 let h = HeartbeatSubmessage {
2132 reader_id: reader_id(),
2133 writer_id: writer_id(),
2134 first_sn: SequenceNumber(1),
2135 last_sn: SequenceNumber(10),
2136 count: 7,
2137 final_flag: true,
2138 liveliness_flag: false,
2139 group_info: None,
2140 };
2141 let (bytes, flags) = h.write_body(true);
2142 assert!(flags & HEARTBEAT_FLAG_FINAL != 0);
2143 assert_eq!(flags & HEARTBEAT_FLAG_LIVELINESS, 0);
2144 assert_eq!(bytes.len(), HeartbeatSubmessage::WIRE_SIZE);
2145 let decoded = HeartbeatSubmessage::read_body(&bytes, true, true, false, false).unwrap();
2146 assert_eq!(decoded, h);
2147 }
2148
2149 #[test]
2150 fn heartbeat_submessage_no_final_flag_when_disabled() {
2151 let h = HeartbeatSubmessage {
2152 reader_id: reader_id(),
2153 writer_id: writer_id(),
2154 first_sn: SequenceNumber(1),
2155 last_sn: SequenceNumber(1),
2156 count: 0,
2157 final_flag: false,
2158 liveliness_flag: false,
2159 group_info: None,
2160 };
2161 let (_, flags) = h.write_body(true);
2162 assert_eq!(flags & HEARTBEAT_FLAG_FINAL, 0);
2163 }
2164
2165 #[test]
2166 fn heartbeat_submessage_liveliness_flag_roundtrip() {
2167 let h = HeartbeatSubmessage {
2168 reader_id: reader_id(),
2169 writer_id: writer_id(),
2170 first_sn: SequenceNumber(1),
2171 last_sn: SequenceNumber(1),
2172 count: 0,
2173 final_flag: false,
2174 liveliness_flag: true,
2175 group_info: None,
2176 };
2177 let (bytes, flags) = h.write_body(true);
2178 assert!(flags & HEARTBEAT_FLAG_LIVELINESS != 0);
2179 let decoded = HeartbeatSubmessage::read_body(&bytes, true, false, true, false).unwrap();
2180 assert_eq!(decoded, h);
2181 assert!(decoded.liveliness_flag);
2182 }
2183
2184 #[test]
2185 fn heartbeat_decode_rejects_truncated() {
2186 let res = HeartbeatSubmessage::read_body(&[0u8; 27], true, false, false, false);
2187 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2188 }
2189
2190 #[test]
2193 fn heartbeat_with_empty_group_info_roundtrip_le() {
2194 let h = HeartbeatSubmessage {
2195 reader_id: reader_id(),
2196 writer_id: writer_id(),
2197 first_sn: SequenceNumber(1),
2198 last_sn: SequenceNumber(5),
2199 count: 3,
2200 final_flag: false,
2201 liveliness_flag: false,
2202 group_info: Some(HeartbeatGroupInfo {
2203 current_gsn: SequenceNumber(100),
2204 first_gsn: SequenceNumber(50),
2205 last_gsn: SequenceNumber(99),
2206 writer_set: vec![],
2207 }),
2208 };
2209 let (bytes, flags) = h.write_body(true);
2210 assert!(flags & HEARTBEAT_FLAG_GROUP_INFO != 0);
2211 let decoded = HeartbeatSubmessage::read_body(&bytes, true, false, false, true).unwrap();
2212 assert_eq!(decoded, h);
2213 }
2214
2215 #[test]
2216 fn heartbeat_with_writer_set_roundtrip_be() {
2217 use crate::wire_types::GuidPrefix;
2218 let h = HeartbeatSubmessage {
2219 reader_id: reader_id(),
2220 writer_id: writer_id(),
2221 first_sn: SequenceNumber(1),
2222 last_sn: SequenceNumber(2),
2223 count: 1,
2224 final_flag: false,
2225 liveliness_flag: false,
2226 group_info: Some(HeartbeatGroupInfo {
2227 current_gsn: SequenceNumber(7),
2228 first_gsn: SequenceNumber(1),
2229 last_gsn: SequenceNumber(7),
2230 writer_set: vec![
2231 GuidPrefix::from_bytes([1; 12]),
2232 GuidPrefix::from_bytes([2; 12]),
2233 GuidPrefix::from_bytes([3; 12]),
2234 ],
2235 }),
2236 };
2237 let (bytes, flags) = h.write_body(false);
2238 assert!(flags & HEARTBEAT_FLAG_GROUP_INFO != 0);
2239 let decoded = HeartbeatSubmessage::read_body(&bytes, false, false, false, true).unwrap();
2240 assert_eq!(decoded, h);
2241 let gi = decoded.group_info.unwrap();
2242 assert_eq!(gi.writer_set.len(), 3);
2243 }
2244
2245 #[test]
2246 fn heartbeat_decode_rejects_oversized_writer_set_length() {
2247 let mut body = Vec::new();
2249 body.extend_from_slice(&reader_id().to_bytes());
2250 body.extend_from_slice(&writer_id().to_bytes());
2251 body.extend_from_slice(&SequenceNumber(1).to_bytes_le());
2252 body.extend_from_slice(&SequenceNumber(1).to_bytes_le());
2253 body.extend_from_slice(&1i32.to_le_bytes());
2254 body.extend_from_slice(&SequenceNumber(0).to_bytes_le());
2256 body.extend_from_slice(&SequenceNumber(0).to_bytes_le());
2257 body.extend_from_slice(&SequenceNumber(0).to_bytes_le());
2258 body.extend_from_slice(&u32::MAX.to_le_bytes());
2260 let res = HeartbeatSubmessage::read_body(&body, true, false, false, true);
2262 assert!(matches!(res, Err(WireError::ValueOutOfRange { .. })));
2263 }
2264
2265 #[test]
2266 fn heartbeat_decode_rejects_truncated_group_info() {
2267 let mut body = Vec::new();
2269 body.extend_from_slice(&reader_id().to_bytes());
2270 body.extend_from_slice(&writer_id().to_bytes());
2271 body.extend_from_slice(&SequenceNumber(1).to_bytes_le());
2272 body.extend_from_slice(&SequenceNumber(1).to_bytes_le());
2273 body.extend_from_slice(&1i32.to_le_bytes());
2274 let res = HeartbeatSubmessage::read_body(&body, true, false, false, true);
2276 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2277 }
2278
2279 #[test]
2282 fn acknack_submessage_roundtrip_le() {
2283 let a = AckNackSubmessage {
2284 reader_id: reader_id(),
2285 writer_id: writer_id(),
2286 reader_sn_state: SequenceNumberSet {
2287 bitmap_base: SequenceNumber(5),
2288 num_bits: 3,
2289 bitmap: vec![0b1010_0000_0000_0000_0000_0000_0000_0000],
2290 },
2291 count: 1,
2292 final_flag: false,
2293 };
2294 let (bytes, flags) = a.write_body(true);
2295 assert_eq!(flags & ACKNACK_FLAG_FINAL, 0);
2296 let decoded = AckNackSubmessage::read_body(&bytes, true, false).unwrap();
2297 assert_eq!(decoded, a);
2298 }
2299
2300 #[test]
2301 fn acknack_submessage_with_final_flag() {
2302 let a = AckNackSubmessage {
2303 reader_id: reader_id(),
2304 writer_id: writer_id(),
2305 reader_sn_state: SequenceNumberSet {
2306 bitmap_base: SequenceNumber(1),
2307 num_bits: 0,
2308 bitmap: vec![],
2309 },
2310 count: 0,
2311 final_flag: true,
2312 };
2313 let (bytes, flags) = a.write_body(true);
2314 assert!(flags & ACKNACK_FLAG_FINAL != 0);
2315 let decoded = AckNackSubmessage::read_body(&bytes, true, true).unwrap();
2316 assert!(decoded.final_flag);
2317 }
2318
2319 #[test]
2322 fn gap_submessage_roundtrip_le() {
2323 let g = GapSubmessage {
2324 reader_id: reader_id(),
2325 writer_id: writer_id(),
2326 gap_start: SequenceNumber(1),
2327 gap_list: SequenceNumberSet {
2328 bitmap_base: SequenceNumber(5),
2329 num_bits: 8,
2330 bitmap: vec![0xFF000000],
2331 },
2332 group_info: None,
2333 filtered_count: None,
2334 };
2335 let (bytes, flags) = g.write_body(true);
2336 assert!(flags & FLAG_E_LITTLE_ENDIAN != 0);
2337 assert_eq!(flags & GAP_FLAG_GROUP_INFO, 0);
2338 assert_eq!(flags & GAP_FLAG_FILTERED_COUNT, 0);
2339 let decoded = GapSubmessage::read_body(&bytes, true, false, false).unwrap();
2340 assert_eq!(decoded, g);
2341 }
2342
2343 #[test]
2344 fn gap_decode_rejects_truncated_header() {
2345 let res = GapSubmessage::read_body(&[0u8; 10], true, false, false);
2346 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2347 }
2348
2349 #[test]
2352 fn gap_with_filtered_count_roundtrip_le() {
2353 let g = GapSubmessage {
2354 reader_id: reader_id(),
2355 writer_id: writer_id(),
2356 gap_start: SequenceNumber(1),
2357 gap_list: SequenceNumberSet {
2358 bitmap_base: SequenceNumber(2),
2359 num_bits: 0,
2360 bitmap: vec![],
2361 },
2362 group_info: None,
2363 filtered_count: Some(3),
2364 };
2365 let (bytes, flags) = g.write_body(true);
2366 assert!(flags & GAP_FLAG_FILTERED_COUNT != 0);
2367 let decoded = GapSubmessage::read_body(&bytes, true, false, true).unwrap();
2368 assert_eq!(decoded, g);
2369 assert_eq!(decoded.filtered_count, Some(3));
2370 }
2371
2372 #[test]
2373 fn gap_with_group_info_roundtrip_be() {
2374 let g = GapSubmessage {
2375 reader_id: reader_id(),
2376 writer_id: writer_id(),
2377 gap_start: SequenceNumber(10),
2378 gap_list: SequenceNumberSet {
2379 bitmap_base: SequenceNumber(11),
2380 num_bits: 0,
2381 bitmap: vec![],
2382 },
2383 group_info: Some(GapGroupInfo {
2384 gap_start_gsn: SequenceNumber(100),
2385 gap_end_gsn: SequenceNumber(110),
2386 }),
2387 filtered_count: None,
2388 };
2389 let (bytes, flags) = g.write_body(false);
2390 assert!(flags & GAP_FLAG_GROUP_INFO != 0);
2391 let decoded = GapSubmessage::read_body(&bytes, false, true, false).unwrap();
2392 assert_eq!(decoded, g);
2393 }
2394
2395 #[test]
2396 fn gap_with_group_info_and_filtered_count_combined() {
2397 let g = GapSubmessage {
2398 reader_id: reader_id(),
2399 writer_id: writer_id(),
2400 gap_start: SequenceNumber(5),
2401 gap_list: SequenceNumberSet {
2402 bitmap_base: SequenceNumber(6),
2403 num_bits: 0,
2404 bitmap: vec![],
2405 },
2406 group_info: Some(GapGroupInfo {
2407 gap_start_gsn: SequenceNumber(50),
2408 gap_end_gsn: SequenceNumber(55),
2409 }),
2410 filtered_count: Some(7),
2411 };
2412 let (bytes, flags) = g.write_body(true);
2413 assert!(flags & GAP_FLAG_GROUP_INFO != 0);
2414 assert!(flags & GAP_FLAG_FILTERED_COUNT != 0);
2415 let decoded = GapSubmessage::read_body(&bytes, true, true, true).unwrap();
2416 assert_eq!(decoded, g);
2417 }
2418
2419 #[test]
2420 fn gap_filtered_count_zero_is_distinct_from_none() {
2421 let zero = GapSubmessage {
2425 reader_id: reader_id(),
2426 writer_id: writer_id(),
2427 gap_start: SequenceNumber(1),
2428 gap_list: SequenceNumberSet {
2429 bitmap_base: SequenceNumber(2),
2430 num_bits: 0,
2431 bitmap: vec![],
2432 },
2433 group_info: None,
2434 filtered_count: Some(0),
2435 };
2436 let (bytes, flags) = zero.write_body(true);
2437 assert!(flags & GAP_FLAG_FILTERED_COUNT != 0);
2438 let decoded = GapSubmessage::read_body(&bytes, true, false, true).unwrap();
2439 assert_eq!(decoded.filtered_count, Some(0));
2440 }
2441
2442 #[test]
2443 fn gap_decode_rejects_truncated_filtered_count() {
2444 let g = GapSubmessage {
2446 reader_id: reader_id(),
2447 writer_id: writer_id(),
2448 gap_start: SequenceNumber(1),
2449 gap_list: SequenceNumberSet {
2450 bitmap_base: SequenceNumber(2),
2451 num_bits: 0,
2452 bitmap: vec![],
2453 },
2454 group_info: None,
2455 filtered_count: None,
2456 };
2457 let (bytes, _) = g.write_body(true);
2458 let res = GapSubmessage::read_body(&bytes, true, false, true);
2460 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2461 }
2462
2463 #[test]
2464 fn gap_decode_rejects_truncated_group_info() {
2465 let g = GapSubmessage {
2466 reader_id: reader_id(),
2467 writer_id: writer_id(),
2468 gap_start: SequenceNumber(1),
2469 gap_list: SequenceNumberSet {
2470 bitmap_base: SequenceNumber(2),
2471 num_bits: 0,
2472 bitmap: vec![],
2473 },
2474 group_info: None,
2475 filtered_count: None,
2476 };
2477 let (bytes, _) = g.write_body(true);
2478 let res = GapSubmessage::read_body(&bytes, true, true, false);
2479 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2480 }
2481
2482 #[test]
2485 fn fnset_wire_size_formula() {
2486 assert_eq!(FragmentNumberSet::wire_size(0), 8);
2487 assert_eq!(FragmentNumberSet::wire_size(1), 12);
2488 assert_eq!(FragmentNumberSet::wire_size(32), 12);
2489 assert_eq!(FragmentNumberSet::wire_size(33), 16);
2490 }
2491
2492 #[test]
2493 fn fnset_from_missing_single() {
2494 let s = FragmentNumberSet::from_missing(
2495 FragmentNumber(1),
2496 &[FragmentNumber(1), FragmentNumber(3)],
2497 );
2498 assert_eq!(s.bitmap_base, FragmentNumber(1));
2499 assert_eq!(s.num_bits, 3);
2500 let set: Vec<_> = s.iter_set().collect();
2501 assert_eq!(set, vec![FragmentNumber(1), FragmentNumber(3)]);
2502 }
2503
2504 #[test]
2505 fn fnset_from_missing_empty() {
2506 let s = FragmentNumberSet::from_missing(FragmentNumber(5), &[]);
2507 assert_eq!(s.num_bits, 0);
2508 assert!(s.iter_set().next().is_none());
2509 }
2510
2511 #[test]
2512 fn fnset_from_missing_caps_num_bits_at_256() {
2513 let missing = [FragmentNumber(1), FragmentNumber(300)];
2521 let s = FragmentNumberSet::from_missing(FragmentNumber(1), &missing);
2522 assert!(
2523 s.num_bits <= 256,
2524 "num_bits {} > 256 (malformed)",
2525 s.num_bits
2526 );
2527 assert_eq!(s.bitmap_base, FragmentNumber(1));
2528 let set: Vec<_> = s.iter_set().collect();
2531 assert!(set.contains(&FragmentNumber(1)));
2532 assert!(!set.contains(&FragmentNumber(300)));
2533 }
2534
2535 #[test]
2536 fn fnset_missing_below_base_is_ignored() {
2537 let s = FragmentNumberSet::from_missing(
2538 FragmentNumber(10),
2539 &[FragmentNumber(5), FragmentNumber(11)],
2540 );
2541 assert_eq!(s.bitmap_base, FragmentNumber(10));
2542 let set: Vec<_> = s.iter_set().collect();
2543 assert_eq!(set, vec![FragmentNumber(11)]);
2544 }
2545
2546 #[test]
2547 fn fnset_roundtrip_le() {
2548 let s = FragmentNumberSet {
2549 bitmap_base: FragmentNumber(100),
2550 num_bits: 35,
2551 bitmap: vec![0xDEAD_BEEF, 0xC000_0000],
2552 };
2553 let mut buf = Vec::new();
2554 s.write_to(&mut buf, true);
2555 assert_eq!(buf.len(), s.encoded_size());
2556 let (decoded, end) = FragmentNumberSet::read_from(&buf, 0, true).unwrap();
2557 assert_eq!(decoded, s);
2558 assert_eq!(end, buf.len());
2559 }
2560
2561 #[test]
2562 fn fnset_roundtrip_be() {
2563 let s = FragmentNumberSet {
2564 bitmap_base: FragmentNumber(1),
2565 num_bits: 8,
2566 bitmap: vec![0xFF00_0000],
2567 };
2568 let mut buf = Vec::new();
2569 s.write_to(&mut buf, false);
2570 let (decoded, _) = FragmentNumberSet::read_from(&buf, 0, false).unwrap();
2571 assert_eq!(decoded, s);
2572 }
2573
2574 #[test]
2575 fn fnset_decode_rejects_truncated() {
2576 let buf = [0u8; 4];
2577 let res = FragmentNumberSet::read_from(&buf, 0, true);
2578 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2579 }
2580
2581 fn dataflag_frag(
2584 writer_sn: i64,
2585 starting: u32,
2586 count: u16,
2587 frag_size: u16,
2588 sample_size: u32,
2589 payload: Vec<u8>,
2590 ) -> DataFragSubmessage {
2591 DataFragSubmessage {
2592 extra_flags: 0,
2593 reader_id: reader_id(),
2594 writer_id: writer_id(),
2595 writer_sn: SequenceNumber(writer_sn),
2596 fragment_starting_num: FragmentNumber(starting),
2597 fragments_in_submessage: count,
2598 fragment_size: frag_size,
2599 sample_size,
2600 serialized_payload: Arc::from(payload),
2601 inline_qos_flag: false,
2602 hash_key_flag: false,
2603 key_flag: false,
2604 non_standard_flag: false,
2605 }
2606 }
2607
2608 #[test]
2609 fn data_frag_roundtrip_le() {
2610 let d = dataflag_frag(1, 1, 1, 4, 12, vec![0xDE, 0xAD, 0xBE, 0xEF]);
2611 let (bytes, flags) = d.write_body(true);
2612 assert!(flags & FLAG_E_LITTLE_ENDIAN != 0);
2613 assert_eq!(bytes.len(), DataFragSubmessage::HEADER_WIRE_SIZE + 4);
2614 let decoded =
2615 DataFragSubmessage::read_body(&bytes, true, false, false, false, false).unwrap();
2616 assert_eq!(decoded, d);
2617 }
2618
2619 #[test]
2620 fn data_frag_roundtrip_be() {
2621 let d = dataflag_frag(7, 2, 1, 8, 16, vec![1, 2, 3, 4, 5, 6, 7, 8]);
2622 let (bytes, flags) = d.write_body(false);
2623 assert_eq!(flags & FLAG_E_LITTLE_ENDIAN, 0);
2624 let decoded =
2625 DataFragSubmessage::read_body(&bytes, false, false, false, false, false).unwrap();
2626 assert_eq!(decoded, d);
2627 }
2628
2629 #[test]
2630 fn data_frag_last_fragment_shorter_than_fragment_size() {
2631 let d = dataflag_frag(1, 3, 1, 4, 10, vec![0xAA, 0xBB]);
2633 let (bytes, _) = d.write_body(true);
2634 let decoded =
2635 DataFragSubmessage::read_body(&bytes, true, false, false, false, false).unwrap();
2636 assert_eq!(decoded.serialized_payload.as_ref(), &[0xAA, 0xBB][..]);
2637 assert_eq!(decoded.sample_size, 10);
2638 assert_eq!(decoded.fragment_size, 4);
2639 }
2640
2641 #[test]
2642 fn data_frag_decode_rejects_truncated() {
2643 let res = DataFragSubmessage::read_body(&[0u8; 20], true, false, false, false, false);
2644 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2645 }
2646
2647 #[test]
2648 fn data_frag_decode_accepts_nonzero_extra_flags_silently() {
2649 let d = dataflag_frag(1, 1, 1, 4, 4, vec![1, 2, 3, 4]);
2652 let (mut bytes, _) = d.write_body(true);
2653 bytes[0..2].copy_from_slice(&0x0042u16.to_le_bytes()); let decoded =
2655 DataFragSubmessage::read_body(&bytes, true, false, false, false, false).unwrap();
2656 assert_eq!(decoded.extra_flags, 0x0042);
2657 }
2658
2659 #[test]
2660 fn seqnumset_rejects_num_bits_above_256() {
2661 let mut buf = Vec::new();
2663 buf.extend_from_slice(&SequenceNumber(1).to_bytes_le());
2664 buf.extend_from_slice(&257u32.to_le_bytes()); let res = SequenceNumberSet::read_from(&buf, 0, true);
2666 assert!(matches!(res, Err(WireError::ValueOutOfRange { .. })));
2667 }
2668
2669 #[test]
2670 fn seqnumset_accepts_exactly_256_bits() {
2671 let mut buf = Vec::new();
2672 buf.extend_from_slice(&SequenceNumber(1).to_bytes_le());
2673 buf.extend_from_slice(&256u32.to_le_bytes());
2674 buf.extend_from_slice(&[0u8; 32]);
2676 let res = SequenceNumberSet::read_from(&buf, 0, true);
2677 assert!(res.is_ok());
2678 }
2679
2680 #[test]
2681 fn fnset_rejects_num_bits_above_256() {
2682 let mut buf = Vec::new();
2683 buf.extend_from_slice(&FragmentNumber(1).to_bytes_le());
2684 buf.extend_from_slice(&1000u32.to_le_bytes()); let res = FragmentNumberSet::read_from(&buf, 0, true);
2686 assert!(matches!(res, Err(WireError::ValueOutOfRange { .. })));
2687 }
2688
2689 #[test]
2690 fn fnset_dos_giant_num_bits_rejected_before_alloc() {
2691 let mut buf = Vec::new();
2694 buf.extend_from_slice(&FragmentNumber(1).to_bytes_le());
2695 buf.extend_from_slice(&u32::MAX.to_le_bytes());
2696 let res = FragmentNumberSet::read_from(&buf, 0, true);
2697 assert!(matches!(res, Err(WireError::ValueOutOfRange { .. })));
2698 }
2699
2700 #[test]
2701 fn data_frag_decode_rejects_wrong_octets_to_inline_qos_when_q_false() {
2702 let d = dataflag_frag(1, 1, 1, 4, 4, vec![1, 2, 3, 4]);
2705 let (mut bytes, _) = d.write_body(true);
2706 bytes[2..4].copy_from_slice(&99u16.to_le_bytes());
2708 let res = DataFragSubmessage::read_body(&bytes, true, false, false, false, false);
2709 assert!(matches!(res, Err(WireError::ValueOutOfRange { .. })));
2710 }
2711
2712 #[test]
2713 fn data_frag_decode_rejects_inline_qos() {
2714 let d = dataflag_frag(1, 1, 1, 4, 4, vec![1, 2, 3, 4]);
2716 let (bytes, _) = d.write_body(true);
2717 let res = DataFragSubmessage::read_body(&bytes, true, true, false, false, false);
2718 assert!(matches!(res, Err(WireError::UnsupportedFeature { .. })));
2719 }
2720
2721 #[test]
2722 fn data_frag_flags_survive_roundtrip() {
2723 let mut d = dataflag_frag(1, 1, 1, 4, 4, vec![1, 2, 3, 4]);
2724 d.hash_key_flag = true;
2725 d.key_flag = true;
2726 d.non_standard_flag = true;
2727 let (bytes, flags) = d.write_body(true);
2728 assert!(flags & DATA_FRAG_FLAG_HASH_KEY != 0);
2729 assert!(flags & DATA_FRAG_FLAG_KEY != 0);
2730 assert!(flags & DATA_FRAG_FLAG_NON_STANDARD != 0);
2731 let decoded = DataFragSubmessage::read_body(&bytes, true, false, true, true, true).unwrap();
2732 assert!(decoded.hash_key_flag);
2733 assert!(decoded.key_flag);
2734 assert!(decoded.non_standard_flag);
2735 }
2736
2737 #[test]
2740 fn heartbeat_frag_roundtrip_le() {
2741 let h = HeartbeatFragSubmessage {
2742 reader_id: reader_id(),
2743 writer_id: writer_id(),
2744 writer_sn: SequenceNumber(42),
2745 last_fragment_num: FragmentNumber(8),
2746 count: 3,
2747 };
2748 let (bytes, flags) = h.write_body(true);
2749 assert!(flags & FLAG_E_LITTLE_ENDIAN != 0);
2750 assert_eq!(bytes.len(), HeartbeatFragSubmessage::WIRE_SIZE);
2751 let decoded = HeartbeatFragSubmessage::read_body(&bytes, true).unwrap();
2752 assert_eq!(decoded, h);
2753 }
2754
2755 #[test]
2756 fn heartbeat_frag_roundtrip_be() {
2757 let h = HeartbeatFragSubmessage {
2758 reader_id: reader_id(),
2759 writer_id: writer_id(),
2760 writer_sn: SequenceNumber(1),
2761 last_fragment_num: FragmentNumber(1),
2762 count: 1,
2763 };
2764 let (bytes, _) = h.write_body(false);
2765 let decoded = HeartbeatFragSubmessage::read_body(&bytes, false).unwrap();
2766 assert_eq!(decoded, h);
2767 }
2768
2769 #[test]
2770 fn heartbeat_frag_decode_rejects_truncated() {
2771 let res = HeartbeatFragSubmessage::read_body(&[0u8; 20], true);
2772 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2773 }
2774
2775 #[test]
2778 fn nack_frag_roundtrip_le() {
2779 let n = NackFragSubmessage {
2780 reader_id: reader_id(),
2781 writer_id: writer_id(),
2782 writer_sn: SequenceNumber(5),
2783 fragment_number_state: FragmentNumberSet {
2784 bitmap_base: FragmentNumber(1),
2785 num_bits: 4,
2786 bitmap: vec![0b1010_0000_0000_0000_0000_0000_0000_0000],
2787 },
2788 count: 2,
2789 };
2790 let (bytes, flags) = n.write_body(true);
2791 assert!(flags & FLAG_E_LITTLE_ENDIAN != 0);
2792 let decoded = NackFragSubmessage::read_body(&bytes, true).unwrap();
2793 assert_eq!(decoded, n);
2794 }
2795
2796 #[test]
2797 fn nack_frag_roundtrip_be() {
2798 let n = NackFragSubmessage {
2799 reader_id: reader_id(),
2800 writer_id: writer_id(),
2801 writer_sn: SequenceNumber(100),
2802 fragment_number_state: FragmentNumberSet {
2803 bitmap_base: FragmentNumber(10),
2804 num_bits: 0,
2805 bitmap: vec![],
2806 },
2807 count: 0,
2808 };
2809 let (bytes, _) = n.write_body(false);
2810 let decoded = NackFragSubmessage::read_body(&bytes, false).unwrap();
2811 assert_eq!(decoded, n);
2812 }
2813
2814 #[test]
2815 fn nack_frag_decode_rejects_truncated() {
2816 let res = NackFragSubmessage::read_body(&[0u8; 20], true);
2817 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2818 }
2819
2820 fn make_info_source() -> InfoSourceSubmessage {
2823 InfoSourceSubmessage {
2824 unused: 0,
2825 protocol_version: crate::wire_types::ProtocolVersion::V2_5,
2826 vendor_id: crate::wire_types::VendorId([0xAB, 0xCD]),
2827 guid_prefix: crate::wire_types::GuidPrefix::from_bytes([0xEE; 12]),
2828 }
2829 }
2830
2831 #[test]
2832 fn info_source_roundtrip_le() {
2833 let i = make_info_source();
2834 let (bytes, flags) = i.write_body(true);
2835 assert!(flags & FLAG_E_LITTLE_ENDIAN != 0);
2836 assert_eq!(bytes.len(), InfoSourceSubmessage::WIRE_SIZE);
2837 let decoded = InfoSourceSubmessage::read_body(&bytes, true).unwrap();
2838 assert_eq!(decoded, i);
2839 }
2840
2841 #[test]
2842 fn info_source_roundtrip_be() {
2843 let i = make_info_source();
2844 let (bytes, flags) = i.write_body(false);
2845 assert_eq!(flags & FLAG_E_LITTLE_ENDIAN, 0);
2846 let decoded = InfoSourceSubmessage::read_body(&bytes, false).unwrap();
2847 assert_eq!(decoded, i);
2848 }
2849
2850 #[test]
2851 fn info_source_wire_size_is_20() {
2852 let i = make_info_source();
2853 let (bytes, _) = i.write_body(true);
2854 assert_eq!(bytes.len(), 20);
2855 }
2856
2857 #[test]
2858 fn info_source_decode_rejects_truncated() {
2859 let res = InfoSourceSubmessage::read_body(&[0u8; 19], true);
2860 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2861 }
2862
2863 #[test]
2864 fn info_source_unused_field_roundtrips() {
2865 let mut i = make_info_source();
2868 i.unused = 0xDEAD_BEEF;
2869 let (bytes, _) = i.write_body(true);
2870 let decoded = InfoSourceSubmessage::read_body(&bytes, true).unwrap();
2871 assert_eq!(decoded.unused, 0xDEAD_BEEF);
2872 }
2873
2874 #[test]
2877 fn info_reply_unicast_only_roundtrip_le() {
2878 use crate::wire_types::Locator;
2879 let i = InfoReplySubmessage {
2880 unicast_locators: vec![
2881 Locator::udp_v4([10, 0, 0, 1], 7411),
2882 Locator::udp_v4([10, 0, 0, 2], 7411),
2883 ],
2884 multicast_locators: None,
2885 };
2886 let (bytes, flags) = i.write_body(true);
2887 assert_eq!(flags & INFO_REPLY_FLAG_MULTICAST, 0);
2888 let decoded = InfoReplySubmessage::read_body(&bytes, true, false).unwrap();
2889 assert_eq!(decoded, i);
2890 }
2891
2892 #[test]
2893 fn info_reply_with_multicast_roundtrip_le() {
2894 use crate::wire_types::Locator;
2895 let i = InfoReplySubmessage {
2896 unicast_locators: vec![Locator::udp_v4([10, 0, 0, 1], 7411)],
2897 multicast_locators: Some(vec![Locator::udp_v4([239, 255, 0, 1], 7400)]),
2898 };
2899 let (bytes, flags) = i.write_body(true);
2900 assert!(flags & INFO_REPLY_FLAG_MULTICAST != 0);
2901 let decoded = InfoReplySubmessage::read_body(&bytes, true, true).unwrap();
2902 assert_eq!(decoded, i);
2903 }
2904
2905 #[test]
2906 fn info_reply_with_multicast_roundtrip_be() {
2907 use crate::wire_types::Locator;
2908 let i = InfoReplySubmessage {
2909 unicast_locators: vec![Locator::udp_v4([10, 0, 0, 5], 7420)],
2910 multicast_locators: Some(vec![Locator::udp_v4([239, 255, 0, 9], 7400)]),
2911 };
2912 let (bytes, _) = i.write_body(false);
2913 let decoded = InfoReplySubmessage::read_body(&bytes, false, true).unwrap();
2914 assert_eq!(decoded, i);
2915 }
2916
2917 #[test]
2918 fn info_reply_empty_unicast_list_is_valid() {
2919 let i = InfoReplySubmessage {
2922 unicast_locators: vec![],
2923 multicast_locators: None,
2924 };
2925 let (bytes, _) = i.write_body(true);
2926 let decoded = InfoReplySubmessage::read_body(&bytes, true, false).unwrap();
2927 assert_eq!(decoded, i);
2928 assert!(decoded.unicast_locators.is_empty());
2929 }
2930
2931 #[test]
2932 fn info_reply_decode_rejects_oversized_locator_list_length() {
2933 let mut body = Vec::new();
2935 body.extend_from_slice(&u32::MAX.to_le_bytes());
2936 let res = InfoReplySubmessage::read_body(&body, true, false);
2938 assert!(matches!(res, Err(WireError::ValueOutOfRange { .. })));
2939 }
2940
2941 #[test]
2942 fn info_reply_decode_rejects_truncated_length_field() {
2943 let res = InfoReplySubmessage::read_body(&[0u8; 3], true, false);
2944 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2945 }
2946
2947 #[test]
2950 fn info_timestamp_roundtrip_le() {
2951 let i = InfoTimestampSubmessage {
2952 timestamp: crate::header_extension::HeTimestamp {
2953 seconds: 0x1234_5678,
2954 fraction: 0x9ABC_DEF0,
2955 },
2956 invalidate: false,
2957 };
2958 let (bytes, flags) = i.write_body(true);
2959 assert_eq!(flags & INFO_TIMESTAMP_FLAG_INVALIDATE, 0);
2960 assert_eq!(bytes.len(), 8);
2961 let decoded = InfoTimestampSubmessage::read_body(&bytes, true, false).unwrap();
2962 assert_eq!(decoded, i);
2963 }
2964
2965 #[test]
2966 fn info_timestamp_roundtrip_be() {
2967 let i = InfoTimestampSubmessage {
2968 timestamp: crate::header_extension::HeTimestamp {
2969 seconds: 1_700_000_000,
2970 fraction: 12345,
2971 },
2972 invalidate: false,
2973 };
2974 let (bytes, flags) = i.write_body(false);
2975 assert_eq!(flags & FLAG_E_LITTLE_ENDIAN, 0);
2976 let decoded = InfoTimestampSubmessage::read_body(&bytes, false, false).unwrap();
2977 assert_eq!(decoded, i);
2978 }
2979
2980 #[test]
2981 fn info_timestamp_invalidate_flag_yields_empty_body() {
2982 let i = InfoTimestampSubmessage {
2983 timestamp: crate::header_extension::HeTimestamp::default(),
2984 invalidate: true,
2985 };
2986 let (bytes, flags) = i.write_body(true);
2987 assert!(flags & INFO_TIMESTAMP_FLAG_INVALIDATE != 0);
2988 assert!(bytes.is_empty(), "I-Flag → empty body");
2989 let decoded = InfoTimestampSubmessage::read_body(&bytes, true, true).unwrap();
2990 assert!(decoded.invalidate);
2991 }
2992
2993 #[test]
2994 fn info_timestamp_decode_rejects_truncated_when_no_invalidate() {
2995 let res = InfoTimestampSubmessage::read_body(&[0u8; 4], true, false);
2996 assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2997 }
2998
2999 #[test]
3000 fn info_timestamp_decode_with_invalidate_ignores_body() {
3001 let res = InfoTimestampSubmessage::read_body(&[0u8; 8], true, true).unwrap();
3004 assert!(res.invalidate);
3005 assert_eq!(
3006 res.timestamp,
3007 crate::header_extension::HeTimestamp::default()
3008 );
3009 }
3010}