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}