Skip to main content

prns_runtime/manifold/
announce_pacer.rs

1use crate::engine::InstantMillis;
2use crate::interfaces::{AnnounceBandwidthCap, BitrateBps};
3use crate::wire::BROADCAST_MTU;
4use core::cmp::Reverse;
5use heapless::Vec as HeaplessVec;
6
7const QUEUED_ANNOUNCE_LIFE_MS: u64 = 24 * 60 * 60 * 1_000;
8
9#[derive(Debug, Clone, Copy, PartialEq, Eq)]
10pub enum PacerReject {
11    FrameTooLarge,
12    QueueFull,
13}
14
15#[derive(Debug, Clone, Copy, PartialEq, Eq)]
16pub enum PacerOffer {
17    Sent,
18    Queued,
19    Rejected(PacerReject),
20}
21
22#[derive(Debug, Clone, Copy, PartialEq, Eq)]
23pub enum PacerRelease {
24    Released,
25    NotDue,
26    Idle,
27}
28
29pub trait PacerQueue<M = ()>: Default {
30    fn insert(
31        &mut self,
32        bytes: &[u8],
33        hops: u8,
34        now: InstantMillis,
35        metadata: M,
36    ) -> Result<(), PacerReject>;
37    fn take_next_with<R>(&mut self, f: impl FnOnce(&[u8], M) -> R) -> Option<R>;
38    fn evict_stale(&mut self, now: InstantMillis, life_ms: u64);
39    fn is_empty(&self) -> bool;
40    fn len(&self) -> usize;
41
42    fn clear(&mut self) -> usize {
43        let mut removed = 0;
44        while self.take_next_with(|_, _| ()).is_some() {
45            removed += 1;
46        }
47        removed
48    }
49}
50
51struct Queued<F, M> {
52    hops: u8,
53    queued_at: InstantMillis,
54    frame: F,
55    metadata: M,
56}
57
58pub struct FixedPacerQueue<const DEPTH: usize, M = ()> {
59    entries: HeaplessVec<Queued<HeaplessVec<u8, BROADCAST_MTU>, M>, DEPTH>,
60}
61
62impl<const DEPTH: usize, M> Default for FixedPacerQueue<DEPTH, M> {
63    fn default() -> Self {
64        Self {
65            entries: HeaplessVec::new(),
66        }
67    }
68}
69
70impl<const DEPTH: usize, M: Copy> PacerQueue<M> for FixedPacerQueue<DEPTH, M> {
71    fn insert(
72        &mut self,
73        bytes: &[u8],
74        hops: u8,
75        now: InstantMillis,
76        metadata: M,
77    ) -> Result<(), PacerReject> {
78        let mut frame = HeaplessVec::new();
79        if frame.extend_from_slice(bytes).is_err() {
80            return Err(PacerReject::FrameTooLarge);
81        }
82        if self.entries.is_full() {
83            match self
84                .entries
85                .iter()
86                .enumerate()
87                .max_by_key(|(_, entry)| (entry.hops, Reverse(entry.queued_at.0)))
88                .map(|(index, entry)| (index, entry.hops))
89            {
90                Some((index, worst_hops)) if hops < worst_hops => {
91                    self.entries.swap_remove(index);
92                }
93                _ => return Err(PacerReject::QueueFull),
94            }
95        }
96        self.entries
97            .push(Queued {
98                hops,
99                queued_at: now,
100                frame,
101                metadata,
102            })
103            .map_err(|_| PacerReject::QueueFull)
104    }
105
106    fn take_next_with<R>(&mut self, f: impl FnOnce(&[u8], M) -> R) -> Option<R> {
107        let index = self
108            .entries
109            .iter()
110            .enumerate()
111            .min_by_key(|(_, entry)| (entry.hops, entry.queued_at.0))
112            .map(|(index, _)| index)?;
113        let entry = self.entries.swap_remove(index);
114        Some(f(entry.frame.as_slice(), entry.metadata))
115    }
116
117    fn evict_stale(&mut self, now: InstantMillis, life_ms: u64) {
118        let mut index = 0;
119        while index < self.entries.len() {
120            if now.0.saturating_sub(self.entries[index].queued_at.0) > life_ms {
121                self.entries.swap_remove(index);
122            } else {
123                index += 1;
124            }
125        }
126    }
127
128    fn is_empty(&self) -> bool {
129        self.entries.is_empty()
130    }
131
132    fn len(&self) -> usize {
133        self.entries.len()
134    }
135}
136
137#[cfg(feature = "alloc")]
138pub use heap::HeapPacerQueue;
139
140#[cfg(feature = "alloc")]
141mod heap {
142    use super::{PacerQueue, PacerReject, Queued};
143    use crate::engine::InstantMillis;
144    use alloc::vec::Vec;
145
146    pub struct HeapPacerQueue<M = ()> {
147        entries: Vec<Queued<Vec<u8>, M>>,
148    }
149
150    impl<M> Default for HeapPacerQueue<M> {
151        fn default() -> Self {
152            Self {
153                entries: Vec::new(),
154            }
155        }
156    }
157
158    impl<M: Copy> PacerQueue<M> for HeapPacerQueue<M> {
159        fn insert(
160            &mut self,
161            bytes: &[u8],
162            hops: u8,
163            now: InstantMillis,
164            metadata: M,
165        ) -> Result<(), PacerReject> {
166            self.entries.push(Queued {
167                hops,
168                queued_at: now,
169                frame: bytes.to_vec(),
170                metadata,
171            });
172            Ok(())
173        }
174
175        fn take_next_with<R>(&mut self, f: impl FnOnce(&[u8], M) -> R) -> Option<R> {
176            let index = self
177                .entries
178                .iter()
179                .enumerate()
180                .min_by_key(|(_, entry)| (entry.hops, entry.queued_at.0))
181                .map(|(index, _)| index)?;
182            let entry = self.entries.swap_remove(index);
183            Some(f(&entry.frame, entry.metadata))
184        }
185
186        fn evict_stale(&mut self, now: InstantMillis, life_ms: u64) {
187            self.entries
188                .retain(|entry| now.0.saturating_sub(entry.queued_at.0) <= life_ms);
189        }
190
191        fn is_empty(&self) -> bool {
192            self.entries.is_empty()
193        }
194
195        fn len(&self) -> usize {
196            self.entries.len()
197        }
198    }
199}
200
201pub struct AnnouncePacer<Q, M = ()>
202where
203    Q: PacerQueue<M>,
204{
205    cap: AnnounceBandwidthCap,
206    bitrate: BitrateBps,
207    allowed_at: InstantMillis,
208    queue: Q,
209    metadata: core::marker::PhantomData<fn(M)>,
210}
211
212impl<Q, M> AnnouncePacer<Q, M>
213where
214    Q: PacerQueue<M>,
215    M: Copy,
216{
217    pub fn new(cap: AnnounceBandwidthCap, bitrate: BitrateBps) -> Self {
218        let allowed_at = match cap {
219            AnnounceBandwidthCap::Limited { cap_per_mille: 0 } => InstantMillis(u64::MAX),
220            AnnounceBandwidthCap::Unlimited | AnnounceBandwidthCap::Limited { .. } => {
221                InstantMillis(0)
222            }
223        };
224        Self {
225            cap,
226            bitrate,
227            allowed_at,
228            queue: Q::default(),
229            metadata: core::marker::PhantomData,
230        }
231    }
232
233    pub fn offer_tagged(
234        &mut self,
235        bytes: &[u8],
236        hops: u8,
237        now: InstantMillis,
238        metadata: M,
239        send: impl FnOnce(&[u8], M),
240    ) -> PacerOffer {
241        self.queue.evict_stale(now, QUEUED_ANNOUNCE_LIFE_MS);
242        if self.queue.is_empty() && self.allowed_at.0 <= now.0 {
243            send(bytes, metadata);
244            self.allowed_at = InstantMillis(
245                now.0
246                    .saturating_add(self.cap.cooldown_after_send_ms(self.bitrate, bytes.len())),
247            );
248            PacerOffer::Sent
249        } else {
250            match self.queue.insert(bytes, hops, now, metadata) {
251                Ok(()) => PacerOffer::Queued,
252                Err(reason) => PacerOffer::Rejected(reason),
253            }
254        }
255    }
256
257    pub fn release_due_tagged(
258        &mut self,
259        now: InstantMillis,
260        send: impl FnOnce(&[u8], M),
261    ) -> PacerRelease {
262        if self.allowed_at.0 > now.0 {
263            return PacerRelease::NotDue;
264        }
265        self.queue.evict_stale(now, QUEUED_ANNOUNCE_LIFE_MS);
266        let cap = self.cap;
267        let bitrate = self.bitrate;
268        match self.queue.take_next_with(|bytes, metadata| {
269            send(bytes, metadata);
270            cap.cooldown_after_send_ms(bitrate, bytes.len())
271        }) {
272            Some(spacing) => {
273                self.allowed_at = InstantMillis(now.0.saturating_add(spacing));
274                PacerRelease::Released
275            }
276            None => PacerRelease::Idle,
277        }
278    }
279
280    pub fn next_release(&self) -> Option<InstantMillis> {
281        (!self.queue.is_empty() && !self.cap.blocks_all()).then_some(self.allowed_at)
282    }
283
284    pub fn is_idle(&self) -> bool {
285        self.queue.is_empty()
286    }
287
288    pub fn queued_len(&self) -> usize {
289        self.queue.len()
290    }
291
292    pub fn clear_queue(&mut self) -> usize {
293        self.queue.clear()
294    }
295}
296
297impl<Q> AnnouncePacer<Q>
298where
299    Q: PacerQueue<()>,
300{
301    pub fn offer(
302        &mut self,
303        bytes: &[u8],
304        hops: u8,
305        now: InstantMillis,
306        send: impl FnOnce(&[u8]),
307    ) -> PacerOffer {
308        self.offer_tagged(bytes, hops, now, (), |frame, ()| send(frame))
309    }
310
311    pub fn release_due(&mut self, now: InstantMillis, send: impl FnOnce(&[u8])) -> PacerRelease {
312        self.release_due_tagged(now, |frame, ()| send(frame))
313    }
314}
315
316#[cfg(test)]
317mod tests {
318    use super::*;
319
320    const SLOW: AnnounceBandwidthCap = AnnounceBandwidthCap::RNS_DEFAULT;
321    const SLOW_BITRATE: BitrateBps = BitrateBps::guess(5_000);
322    const SPACING_MS: u64 = 800;
323
324    fn frame(tag: u8) -> [u8; 10] {
325        [tag; 10]
326    }
327
328    fn capture() -> std::vec::Vec<std::vec::Vec<u8>> {
329        std::vec::Vec::new()
330    }
331
332    #[test]
333    fn an_unlimited_link_emits_immediately_and_never_queues() {
334        let mut pacer =
335            AnnouncePacer::<FixedPacerQueue<4>>::new(AnnounceBandwidthCap::Unlimited, SLOW_BITRATE);
336        let mut sent = capture();
337        for at in [0, 1, 2, 3] {
338            pacer.offer(&frame(at as u8), 1, InstantMillis(at), |b| {
339                sent.push(b.to_vec())
340            });
341        }
342        assert_eq!(sent.len(), 4);
343        assert!(pacer.is_idle());
344        assert_eq!(pacer.next_release(), None);
345    }
346
347    #[test]
348    fn an_idle_pacer_emits_the_first_announce_now() {
349        let mut pacer = AnnouncePacer::<FixedPacerQueue<4>>::new(SLOW, SLOW_BITRATE);
350        let mut sent = capture();
351        pacer.offer(&frame(0), 1, InstantMillis(1_000), |b| {
352            sent.push(b.to_vec())
353        });
354        assert_eq!(sent.len(), 1);
355        assert_eq!(pacer.next_release(), None);
356    }
357
358    #[test]
359    fn a_zero_cap_queues_without_emitting() {
360        let mut pacer = AnnouncePacer::<FixedPacerQueue<4>>::new(
361            AnnounceBandwidthCap::Limited { cap_per_mille: 0 },
362            SLOW_BITRATE,
363        );
364        let mut sent = capture();
365        assert_eq!(
366            pacer.offer(&frame(0), 1, InstantMillis(0), |bytes| {
367                sent.push(bytes.to_vec())
368            }),
369            PacerOffer::Queued
370        );
371        assert!(sent.is_empty());
372        assert_eq!(pacer.next_release(), None);
373    }
374
375    #[test]
376    fn a_second_announce_within_the_window_queues() {
377        let mut pacer = AnnouncePacer::<FixedPacerQueue<4>>::new(SLOW, SLOW_BITRATE);
378        let mut sent = capture();
379        pacer.offer(&frame(0), 1, InstantMillis(1_000), |b| {
380            sent.push(b.to_vec())
381        });
382        pacer.offer(&frame(1), 1, InstantMillis(1_500), |b| {
383            sent.push(b.to_vec())
384        });
385        assert_eq!(sent.len(), 1, "the second is held, not emitted");
386        assert_eq!(
387            pacer.next_release(),
388            Some(InstantMillis(1_000 + SPACING_MS))
389        );
390    }
391
392    #[test]
393    fn the_queue_releases_lowest_hops_first() {
394        let mut pacer = AnnouncePacer::<FixedPacerQueue<8>>::new(SLOW, SLOW_BITRATE);
395        let mut sent = capture();
396        pacer.offer(&frame(9), 9, InstantMillis(0), |b| sent.push(b.to_vec()));
397        pacer.offer(&frame(5), 5, InstantMillis(0), |b| sent.push(b.to_vec()));
398        pacer.offer(&frame(1), 1, InstantMillis(0), |b| sent.push(b.to_vec()));
399        pacer.offer(&frame(3), 3, InstantMillis(0), |b| sent.push(b.to_vec()));
400        assert_eq!(sent, std::vec![frame(9).to_vec()], "hops-9 went out idle");
401
402        let mut now = 0;
403        for expected in [frame(1), frame(3), frame(5)] {
404            now += SPACING_MS;
405            assert_eq!(
406                pacer.release_due(InstantMillis(now), |b| sent.push(b.to_vec())),
407                PacerRelease::Released
408            );
409            assert_eq!(*sent.last().unwrap(), expected.to_vec());
410        }
411        assert!(pacer.is_idle());
412    }
413
414    #[test]
415    fn a_burst_drains_one_per_spacing_interval() {
416        let mut pacer = AnnouncePacer::<FixedPacerQueue<8>>::new(SLOW, SLOW_BITRATE);
417        let mut sent = capture();
418        for n in 0..4 {
419            pacer.offer(&frame(n), 1, InstantMillis(0), |b| sent.push(b.to_vec()));
420        }
421        assert_eq!(sent.len(), 1, "first goes now, the rest queue");
422
423        let mut now = 0;
424        for expected in 2..=4 {
425            now += SPACING_MS;
426            assert_eq!(
427                pacer.release_due(InstantMillis(now), |b| sent.push(b.to_vec())),
428                PacerRelease::Released
429            );
430            assert_eq!(sent.len(), expected);
431            assert_eq!(
432                pacer.release_due(InstantMillis(now), |b| sent.push(b.to_vec())),
433                PacerRelease::NotDue,
434                "only one releases per interval"
435            );
436        }
437        assert!(pacer.is_idle());
438    }
439
440    #[test]
441    fn a_full_fixed_queue_evicts_the_worst_hops() {
442        let mut pacer = AnnouncePacer::<FixedPacerQueue<2>>::new(SLOW, SLOW_BITRATE);
443        let mut sent = capture();
444        pacer.offer(&frame(5), 5, InstantMillis(0), |b| sent.push(b.to_vec()));
445        pacer.offer(&frame(5), 5, InstantMillis(0), |b| sent.push(b.to_vec()));
446        pacer.offer(&frame(5), 5, InstantMillis(0), |b| sent.push(b.to_vec()));
447        pacer.offer(&frame(1), 1, InstantMillis(0), |b| sent.push(b.to_vec()));
448        assert_eq!(
449            pacer.offer(&frame(9), 9, InstantMillis(0), |b| sent.push(b.to_vec())),
450            PacerOffer::Rejected(PacerReject::QueueFull),
451            "hops-9 is worse than every held announce, so the full gate rejects it",
452        );
453
454        let mut drained = capture();
455        let mut now = 0;
456        while !pacer.is_idle() {
457            now += SPACING_MS;
458            pacer.release_due(InstantMillis(now), |b| drained.push(b.to_vec()));
459        }
460        assert_eq!(
461            drained[0],
462            frame(1).to_vec(),
463            "the best-hops survivor goes first"
464        );
465        assert_eq!(drained[1], frame(5).to_vec());
466        assert!(
467            !drained.contains(&frame(9).to_vec()),
468            "the worse-than-queued hops-9 was dropped at the full gate"
469        );
470    }
471
472    #[cfg(feature = "alloc")]
473    #[test]
474    fn a_heap_queue_grows_without_dropping() {
475        let mut pacer = AnnouncePacer::<HeapPacerQueue>::new(SLOW, SLOW_BITRATE);
476        let mut sent = capture();
477        for n in 0..64u8 {
478            pacer.offer(&frame(n), 1, InstantMillis(0), |b| sent.push(b.to_vec()));
479        }
480        assert_eq!(sent.len(), 1, "first goes now, 63 queue and none drop");
481
482        let mut released = 0;
483        let mut now = 0;
484        while !pacer.is_idle() {
485            now += SPACING_MS;
486            if matches!(
487                pacer.release_due(InstantMillis(now), |_| {}),
488                PacerRelease::Released
489            ) {
490                released += 1;
491            }
492        }
493        assert_eq!(released, 63);
494    }
495
496    #[cfg(feature = "alloc")]
497    #[test]
498    fn clearing_a_queue_reports_every_removed_announce_without_resetting_cadence() {
499        let mut pacer = AnnouncePacer::<HeapPacerQueue>::new(SLOW, SLOW_BITRATE);
500        let mut sent = capture();
501        pacer.offer(&frame(0), 1, InstantMillis(0), |bytes| {
502            sent.push(bytes.to_vec())
503        });
504        pacer.offer(&frame(1), 1, InstantMillis(100), |bytes| {
505            sent.push(bytes.to_vec())
506        });
507        pacer.offer(&frame(2), 1, InstantMillis(200), |bytes| {
508            sent.push(bytes.to_vec())
509        });
510
511        assert_eq!(pacer.clear_queue(), 2);
512        assert_eq!(pacer.queued_len(), 0);
513        assert_eq!(pacer.next_release(), None);
514
515        pacer.offer(&frame(3), 1, InstantMillis(300), |bytes| {
516            sent.push(bytes.to_vec())
517        });
518        assert_eq!(sent, std::vec![frame(0).to_vec()]);
519        assert_eq!(pacer.queued_len(), 1);
520        assert_eq!(pacer.next_release(), Some(InstantMillis(SPACING_MS)));
521    }
522
523    #[test]
524    fn equal_hops_release_in_time_order_despite_internal_reordering() {
525        let mut pacer = AnnouncePacer::<FixedPacerQueue<8>>::new(SLOW, SLOW_BITRATE);
526        let mut sent = capture();
527        pacer.offer(&frame(0), 2, InstantMillis(0), |b| sent.push(b.to_vec()));
528        for (tag, queued_at) in [(1u8, 100u64), (2, 200), (3, 300), (4, 400)] {
529            pacer.offer(&frame(tag), 2, InstantMillis(queued_at), |b| {
530                sent.push(b.to_vec())
531            });
532        }
533        assert_eq!(
534            sent,
535            std::vec![frame(0).to_vec()],
536            "the first went out idle"
537        );
538
539        let mut now = 0;
540        for expected in [frame(1), frame(2), frame(3), frame(4)] {
541            now += SPACING_MS;
542            assert_eq!(
543                pacer.release_due(InstantMillis(now), |b| sent.push(b.to_vec())),
544                PacerRelease::Released
545            );
546            assert_eq!(
547                *sent.last().unwrap(),
548                expected.to_vec(),
549                "same-hops announces leave oldest-first even as swap_remove shuffles storage",
550            );
551        }
552        assert!(pacer.is_idle());
553    }
554
555    #[test]
556    fn a_stale_queued_announce_is_swept_and_a_fresh_one_sends() {
557        let mut pacer = AnnouncePacer::<FixedPacerQueue<8>>::new(SLOW, SLOW_BITRATE);
558        let mut sent = capture();
559        pacer.offer(&frame(0), 1, InstantMillis(0), |b| sent.push(b.to_vec()));
560        pacer.offer(&frame(1), 1, InstantMillis(400), |b| sent.push(b.to_vec()));
561        assert_eq!(sent, std::vec![frame(0).to_vec()], "the second is held");
562        assert!(!pacer.is_idle());
563
564        let long_after = 400 + QUEUED_ANNOUNCE_LIFE_MS + 1;
565        pacer.offer(&frame(2), 1, InstantMillis(long_after), |b| {
566            sent.push(b.to_vec())
567        });
568        assert_eq!(
569            sent,
570            std::vec![frame(0).to_vec(), frame(2).to_vec()],
571            "the day-old held announce was swept, never sent; the fresh one goes out",
572        );
573        assert!(pacer.is_idle());
574    }
575
576    #[test]
577    fn release_sweeps_a_stale_queue_and_sends_nothing() {
578        let mut pacer = AnnouncePacer::<FixedPacerQueue<8>>::new(SLOW, SLOW_BITRATE);
579        let mut sent = capture();
580        pacer.offer(&frame(0), 1, InstantMillis(0), |b| sent.push(b.to_vec()));
581        pacer.offer(&frame(1), 1, InstantMillis(400), |b| sent.push(b.to_vec()));
582
583        let long_after = 400 + QUEUED_ANNOUNCE_LIFE_MS + 1;
584        assert_eq!(
585            pacer.release_due(InstantMillis(long_after), |b| sent.push(b.to_vec())),
586            PacerRelease::Idle,
587            "the only held announce aged out, so the release finds nothing to send",
588        );
589        assert_eq!(sent, std::vec![frame(0).to_vec()]);
590        assert!(pacer.is_idle(), "the stale entry was swept from the queue");
591    }
592}