Skip to main content

zerodds_rtps/
submessages.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright 2026 ZeroDDS Contributors
3//! RTPS submessages — DDSI-RTPS 2.5 §8.3.7.
4//!
5//! Implemented submessages:
6//! - **DATA** (`DataSubmessage`) — §8.3.7.2.
7//! - **HEARTBEAT** (`HeartbeatSubmessage`) — §8.3.7.5.
8//! - **ACKNACK** (`AckNackSubmessage`) — §8.3.7.1.
9//! - **GAP** (`GapSubmessage`) — §8.3.7.4.
10//! - **DATA_FRAG** (`DataFragSubmessage`) — §8.3.7.3.
11//! - **HEARTBEAT_FRAG** (`HeartbeatFragSubmessage`) — §8.3.7.6.
12//! - **NACK_FRAG** (`NackFragSubmessage`) — §8.3.7.10.
13//! - **INFO_SRC** (`InfoSourceSubmessage`) — §8.3.7.9.
14//! - **INFO_TS** (`InfoTimestampSubmessage`) — §8.3.7.5/§8.3.8.5.
15//! - **INFO_REPLY** (`InfoReplySubmessage`) — §8.3.7.8.
16//!
17//! ParameterList (inline QoS) lives in the separate module
18//! [`crate::parameter_list`]; SecuredPayload wrapping is in the
19//! `zerodds-security` crate (DDS-Security 1.2 §7.4).
20//!
21//! # Endianness
22//!
23//! Submessage bodies are written in the endianness of the submessage
24//! header (E-flag). The `to_bytes_*`/`from_bytes_*` functions given here
25//! take explicit endianness as a parameter — the caller must choose it
26//! consistently with the header.
27
28extern 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
36/// Hard cap for `numBits` in `SequenceNumberSet` and
37/// `FragmentNumberSet`. DDSI-RTPS gives no specific limit, but both
38/// Cyclone DDS (`ddsi_radmin.c`) and Fast-DDS (`BitmapRange<..., 256>`)
39/// cap at 256. We follow — prevents DoS via a `numBits=2^32-1` bitmap
40/// alloc.
41pub const RTPS_BITMAP_MAX_BITS: u32 = 256;
42
43// ============================================================================
44// SequenceNumberSet (§9.4.2.6)
45// ============================================================================
46
47/// Bitset of sequence numbers from `bitmap_base`. Used in HEARTBEAT/
48/// ACKNACK/GAP to signal sets of lost or known packets.
49///
50/// Wire layout (variable length):
51///   bitmapBase: 8 byte (SequenceNumber, big or little per header)
52///   numBits:    4 byte (u32)
53///   bitmap:     ceil(numBits/32) * 4 byte (u32 words)
54#[derive(Debug, Clone, PartialEq, Eq)]
55pub struct SequenceNumberSet {
56    /// First sequence number that the first bit is responsible for.
57    pub bitmap_base: SequenceNumber,
58    /// Number of valid bits.
59    pub num_bits: u32,
60    /// `ceil(num_bits/32)` u32 words.
61    pub bitmap: Vec<u32>,
62}
63
64impl SequenceNumberSet {
65    /// Computes the wire size in bytes based on `num_bits`.
66    #[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    /// Builds a `SequenceNumberSet` from a sorted list of missing SNs.
73    ///
74    /// `base` is the smallest not-yet-acked SN (the AckNack base).
75    /// `missing` must be sorted ascending and all SNs ≥ `base`. Bits are
76    /// set in RTPS convention: bit 0 is the most-significant bit (MSB) of
77    /// `bitmap[0]`.
78    #[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    /// Iterates over all SNs whose bit is set.
116    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    /// Tatsaechliche Wire-Size dieses Sets.
129    #[must_use]
130    pub fn encoded_size(&self) -> usize {
131        Self::wire_size(self.num_bits)
132    }
133
134    /// Encodes the set into `out` with the given endianness.
135    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    /// Decodes a set from `bytes` at `offset`. Returns (set, new position).
152    ///
153    /// # Errors
154    /// `UnexpectedEof`.
155    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// ============================================================================
225// FragmentNumberSet (§9.4.2.8)
226// ============================================================================
227
228/// Bitset of `FragmentNumber` values from `bitmap_base`. Analogous to
229/// [`SequenceNumberSet`], but with `FragmentNumber` (u32) as the base
230/// instead of `SequenceNumber`.
231///
232/// Wire layout:
233///   bitmapBase: 4 byte (FragmentNumber, LE or BE per header)
234///   numBits:    4 byte (u32)
235///   bitmap:     ceil(numBits/32) * 4 byte
236#[derive(Debug, Clone, PartialEq, Eq)]
237pub struct FragmentNumberSet {
238    /// First fragment that the first bit is responsible for.
239    pub bitmap_base: FragmentNumber,
240    /// Number of valid bits.
241    pub num_bits: u32,
242    /// `ceil(num_bits/32)` u32 words.
243    pub bitmap: Vec<u32>,
244}
245
246impl FragmentNumberSet {
247    /// Wire size in bytes.
248    #[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    /// Builds the set from a sorted list of missing fragments.
255    /// `base` = smallest not-yet-acknowledged FragmentNumber.
256    #[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        // DDSI-RTPS §8.3.5.4: numBits MUST be <= 256. A gap over more
273        // than 256 fragments (large samples under packet loss) would
274        // otherwise produce a num_bits > 256, which every spec-conformant
275        // receiver discards as malformed → the NACK_FRAG is lost →
276        // fragment stall. We cover only the first 256; the rest follows
277        // in the next NACK_FRAG once bitmap_base has advanced.
278        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            // Skip fragments beyond the 256-bit window (follow-up NACK_FRAG).
287            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    /// Iterates over all set FragmentNumbers.
304    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    /// Tatsaechliche Wire-Size dieses Sets.
317    #[must_use]
318    pub fn encoded_size(&self) -> usize {
319        Self::wire_size(self.num_bits)
320    }
321
322    /// Encodes the set into `out`.
323    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    /// Decodes a set from `bytes` at `offset`.
340    ///
341    /// # Errors
342    /// `UnexpectedEof`.
343    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
412// ============================================================================
413// DATA Submessage (§8.3.7.2)
414// ============================================================================
415
416/// DATA-Submessage Flag: Q (Inline-QoS present).
417pub const DATA_FLAG_INLINE_QOS: u8 = 0x02;
418/// DATA-Submessage Flag: D (data payload present).
419pub const DATA_FLAG_DATA: u8 = 0x04;
420/// DATA-Submessage Flag: K (key payload present, Q-flag mutually exclusive with D).
421pub const DATA_FLAG_KEY: u8 = 0x08;
422/// DATA-Submessage Flag: N (non-standard payload).
423pub const DATA_FLAG_NON_STANDARD: u8 = 0x10;
424
425/// DATA submessage. Phase 0 supports only the D-flag (data), no Q
426/// (no inline QoS), no K, no N.
427///
428/// `serialized_payload` is `Arc<[u8]>` (WP 2.0a zero-copy spike).
429/// Writers share the payload allocation with `CacheChange` and all
430/// DATA/DATA_FRAG datagrams — `clone()` on this struct is a pure
431/// refcount bump.
432#[derive(Debug, Clone, PartialEq, Eq)]
433pub struct DataSubmessage {
434    /// Reserved extra flags (uint16, mostly 0).
435    pub extra_flags: u16,
436    /// Receiver EntityId.
437    pub reader_id: EntityId,
438    /// Sender EntityId.
439    pub writer_id: EntityId,
440    /// Sequence number of this DATA.
441    pub writer_sn: SequenceNumber,
442    /// Inline-QoS ParameterList (Q-flag, §9.4.5.3.2). `None` = no
443    /// Q-flag, no inline-QoS block. Carrier for PID_KEY_HASH (WP 1.B),
444    /// PID_STATUS_INFO, PID_COHERENT_SET etc.
445    pub inline_qos: Option<crate::parameter_list::ParameterList>,
446    /// K-flag (spec §8.3.8.2 Tab. 8.43). `true`: `serialized_payload`
447    /// contains only the @key fields (key-only sample, e.g. a dispose
448    /// marker). The D-flag can be false at the same time when only the
449    /// key is sent; in that case `serialized_payload` is an
450    /// XCDR-encoded key holder.
451    pub key_flag: bool,
452    /// N-flag (spec §8.3.8.2 Tab. 8.43, NonStandardPayloadFlag).
453    /// `true`: `serialized_payload` is NOT encoded in the CDR variant
454    /// implied by `representation_identifier` (e.g. for
455    /// DDS-Security-encrypted payloads).
456    pub non_standard_flag: bool,
457    /// Serialized payload (XCDR2-encoded or vendor-specific).
458    pub serialized_payload: Arc<[u8]>,
459}
460
461impl DataSubmessage {
462    /// Encodes the DATA body (without the submessage header) into a Vec.
463    /// Automatically sets the D-flag and, if applicable, the Q-flag in
464    /// the `flags` output (returned), so the caller can fill the
465    /// submessage header correctly.
466    ///
467    /// Layout: extraFlags(2) + octetsToInlineQos(2) + readerId(4) +
468    /// writerId(4) + writerSN(8) + [optional InlineQoS PL] + payload.
469    #[must_use]
470    pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
471        // Encode inline QoS once (instead of amortizing it with
472        // `out.extend_from_slice`) and thereby set the Vec capacity
473        // *exactly*. This eliminates the reallocation cascade with
474        // `extend_from_slice` — visible in the macOS recv-thread profile
475        // (`finish_grow`/`reserve` ~17% of samples before this refactor).
476        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        // 20 = extraFlags(2)+octetsToInlineQos(2)+readerId(4)+writerId(4)
482        // +writerSN(8).
483        let mut out = Vec::with_capacity(20 + inline_qos_len + self.serialized_payload.len());
484        // extraFlags (2 byte)
485        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        // octetsToInlineQos (2 byte) — distance from the end of this field to
492        // the start of readerId. Constant 16 (4 readerId + 4 writerId
493        // + 8 writerSN), independent of the Q flag.
494        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        // readerId, writerId
502        out.extend_from_slice(&self.reader_id.to_bytes());
503        out.extend_from_slice(&self.writer_id.to_bytes());
504        // writerSN
505        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        // Inline-QoS ParameterList (Q-flag) — if present.
511        if let Some(qos_bytes) = inline_qos_buf {
512            out.extend_from_slice(&qos_bytes);
513        }
514        // serializedPayload
515        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    /// Decodes the DATA body from a slice. A backward-compat wrapper for
535    /// callers that carry no Q-flag — inline QoS is ignored.
536    ///
537    /// # Errors
538    /// `UnexpectedEof` on a too-short body.
539    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    /// Decodes the DATA body taking the submessage flags into account.
544    /// If `flags & DATA_FLAG_INLINE_QOS != 0` the decoder parses the
545    /// ParameterList after the writerSN, before the payload.
546    ///
547    /// # Errors
548    /// `UnexpectedEof` on a too-short body. ParameterList errors are
549    /// passed through as `WireError`.
550    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        // Read octetsToInlineQos — we do not use it directly (always 16),
571        // but we consume the field.
572        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        // Inline QoS — Q-flag (DATA_FLAG_INLINE_QOS = 0x02).
591        let inline_qos = if flags & DATA_FLAG_INLINE_QOS != 0 {
592            // ParameterList::from_bytes parses up to the sentinel and
593            // returns the rest of the buffer to us via consumed bytes.
594            // Since from_bytes only returns the list, we must track the
595            // consumed length ourselves — we compute it from the list's
596            // encode output.
597            let pl = crate::parameter_list::ParameterList::from_bytes(&body[pos..], little_endian)?;
598            // Re-encode to determine the consumed byte length. This is
599            // robust and avoids drift with the from_bytes parser.
600            let consumed = pl.to_bytes(little_endian).len();
601            pos += consumed;
602            Some(pl)
603        } else {
604            None
605        };
606
607        // The rest is serializedPayload.
608        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
624// ============================================================================
625// HEARTBEAT Submessage (§8.3.7.5)
626// ============================================================================
627
628/// HEARTBEAT Flag: F (Final).
629pub const HEARTBEAT_FLAG_FINAL: u8 = 0x02;
630/// HEARTBEAT Flag: L (Liveliness).
631pub const HEARTBEAT_FLAG_LIVELINESS: u8 = 0x04;
632/// HEARTBEAT flag: G (GroupInfo present). A vendor extension for
633/// group-ordered access (§8.3.8.6.2). The trailer contains currentGSN,
634/// firstGSN, lastGSN and a `writerSet` (list of the GUID prefixes of the
635/// group members).
636pub const HEARTBEAT_FLAG_GROUP_INFO: u8 = 0x08;
637
638/// Optional GroupInfo trailer of a HEARTBEAT submessage (§8.3.8.6.2).
639///
640/// Wire layout:
641/// - currentGSN: i64
642/// - firstGSN:   i64
643/// - lastGSN:    i64
644/// - writerSet:  u32 length + length × GuidPrefix(12 byte)
645#[derive(Debug, Clone, PartialEq, Eq)]
646pub struct HeartbeatGroupInfo {
647    /// Current group SN (highest assigned by the group coordinator).
648    pub current_gsn: SequenceNumber,
649    /// First relevant group SN (cache_min of the group).
650    pub first_gsn: SequenceNumber,
651    /// Last available group SN (= currentGSN minus pending, in practice
652    /// identical to currentGSN at steady state).
653    pub last_gsn: SequenceNumber,
654    /// GuidPrefix set of the participating writers of this group.
655    pub writer_set: Vec<crate::wire_types::GuidPrefix>,
656}
657
658/// HEARTBEAT submessage.
659///
660/// `final_flag`, `liveliness_flag` and `group_info_flag` (via `Some` of
661/// `group_info`) correspond to the F-/L-/G bits in the submessage header
662/// (spec §8.3.7.5.1, §8.3.8.6.2) — they are **not** in the body, but are
663/// carried here as a semantic part of the message.
664#[derive(Debug, Clone, PartialEq, Eq)]
665pub struct HeartbeatSubmessage {
666    /// Reader EntityId (target).
667    pub reader_id: EntityId,
668    /// Writer EntityId (source).
669    pub writer_id: EntityId,
670    /// First available sequence number in the history cache.
671    pub first_sn: SequenceNumber,
672    /// Last sent sequence number.
673    pub last_sn: SequenceNumber,
674    /// Count_t (i32) — heartbeat sequence number (for ACK correlation).
675    pub count: i32,
676    /// F-flag: `true` = the reader need not send a response when complete.
677    pub final_flag: bool,
678    /// L-flag: liveliness announce (without history semantics).
679    pub liveliness_flag: bool,
680    /// G-flag (§8.3.8.6.2): optional GroupInfo trailer.
681    pub group_info: Option<HeartbeatGroupInfo>,
682}
683
684impl HeartbeatSubmessage {
685    /// Minimal wire size (body without GroupInfo): 28 bytes (4+4+8+8+4).
686    /// Flags are in the submessage header.
687    pub const WIRE_SIZE: usize = 28;
688
689    /// Encodes the body. Returns (bytes, flags), where `flags` is the
690    /// submessage-header flag byte incl. E/F/L/G.
691    #[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    /// Decodes the body. `final_flag`, `liveliness_flag`,
744    /// `group_info_flag` are extracted by the caller from the submessage header.
745    ///
746    /// # Errors
747    /// `UnexpectedEof`, `ValueOutOfRange` (writerSet length bizarrely large).
748    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            // 3 × i64 + u32 = 28 byte minimum
795            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            // Cap: writer_set must not be larger than what the body has
832            // left. Protection against DoS via a huge length field.
833            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
868// ============================================================================
869// ACKNACK Submessage (§8.3.7.1)
870// ============================================================================
871
872/// ACKNACK Flag: F (Final).
873pub const ACKNACK_FLAG_FINAL: u8 = 0x02;
874
875/// ACKNACK submessage.
876///
877/// `final_flag` corresponds to the F-bit in the submessage header (spec
878/// §8.3.7.1.1). `final=false` requires a timely HEARTBEAT response from
879/// the writer.
880#[derive(Debug, Clone, PartialEq, Eq)]
881pub struct AckNackSubmessage {
882    /// Reader EntityId (source).
883    pub reader_id: EntityId,
884    /// Writer EntityId (target).
885    pub writer_id: EntityId,
886    /// Bitset of the not-yet-received sequence numbers.
887    pub reader_sn_state: SequenceNumberSet,
888    /// Count_t (for correlation with HEARTBEAT.count).
889    pub count: i32,
890    /// F-flag: `false` = the writer should answer with a timely HEARTBEAT.
891    pub final_flag: bool,
892}
893
894impl AckNackSubmessage {
895    /// Encodes the body. Returns (bytes, flags).
896    #[must_use]
897    pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
898        // 4 readerId + 4 writerId + 12 SN base + 4*num_words bitmap +
899        // 4 count. The SN set can be variable; 12 + words*4 is the upper
900        // bound. Pre-allocating saves the realloc step on the recv-thread
901        // hot path (HEARTBEAT response).
902        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    /// Decodes the body. `final_flag` is extracted by the caller from
923    /// the submessage header.
924    ///
925    /// # Errors
926    /// `UnexpectedEof`.
927    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
972// ============================================================================
973// GAP Submessage (§8.3.7.4 / §8.3.8.4.2)
974// ============================================================================
975
976/// GAP Flag: G (GroupInfo present — `gapStartGSN`/`gapEndGSN` Trailer).
977/// A vendor extension for group-ordered access (§8.3.8.4.2). The ZeroDDS encoder does not set it; the decoder accepts it on read.
978pub const GAP_FLAG_GROUP_INFO: u8 = 0x04;
979
980/// GAP flag: K (FilteredCount present). Spec §8.3.8.4.2 introduces an
981/// optional `Count_t filteredCount` trailer field that lets the reader
982/// distinguish "discarded via content filter" from "really removed" — a
983/// prerequisite for correct instance-state transitions per §8.7.4
984/// (NOT_ALIVE_FILTERED vs. NOT_ALIVE_DISPOSED).
985pub const GAP_FLAG_FILTERED_COUNT: u8 = 0x08;
986
987/// Optional trailer of a GAP submessage with GroupInfo (G-flag, §8.3.8.4.2).
988#[derive(Debug, Clone, Copy, PartialEq, Eq)]
989pub struct GapGroupInfo {
990    /// Group SN of the first skipped sample in the group.
991    pub gap_start_gsn: SequenceNumber,
992    /// Group SN of the last skipped sample in the group.
993    pub gap_end_gsn: SequenceNumber,
994}
995
996/// GAP submessage. Signals the reader that the writer will never send
997/// sequence numbers `[gap_start, gap_list.bitmap_base)` (all before
998/// `gap_list.bitmap_base` are gaps; the bits in `gap_list` mark
999/// individual further gaps from `bitmap_base`).
1000#[derive(Debug, Clone, PartialEq, Eq)]
1001pub struct GapSubmessage {
1002    /// Reader EntityId (target).
1003    pub reader_id: EntityId,
1004    /// Writer EntityId (source).
1005    pub writer_id: EntityId,
1006    /// First irreversible gap SN.
1007    pub gap_start: SequenceNumber,
1008    /// Bitset of the further gaps from `gap_list.bitmap_base`.
1009    pub gap_list: SequenceNumberSet,
1010    /// Optional GroupInfo (§8.3.8.4.2). `Some` ⇒ G flag set in the header.
1011    pub group_info: Option<GapGroupInfo>,
1012    /// Optional `filteredCount` trailer (§8.3.8.4.2). `Some` ⇒
1013    /// K flag set in the header. `0` is explicitly "nothing filtered,
1014    /// everything really removed"; `1+` means "n samples discarded via
1015    /// content filter".
1016    pub filtered_count: Option<u32>,
1017}
1018
1019impl GapSubmessage {
1020    /// Encodes the body. Returns (bytes, flags) incl. possible G/K bit.
1021    #[must_use]
1022    pub fn write_body(&self, little_endian: bool) -> (Vec<u8>, u8) {
1023        // 4 readerId + 4 writerId + 8 gap_start + 12 SN-Set-Base +
1024        // 4*words. Optional GroupInfo (24) + filteredCount (4) on top.
1025        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    /// Decodes the body. Flags G/K are passed by the caller from the
1066    /// submessage header (see `decode_datagram`).
1067    ///
1068    /// # Errors
1069    /// `UnexpectedEof`.
1070    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
1159// ============================================================================
1160// DATA_FRAG Submessage (§8.3.7.3)
1161// ============================================================================
1162
1163/// DATA_FRAG Flag: Q (Inline-QoS present).
1164pub const DATA_FRAG_FLAG_INLINE_QOS: u8 = 0x02;
1165/// DATA_FRAG Flag: H (hash key).
1166pub const DATA_FRAG_FLAG_HASH_KEY: u8 = 0x04;
1167/// DATA_FRAG flag: K (key flag — serialized_payload is key instead of data).
1168pub const DATA_FRAG_FLAG_KEY: u8 = 0x08;
1169/// DATA_FRAG Flag: N (non-standard payload).
1170pub const DATA_FRAG_FLAG_NON_STANDARD: u8 = 0x10;
1171
1172/// DATA_FRAG submessage. Carries a section (fragments) of a sample
1173/// whose total size is in `sample_size`.
1174///
1175/// Flags (Q/H/K/N) are mirrored from the submessage header. The encoder does not currently set these flags; the decoder accepts them on read.
1176#[derive(Debug, Clone, PartialEq, Eq)]
1177pub struct DataFragSubmessage {
1178    /// octetsToInlineQos analogue to DATA (§8.3.7.2 speaks of extraFlags+
1179    /// octetsToInlineQos; this variant carries 0).
1180    pub extra_flags: u16,
1181    /// Reader EntityId (target).
1182    pub reader_id: EntityId,
1183    /// Writer EntityId (source).
1184    pub writer_id: EntityId,
1185    /// Sequence number of the sample whose fragments this carries.
1186    pub writer_sn: SequenceNumber,
1187    /// First fragment in this submessage (1-based).
1188    pub fragment_starting_num: FragmentNumber,
1189    /// Number of fragments in this submessage. Writer: always 1.
1190    pub fragments_in_submessage: u16,
1191    /// Size of a single fragment (the last may be shorter).
1192    pub fragment_size: u16,
1193    /// Total size of the sample in bytes.
1194    pub sample_size: u32,
1195    /// Fragmented payload section. Arc-shared:
1196    /// writer re-sends are just refcount bumps, no copy.
1197    pub serialized_payload: Arc<[u8]>,
1198    /// Q-flag from the submessage header (inline_qos present).
1199    pub inline_qos_flag: bool,
1200    /// H-flag from the submessage header (hash_key).
1201    pub hash_key_flag: bool,
1202    /// K-flag from the submessage header (serialized_payload = key).
1203    pub key_flag: bool,
1204    /// N-flag from the submessage header (non-standard payload).
1205    pub non_standard_flag: bool,
1206}
1207
1208impl DataFragSubmessage {
1209    /// Minimal body size without payload: extraFlags(2) + octetsToInlineQos(2)
1210    /// + readerId(4) + writerId(4) + writerSN(8) + fragmentStartingNum(4)
1211    /// + fragmentsInSubmessage(2) + fragmentSize(2) + sampleSize(4) = 32.
1212    pub const HEADER_WIRE_SIZE: usize = 32;
1213
1214    /// octetsToInlineQos: offset from the end of this field to the start
1215    /// of inlineQos or serializedPayload. Variant with
1216    /// Q=false: offset = 28 (readerId..sampleSize).
1217    pub const OCTETS_TO_INLINE_QOS: u16 = 28;
1218
1219    /// Encodes the body. Returns (bytes, flags).
1220    #[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    /// Decodes the body. Flags are passed by the caller from the
1272    /// submessage header.
1273    ///
1274    /// # Errors
1275    /// `UnexpectedEof`.
1276    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        // extra_flags: see DATA. 2.1 readers must ignore them
1300        // (Cyclone + Fast-DDS do), we do too — only read and pass on.
1301        // octetsToInlineQos (2 byte): offset from the start of readerId
1302        // to the inline QoS, or (with Q=false) to the serializedPayload.
1303        // The spec requires 28 = header size − 2 (extra_flags) − 2 (this
1304        // field). A deviation at Q=false we catch here — a frequent
1305        // interop bug that we reject spec-faithfully.
1306        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        // Q-flag = false, so no inline-QoS block.
1368        // With Q=true, ParameterList bytes would follow here — we do not
1369        // currently accept that.
1370        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// ============================================================================
1395// InfoSource Submessage (§8.3.7.9 / §8.3.8.9.4) — submessageId 0x0c (legacy
1396// table) or 0x0A in the 2.5 PSM. We follow 2.5: id=0x0A.
1397// ============================================================================
1398
1399/// InfoSource submessage (§8.3.8.9.4). Resets `sourceProtocolVersion`,
1400/// `sourceVendorId`, `sourceGuidPrefix` in the ReceiverState — all
1401/// subsequent submessages are attributed to this source (not the
1402/// datagram header).
1403///
1404/// Wire layout (body, 20 byte):
1405/// - unused (4 byte, "Long unused" in the spec)
1406/// - ProtocolVersion (2 byte)
1407/// - VendorId (2 byte)
1408/// - GuidPrefix (12 byte)
1409#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1410pub struct InfoSourceSubmessage {
1411    /// Reserved 4 byte (spec says: must be 0 from the sender, ignored by
1412    /// the receiver).
1413    pub unused: u32,
1414    /// Source ProtocolVersion (e.g. 2.5).
1415    pub protocol_version: crate::wire_types::ProtocolVersion,
1416    /// Source-VendorId (Hersteller-Kennung).
1417    pub vendor_id: crate::wire_types::VendorId,
1418    /// Source-GuidPrefix (12 byte).
1419    pub guid_prefix: crate::wire_types::GuidPrefix,
1420}
1421
1422impl InfoSourceSubmessage {
1423    /// Wire-Size: 20 Bytes (4+2+2+12).
1424    pub const WIRE_SIZE: usize = 20;
1425
1426    /// Encodes the body. Returns (bytes, flags) incl. E bit.
1427    #[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    /// Decoded den Body.
1446    ///
1447    /// # Errors
1448    /// `UnexpectedEof`.
1449    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
1485// ============================================================================
1486// InfoTimestamp Submessage (§8.3.8.5 / §8.3.7.5) — submessageId 0x09
1487// ============================================================================
1488
1489/// InfoTimestamp flag: I (Invalidate). When set: the body is empty and
1490/// `haveTimestamp` is set to `false` in the ReceiverState.
1491pub const INFO_TIMESTAMP_FLAG_INVALIDATE: u8 = 0x02;
1492
1493/// InfoTimestamp submessage (§8.3.7.5 / §8.3.8.5). Sets the `timestamp`
1494/// field + `haveTimestamp` flag in the ReceiverState.
1495/// Inverted via `INFO_TIMESTAMP_FLAG_INVALIDATE` (I-flag).
1496#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
1497pub struct InfoTimestampSubmessage {
1498    /// Timestamp (8 byte: i32 sec + u32 fraction). Ignored when
1499    /// `invalidate=true`.
1500    pub timestamp: crate::header_extension::HeTimestamp,
1501    /// `true` = the I-flag is set in the submessage → the body is empty
1502    /// and the receiver sets `haveTimestamp = false`.
1503    pub invalidate: bool,
1504}
1505
1506impl InfoTimestampSubmessage {
1507    /// Encodes the body. If `invalidate=true`: body empty (0 byte).
1508    /// Sonst: 8 byte Time_t.
1509    #[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    /// Decodes the body. When `invalidate_flag=true`: expects an empty
1536    /// body and returns `timestamp = default()`.
1537    ///
1538    /// # Errors
1539    /// `UnexpectedEof` if `invalidate_flag=false` and the body is < 8 byte.
1540    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
1578// ============================================================================
1579// InfoReply Submessage (§8.3.7.10 / §8.3.8.10.4) — submessageId 0x0F
1580// ============================================================================
1581
1582/// InfoReply flag: M (multicast). If set: a second LocatorList
1583/// (multicastReplyLocatorList) folgt im Body.
1584pub const INFO_REPLY_FLAG_MULTICAST: u8 = 0x02;
1585
1586/// InfoReply submessage (§8.3.8.10.4). Sets `unicastReplyLocatorList`
1587/// (mandatory) and, if applicable, `multicastReplyLocatorList` (with the
1588/// M-flag) in the ReceiverState.
1589///
1590/// Wire layout (body):
1591/// - unicastLocatorList: u32 length + N × 24 byte locator
1592/// - (M-flag) multicastLocatorList: u32 length + N × 24 byte locator
1593#[derive(Debug, Clone, PartialEq, Eq)]
1594pub struct InfoReplySubmessage {
1595    /// Unicast reply locators (at least 1 sensible, empty list allowed).
1596    pub unicast_locators: Vec<crate::wire_types::Locator>,
1597    /// Multicast reply locators (`Some` ⇒ M flag set in the header).
1598    pub multicast_locators: Option<Vec<crate::wire_types::Locator>>,
1599}
1600
1601impl InfoReplySubmessage {
1602    /// Encodes the body. Returns (bytes, flags) incl. E and possibly M bit.
1603    #[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            // The locator has its own wire format (24 byte). We take the
1631            // LE path here — the locator is always LE in RTPS in the
1632            // ParameterList paths; for the submessage body we follow the
1633            // submessage endianness.
1634            if little_endian {
1635                out.extend_from_slice(&loc.to_bytes_le());
1636            } else {
1637                // BE variant: kind (4 byte BE), port (4 byte BE), addr (16 byte raw)
1638                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    /// Decodes the body. The M-flag is extracted by the caller from the
1646    /// submessage header.
1647    ///
1648    /// # Errors
1649    /// `UnexpectedEof`, `ValueOutOfRange` (locator length bizarrely large).
1650    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            // BE decode: build the locator manually, since from_bytes_le
1698            // strikt LE annimmt.
1699            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// ============================================================================
1725// HEARTBEAT_FRAG Submessage (§8.3.7.7)
1726// ============================================================================
1727
1728/// HEARTBEAT_FRAG submessage. Sent by the writer to inform the reader
1729/// that fragments up to `last_fragment_num` are available for
1730/// `writer_sn`. The writer does not send these; the decoder is kept
1731/// ready anyway for interop.
1732#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1733pub struct HeartbeatFragSubmessage {
1734    /// Reader EntityId (target).
1735    pub reader_id: EntityId,
1736    /// Writer EntityId (source).
1737    pub writer_id: EntityId,
1738    /// Associated sample.
1739    pub writer_sn: SequenceNumber,
1740    /// Highest available FragmentNumber.
1741    pub last_fragment_num: FragmentNumber,
1742    /// Count_t (for correlation with NACK_FRAG).
1743    pub count: i32,
1744}
1745
1746impl HeartbeatFragSubmessage {
1747    /// Wire-Size: 24 Bytes (4+4+8+4+4).
1748    pub const WIRE_SIZE: usize = 24;
1749
1750    /// Encoded den Body.
1751    #[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    /// Decoded den Body.
1779    ///
1780    /// # Errors
1781    /// `UnexpectedEof`.
1782    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// ============================================================================
1832// NACK_FRAG Submessage (§8.3.7.6)
1833// ============================================================================
1834
1835/// NACK_FRAG submessage. The reader reports missing fragments for a
1836/// specific `writer_sn`. No flags on the wire except E.
1837#[derive(Debug, Clone, PartialEq, Eq)]
1838pub struct NackFragSubmessage {
1839    /// Reader EntityId (source).
1840    pub reader_id: EntityId,
1841    /// Writer EntityId (target).
1842    pub writer_id: EntityId,
1843    /// Associated sample.
1844    pub writer_sn: SequenceNumber,
1845    /// Bitset of missing fragments.
1846    pub fragment_number_state: FragmentNumberSet,
1847    /// Count_t (for correlation).
1848    pub count: i32,
1849}
1850
1851impl NackFragSubmessage {
1852    /// Encoded den Body.
1853    #[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    /// Decoded den Body.
1877    ///
1878    /// # Errors
1879    /// `UnexpectedEof`.
1880    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    // ---- SequenceNumberSet ----
1944
1945    #[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        // numBits=64 → 8 byte bitmap expected; only 4 present.
1990        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]); // only 4 instead of 8
1994        let res = SequenceNumberSet::read_from(&buf, 0, true);
1995        assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
1996    }
1997
1998    // ---- DATA Submessage ----
1999
2000    #[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        // Spec §8.3.8.2 K-flag: serialized_payload contains only the key.
2040        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        // Spec §8.3.8.2 N-Flag: NonStandardPayload (z.B. Encrypted).
2061        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        // E + Q + D + K + N all set — the full 5-flag roundtrip.
2082        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        // bytes[2..4] = octetsToInlineQos LE
2118        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    // ---- HEARTBEAT Submessage ----
2128
2129    #[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    // ---- WP 1.E stage D: HEARTBEAT GroupInfo ----
2191
2192    #[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        // length=u32::MAX waere 12 × MAX byte → DoS-Schutz.
2248        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        // 3 × i64 GSN
2255        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        // bizarre length
2259        body.extend_from_slice(&u32::MAX.to_le_bytes());
2260        // no body for prefixes
2261        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        // body ends before the 3 GSN fields
2268        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        // GroupInfo trailer missing → UnexpectedEof
2275        let res = HeartbeatSubmessage::read_body(&body, true, false, false, true);
2276        assert!(matches!(res, Err(WireError::UnexpectedEof { .. })));
2277    }
2278
2279    // ---- ACKNACK Submessage ----
2280
2281    #[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    // ---- GAP Submessage ----
2320
2321    #[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    // ---- WP 1.E stage C: GAP filteredCount + GroupInfo ----
2350
2351    #[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        // filtered_count=Some(0) means "K flag set, but 0 filtered"
2422        // (= everything really removed). None = trailer completely missing.
2423        // Both must round-trip.
2424        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        // body ends before filtered_count → UnexpectedEof
2445        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        // Decoder with filtered_count_flag=true expects a 4-byte trailer
2459        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    // ---- FragmentNumberSet ----
2483
2484    #[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        // Regression M-4 / DDSI-RTPS §8.3.5.4: numBits MUST be <= 256. A
2514        // gap over > 256 fragments (e.g. fragment 1 AND 300 missing at
2515        // ~781 fragments under packet loss) must not build the set with
2516        // num_bits=300 — a spec-conformant receiver discards that as
2517        // malformed → the NACK_FRAG is lost → fragments are never resent
2518        // → sample stall. The set covers only the first 256; the rest
2519        // follows in the next NACK_FRAG once bitmap_base has advanced.
2520        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        // Fragment 1 (within the 256 window) is set, 300 (outside) is
2529        // not — it is re-requested in a follow-up NACK_FRAG.
2530        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    // ---- DATA_FRAG Submessage ----
2582
2583    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        // sample_size=10, fragment_size=4, fragment 3 carries only 2 bytes
2632        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        // B3 research: Cyclone/Fast-DDS ignore non-zero extra_flags.
2650        // So do we.
2651        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()); // extra_flags nonzero
2654        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        // B7: hard cap against DoS via a huge bitmap.
2662        let mut buf = Vec::new();
2663        buf.extend_from_slice(&SequenceNumber(1).to_bytes_le());
2664        buf.extend_from_slice(&257u32.to_le_bytes()); // num_bits > 256
2665        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        // 256 bits = 8 words = 32 byte bitmap
2675        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()); // num_bits far above 256
2685        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        // Pathological: num_bits = u32::MAX would allocate ~512 MB if we
2692        // do not cap beforehand.
2693        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        // We craft a DATA_FRAG body with a wrong octetsToInlineQos=99 and
2703        // Q=false. The decoder must reject it.
2704        let d = dataflag_frag(1, 1, 1, 4, 4, vec![1, 2, 3, 4]);
2705        let (mut bytes, _) = d.write_body(true);
2706        // octetsToInlineQos sits in bytes [2..4] (after extra_flags).
2707        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        // Q-flag true is rejected (feature not implemented).
2715        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    // ---- HEARTBEAT_FRAG Submessage ----
2738
2739    #[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    // ---- NACK_FRAG Submessage ----
2776
2777    #[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    // ---- WP 1.E stage E: InfoSource ----
2821
2822    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        // unused MUST roundtrip byte-identisch — manche Vendor-Implementations
2866        // enter diagnostic data there.
2867        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    // ---- WP 1.E stage F: InfoReply ----
2875
2876    #[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        // The spec allows an empty list (e.g. "forget all previous
2920        // reply locators"). The decoder must accept it.
2921        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        // length=u32::MAX → DoS-Schutz greift
2934        let mut body = Vec::new();
2935        body.extend_from_slice(&u32::MAX.to_le_bytes());
2936        // no locator body
2937        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    // ---- §8.3.7.5 / §8.3.8.5 InfoTimestamp ----
2948
2949    #[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        // I-flag → the body is ignored; even when it is full,
3002        // invalidate=true holds.
3003        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}