zerodds_rtps/
writer_proxy.rs1extern crate alloc;
15use alloc::collections::BTreeSet;
16use alloc::vec::Vec;
17
18use crate::wire_types::{Guid, Locator, SequenceNumber};
19
20#[derive(Debug, Clone)]
22pub struct WriterProxy {
23 pub remote_writer_guid: Guid,
25 pub unicast_locators: Vec<Locator>,
27 pub multicast_locators: Vec<Locator>,
29 pub is_reliable: bool,
31 first_available_sn: SequenceNumber,
33 last_available_sn: SequenceNumber,
35 highest_received_sn: SequenceNumber,
37 received: BTreeSet<SequenceNumber>,
39 irrelevant: BTreeSet<SequenceNumber>,
41}
42
43impl WriterProxy {
44 #[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 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 pub fn update_from_heartbeat(&mut self, first_sn: SequenceNumber, last_sn: SequenceNumber) {
84 if first_sn > self.first_available_sn {
86 self.first_available_sn = first_sn;
87 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 pub fn received_change_set(&mut self, sn: SequenceNumber) {
102 if sn < self.first_available_sn {
103 return;
105 }
106 self.received.insert(sn);
107 if sn > self.highest_received_sn {
108 self.highest_received_sn = sn;
109 }
110 }
111
112 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 #[must_use]
122 pub fn is_known(&self, sn: SequenceNumber) -> bool {
123 self.received.contains(&sn) || self.irrelevant.contains(&sn)
124 }
125
126 #[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 #[must_use]
149 pub fn has_missing_changes(&self) -> bool {
150 !self.missing_changes(1).is_empty()
151 }
152
153 #[must_use]
155 pub fn first_available_sn(&self) -> SequenceNumber {
156 self.first_available_sn
157 }
158
159 #[must_use]
161 pub fn last_available_sn(&self) -> SequenceNumber {
162 self.last_available_sn
163 }
164
165 #[must_use]
167 pub fn highest_received_sn(&self) -> SequenceNumber {
168 self.highest_received_sn
169 }
170
171 #[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 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 p.update_from_heartbeat(sn(5), sn(10));
267 assert_eq!(p.first_available_sn(), sn(5));
268 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}