Skip to main content

zerodds_rtps/
writer_proxy.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright 2026 ZeroDDS Contributors
3//! `WriterProxy` — reader-side state over **one** remote writer.
4//!
5//! DDSI-RTPS 2.5 §8.4.6.5 (stateful reader behavior). The reader keeps
6//! a `WriterProxy` per matched writer, in which it tracks the range
7//! `[first_available_sn, last_available_sn]` from HEARTBEATs,
8//! marks already-received SNs and recognizes missing ones as **missing**.
9//! The missing set feeds the AckNack bitmap.
10//!
11//! A reader currently has one writer (single-writer
12//! assumption).
13
14extern crate alloc;
15use alloc::collections::BTreeSet;
16use alloc::vec::Vec;
17
18use crate::wire_types::{Guid, Locator, SequenceNumber};
19
20/// Reader-side state for one remote writer.
21#[derive(Debug, Clone)]
22pub struct WriterProxy {
23    /// GUID of the remote writer endpoint.
24    pub remote_writer_guid: Guid,
25    /// Unicast locators of the writer (e.g. for directed re-sends).
26    pub unicast_locators: Vec<Locator>,
27    /// Multicast locators.
28    pub multicast_locators: Vec<Locator>,
29    /// Reliable kind.
30    pub is_reliable: bool,
31    /// Smallest SN the writer **still** holds in the cache (from HEARTBEAT.first_sn).
32    first_available_sn: SequenceNumber,
33    /// Largest SN the writer announced (from HEARTBEAT.last_sn).
34    last_available_sn: SequenceNumber,
35    /// Highest SN this reader has actually **received**.
36    highest_received_sn: SequenceNumber,
37    /// Already-received SNs (for dup rejection + in-order delivery).
38    received: BTreeSet<SequenceNumber>,
39    /// SNs marked irrelevant by GAP submessages.
40    irrelevant: BTreeSet<SequenceNumber>,
41}
42
43impl WriterProxy {
44    /// Creates a fresh proxy.
45    #[must_use]
46    pub fn new(
47        remote_writer_guid: Guid,
48        unicast_locators: Vec<Locator>,
49        multicast_locators: Vec<Locator>,
50        is_reliable: bool,
51    ) -> Self {
52        Self {
53            remote_writer_guid,
54            unicast_locators,
55            multicast_locators,
56            is_reliable,
57            first_available_sn: SequenceNumber(1),
58            last_available_sn: SequenceNumber(0),
59            highest_received_sn: SequenceNumber(0),
60            received: BTreeSet::new(),
61            irrelevant: BTreeSet::new(),
62        }
63    }
64
65    /// Updates only the locators (re-discovery of the same writer), without
66    /// touching the reliability state (SN bounds, received/irrelevant).
67    /// A renewed SPDP/SEDP announce must not discard the reader progress
68    /// — otherwise the reader falsely reports "nothing missing" after a
69    /// HEARTBEAT and the writer never delivers the DATA.
70    pub fn refresh_locators(
71        &mut self,
72        unicast_locators: Vec<Locator>,
73        multicast_locators: Vec<Locator>,
74    ) {
75        self.unicast_locators = unicast_locators;
76        self.multicast_locators = multicast_locators;
77    }
78
79    /// Processes a HEARTBEAT.
80    ///
81    /// Per §8.4.15: `first_sn` is the smallest SN the writer
82    /// can re-deliver; `last_sn` the largest announced.
83    pub fn update_from_heartbeat(&mut self, first_sn: SequenceNumber, last_sn: SequenceNumber) {
84        // Monotonically growing bounds.
85        if first_sn > self.first_available_sn {
86            self.first_available_sn = first_sn;
87            // SNs **before** first_sn are lost and can no longer be
88            // requested — we drop them from received/irrelevant; they
89            // will also no longer be missing.
90            let split = self.received.split_off(&first_sn);
91            self.received = split;
92            let split = self.irrelevant.split_off(&first_sn);
93            self.irrelevant = split;
94        }
95        if last_sn > self.last_available_sn {
96            self.last_available_sn = last_sn;
97        }
98    }
99
100    /// Marks an SN as received.
101    pub fn received_change_set(&mut self, sn: SequenceNumber) {
102        if sn < self.first_available_sn {
103            // Lies before the announced range — ignore.
104            return;
105        }
106        self.received.insert(sn);
107        if sn > self.highest_received_sn {
108            self.highest_received_sn = sn;
109        }
110    }
111
112    /// Marks an SN as irrelevant (via GAP).
113    pub fn irrelevant_change_set(&mut self, sn: SequenceNumber) {
114        if sn < self.first_available_sn {
115            return;
116        }
117        self.irrelevant.insert(sn);
118    }
119
120    /// True if the SN is already received or marked irrelevant.
121    #[must_use]
122    pub fn is_known(&self, sn: SequenceNumber) -> bool {
123        self.received.contains(&sn) || self.irrelevant.contains(&sn)
124    }
125
126    /// Returns all **missing** SNs (neither received nor irrelevant) in the
127    /// range `[first_available_sn, last_available_sn]`.
128    ///
129    /// The vector is sorted ascending by SN. Limited to `max_count`
130    /// entries — the expected RTPS bitmap window is 256 SNs.
131    #[must_use]
132    pub fn missing_changes(&self, max_count: usize) -> Vec<SequenceNumber> {
133        let mut out = Vec::new();
134        if self.last_available_sn < self.first_available_sn {
135            return out;
136        }
137        let mut sn = self.first_available_sn;
138        while sn <= self.last_available_sn && out.len() < max_count {
139            if !self.is_known(sn) {
140                out.push(sn);
141            }
142            sn = SequenceNumber(sn.0 + 1);
143        }
144        out
145    }
146
147    /// True if there are missing SNs.
148    #[must_use]
149    pub fn has_missing_changes(&self) -> bool {
150        !self.missing_changes(1).is_empty()
151    }
152
153    /// Getter: smallest announced SN.
154    #[must_use]
155    pub fn first_available_sn(&self) -> SequenceNumber {
156        self.first_available_sn
157    }
158
159    /// Getter: largest announced SN.
160    #[must_use]
161    pub fn last_available_sn(&self) -> SequenceNumber {
162        self.last_available_sn
163    }
164
165    /// Getter: highest received SN.
166    #[must_use]
167    pub fn highest_received_sn(&self) -> SequenceNumber {
168        self.highest_received_sn
169    }
170
171    /// Matching AckNack base: smallest not-yet-acked SN.
172    ///
173    /// Convention: all SN < `acknack_base` are acked. We return
174    /// the smallest not-yet-received-or-irrelevant SN in `[first, last+1]`.
175    #[must_use]
176    pub fn acknack_base(&self) -> SequenceNumber {
177        let mut sn = self.first_available_sn;
178        while sn <= self.last_available_sn {
179            if !self.is_known(sn) {
180                return sn;
181            }
182            sn = SequenceNumber(sn.0 + 1);
183        }
184        SequenceNumber(self.last_available_sn.0 + 1)
185    }
186}
187
188#[cfg(test)]
189#[allow(clippy::expect_used, clippy::unwrap_used)]
190mod tests {
191    use super::*;
192    use crate::wire_types::{EntityId, GuidPrefix};
193
194    fn sn(n: i64) -> SequenceNumber {
195        SequenceNumber(n)
196    }
197
198    fn proxy() -> WriterProxy {
199        let guid = Guid::new(
200            GuidPrefix::from_bytes([2; 12]),
201            EntityId::user_writer_with_key([0x10, 0x20, 0x30]),
202        );
203        WriterProxy::new(guid, alloc::vec![], alloc::vec![], true)
204    }
205
206    #[test]
207    fn fresh_proxy_has_no_missing() {
208        let p = proxy();
209        assert!(!p.has_missing_changes());
210        assert_eq!(p.missing_changes(10), alloc::vec![]);
211        assert_eq!(p.acknack_base(), sn(1));
212    }
213
214    #[test]
215    fn heartbeat_sets_available_range() {
216        let mut p = proxy();
217        p.update_from_heartbeat(sn(1), sn(5));
218        assert_eq!(p.first_available_sn(), sn(1));
219        assert_eq!(p.last_available_sn(), sn(5));
220        // Nothing received yet → everything missing
221        assert_eq!(
222            p.missing_changes(10),
223            alloc::vec![sn(1), sn(2), sn(3), sn(4), sn(5)]
224        );
225    }
226
227    #[test]
228    fn received_removes_from_missing() {
229        let mut p = proxy();
230        p.update_from_heartbeat(sn(1), sn(5));
231        p.received_change_set(sn(2));
232        p.received_change_set(sn(4));
233        assert_eq!(p.missing_changes(10), alloc::vec![sn(1), sn(3), sn(5)]);
234        assert_eq!(p.acknack_base(), sn(1));
235    }
236
237    #[test]
238    fn gap_marks_irrelevant() {
239        let mut p = proxy();
240        p.update_from_heartbeat(sn(1), sn(5));
241        p.irrelevant_change_set(sn(3));
242        assert_eq!(
243            p.missing_changes(10),
244            alloc::vec![sn(1), sn(2), sn(4), sn(5)]
245        );
246    }
247
248    #[test]
249    fn acknack_base_walks_up() {
250        let mut p = proxy();
251        p.update_from_heartbeat(sn(1), sn(3));
252        p.received_change_set(sn(1));
253        p.received_change_set(sn(2));
254        assert_eq!(p.acknack_base(), sn(3));
255        p.received_change_set(sn(3));
256        assert_eq!(p.acknack_base(), sn(4));
257    }
258
259    #[test]
260    fn heartbeat_advancing_first_prunes_old_state() {
261        let mut p = proxy();
262        p.update_from_heartbeat(sn(1), sn(10));
263        p.received_change_set(sn(3));
264        p.received_change_set(sn(7));
265        // Writer rotates the cache → first now at 5
266        p.update_from_heartbeat(sn(5), sn(10));
267        assert_eq!(p.first_available_sn(), sn(5));
268        // sn(3) removed from received, sn(7) stays
269        assert!(!p.is_known(sn(3)));
270        assert!(p.is_known(sn(7)));
271    }
272
273    #[test]
274    fn highest_received_tracks_max() {
275        let mut p = proxy();
276        p.update_from_heartbeat(sn(1), sn(10));
277        p.received_change_set(sn(3));
278        p.received_change_set(sn(7));
279        p.received_change_set(sn(5));
280        assert_eq!(p.highest_received_sn(), sn(7));
281    }
282
283    #[test]
284    fn received_before_first_is_ignored() {
285        let mut p = proxy();
286        p.update_from_heartbeat(sn(5), sn(10));
287        p.received_change_set(sn(2));
288        assert!(!p.is_known(sn(2)));
289        assert_eq!(p.highest_received_sn(), sn(0));
290    }
291
292    #[test]
293    fn missing_changes_respects_max_count() {
294        let mut p = proxy();
295        p.update_from_heartbeat(sn(1), sn(100));
296        let m = p.missing_changes(5);
297        assert_eq!(m, alloc::vec![sn(1), sn(2), sn(3), sn(4), sn(5)]);
298    }
299
300    #[test]
301    fn acknack_base_when_all_received_is_last_plus_one() {
302        let mut p = proxy();
303        p.update_from_heartbeat(sn(1), sn(3));
304        p.received_change_set(sn(1));
305        p.received_change_set(sn(2));
306        p.received_change_set(sn(3));
307        assert_eq!(p.acknack_base(), sn(4));
308        assert!(!p.has_missing_changes());
309    }
310}