use alloc::collections::BTreeMap;
use alloc::vec::Vec;
use core::time::Duration;
use crate::arq::seq;
#[allow(dead_code)]
const MAX_TIMESTAMP: u64 = 0xFFFF_FFFF;
const TSBPD_DELAY_MIN_MS: u64 = 120;
#[derive(Debug, Clone, PartialEq, Eq, Default)]
#[non_exhaustive]
pub struct TickOutcome {
pub delivered: Vec<u32>,
pub dropped: Vec<u32>,
}
#[derive(Debug)]
pub struct TsbpdScheduler {
tsbpd_time_base: u64,
tsbpd_delay_ms: u64,
drift_us: u64,
tlpktdrop_threshold_us: u64,
tlpktdrop_enabled: bool,
next_release: u32,
buffer: BTreeMap<u32, u64>,
highest_fed: Option<u32>,
}
impl TsbpdScheduler {
#[allow(clippy::too_many_arguments)]
pub fn new(
initial_seq: u32,
tsbpd_time_base: u64,
tsbpd_delay_ms: u64,
drift_us: u64,
tlpktdrop_enabled: bool,
tlpktdrop_threshold_us: Option<u64>,
) -> Self {
let tsbpd_delay_ms = tsbpd_delay_ms.max(TSBPD_DELAY_MIN_MS);
let tlpktdrop_threshold_us = tlpktdrop_threshold_us.unwrap_or_else(|| {
(tsbpd_delay_ms * 5).div_ceil(4) * 1000
});
TsbpdScheduler {
tsbpd_time_base,
tsbpd_delay_ms,
drift_us,
tlpktdrop_threshold_us,
tlpktdrop_enabled,
next_release: initial_seq,
buffer: BTreeMap::new(),
highest_fed: None,
}
}
fn pkt_tsbpd_time(&self, pkt_timestamp: u32) -> u64 {
self.tsbpd_time_base + u64::from(pkt_timestamp) + self.tsbpd_delay_ms * 1000 + self.drift_us
}
pub fn feed_data(&mut self, seq_number: u32, pkt_timestamp: u32, now: Duration) -> TickOutcome {
let now_us = now.as_micros() as u64;
let pkt_tsbpd_time = self.pkt_tsbpd_time(pkt_timestamp);
match self.highest_fed {
None => self.highest_fed = Some(seq_number),
Some(h) if seq::seq_gt(seq_number, h) => self.highest_fed = Some(seq_number),
_ => {}
}
if self.tlpktdrop_enabled {
let drop_before = now_us.saturating_sub(self.tlpktdrop_threshold_us);
if pkt_tsbpd_time < drop_before {
let mut outcome = TickOutcome {
dropped: alloc::vec![seq_number],
..TickOutcome::default()
};
if seq_number == self.next_release {
self.next_release = seq::seq_next(seq_number);
}
self.buffer.remove(&seq_number);
outcome.delivered = self.release_ready(now_us);
return outcome;
}
}
if !seq::seq_lt(seq_number, self.next_release) {
self.buffer.insert(seq_number, pkt_tsbpd_time);
}
TickOutcome {
delivered: self.release_ready(now_us),
dropped: Vec::new(),
}
}
pub fn tick(&mut self, now: Duration) -> TickOutcome {
let now_us = now.as_micros() as u64;
TickOutcome {
delivered: self.release_ready(now_us),
dropped: Vec::new(), }
}
fn release_ready(&mut self, now_us: u64) -> Vec<u32> {
let mut delivered = Vec::new();
loop {
if !self.buffer.contains_key(&self.next_release) {
break; }
let tsbpd_time = self.buffer[&self.next_release];
if now_us < tsbpd_time {
break; }
if self.tlpktdrop_enabled {
let drop_before = now_us.saturating_sub(self.tlpktdrop_threshold_us);
if tsbpd_time < drop_before {
self.buffer.remove(&self.next_release);
self.next_release = seq::seq_next(self.next_release);
continue;
}
}
self.buffer.remove(&self.next_release);
delivered.push(self.next_release);
self.next_release = seq::seq_next(self.next_release);
}
delivered
}
pub fn next_release(&self) -> u32 {
self.next_release
}
pub fn buffered_count(&self) -> usize {
self.buffer.len()
}
pub fn has_gap(&self) -> bool {
!self.buffer.contains_key(&self.next_release)
}
}
#[cfg(test)]
mod tests {
use super::*;
use alloc::vec;
use core::time::Duration;
const TIME_BASE: u64 = 1_000_000; const DELAY_MS: u64 = 120; const ISN: u32 = 0;
fn sched() -> TsbpdScheduler {
TsbpdScheduler::new(ISN, TIME_BASE, DELAY_MS, 0, true, None)
}
#[test]
fn pkt_tsbpd_time_formula() {
let s = sched();
let ts = 5000u32;
let expected = TIME_BASE + u64::from(ts) + DELAY_MS * 1000;
assert_eq!(s.pkt_tsbpd_time(ts), expected);
}
#[test]
fn in_order_after_delay() {
let mut s = sched();
let ts_a = 0u32;
let ts_b = 10_000u32;
let outcome = s.feed_data(0, ts_a, Duration::ZERO);
assert!(outcome.delivered.is_empty());
assert!(outcome.dropped.is_empty());
assert_eq!(s.buffered_count(), 1);
let outcome = s.feed_data(1, ts_b, Duration::ZERO);
assert!(outcome.delivered.is_empty());
assert!(outcome.dropped.is_empty());
assert_eq!(s.buffered_count(), 2);
let pkt1_tsbpd = TIME_BASE + u64::from(ts_a) + DELAY_MS * 1000;
let outcome = s.tick(Duration::from_micros(pkt1_tsbpd));
assert_eq!(outcome.delivered, vec![0]);
let pkt2_tsbpd = TIME_BASE + u64::from(ts_b) + DELAY_MS * 1000;
let outcome = s.tick(Duration::from_micros(pkt2_tsbpd));
assert_eq!(outcome.delivered, vec![1]);
assert!(s.buffer.is_empty());
}
#[test]
fn out_of_order_arrival() {
let mut s = sched();
let ts = 10_000u32;
let pkt1_tsbpd = TIME_BASE + 10_000 + DELAY_MS * 1000;
let outcome = s.feed_data(1, ts, Duration::from_micros(pkt1_tsbpd));
assert!(outcome.delivered.is_empty());
assert_eq!(s.buffered_count(), 1);
let outcome = s.feed_data(0, 0, Duration::from_micros(pkt1_tsbpd));
assert_eq!(outcome.delivered, vec![0, 1]);
assert!(s.buffer.is_empty());
}
#[test]
fn too_late_drop_on_arrival() {
let mut s = sched();
let pkt_tsbpd = TIME_BASE + DELAY_MS * 1000;
let threshold_us = (DELAY_MS * 5).div_ceil(4) * 1000;
let very_late_now = Duration::from_micros(pkt_tsbpd + threshold_us + 1);
let outcome = s.feed_data(0, 0, very_late_now);
assert!(
outcome.delivered.is_empty(),
"should not deliver late packet"
);
assert_eq!(outcome.dropped, vec![0]);
}
#[test]
fn too_late_drop_buffered() {
let mut s = TsbpdScheduler::new(0, TIME_BASE, DELAY_MS, 0, true, None);
s.feed_data(
1,
10_000,
Duration::from_micros(TIME_BASE + 10_000 + DELAY_MS * 1000),
);
s.feed_data(
2,
20_000,
Duration::from_micros(TIME_BASE + 20_000 + DELAY_MS * 1000),
);
assert_eq!(s.buffered_count(), 2);
let pkt1_tsbpd = TIME_BASE + 10_000 + DELAY_MS * 1000;
let threshold = (DELAY_MS * 5).div_ceil(4) * 1000;
let far_future = Duration::from_micros(pkt1_tsbpd + threshold + 1);
let outcome = s.tick(far_future);
assert!(outcome.delivered.is_empty());
assert!(outcome.dropped.is_empty());
let pkt0_tsbpd = TIME_BASE + DELAY_MS * 1000;
let very_far = Duration::from_micros(pkt0_tsbpd + threshold * 10);
let outcome = s.feed_data(0, 0, very_far);
assert!(outcome.delivered.is_empty());
assert_eq!(s.next_release(), 3);
assert_eq!(s.buffered_count(), 0);
}
#[test]
fn sequence_order_preserved() {
let mut s = sched();
let ts_step = 10_000u32;
let num_packets = 5;
for i in 0..num_packets {
s.feed_data(i, ts_step * i, Duration::ZERO);
}
assert_eq!(s.buffered_count(), num_packets as usize);
let last_tsbpd =
TIME_BASE + u64::from(ts_step) * (num_packets - 1) as u64 + DELAY_MS * 1000;
let outcome = s.tick(Duration::from_micros(last_tsbpd));
assert_eq!(outcome.delivered, (0..num_packets).collect::<Vec<_>>());
}
#[test]
fn timestamp_wrap_smoke() {
let mut s = TsbpdScheduler::new(0, TIME_BASE, DELAY_MS, 0, false, None);
let near_wrap: u32 = 0xFFFF_FF00u32;
let past_wrap: u32 = 500;
let outcome = s.feed_data(0, near_wrap, Duration::ZERO);
assert!(outcome.delivered.is_empty());
let outcome = s.feed_data(1, past_wrap, Duration::ZERO);
assert!(outcome.delivered.is_empty());
assert_eq!(s.buffered_count(), 2);
let pkt0_tsbpd = TIME_BASE + u64::from(near_wrap) + DELAY_MS * 1000;
let pkt1_tsbpd = TIME_BASE + u64::from(past_wrap) + DELAY_MS * 1000;
assert!(pkt0_tsbpd > pkt1_tsbpd);
let outcome = s.tick(Duration::from_micros(pkt0_tsbpd));
assert_eq!(outcome.delivered, vec![0, 1]);
}
#[test]
fn minimum_delay_floor_applied() {
let s = TsbpdScheduler::new(0, TIME_BASE, 10, 0, false, None);
let ts = 0u32;
let expected = TIME_BASE + u64::from(ts) + TSBPD_DELAY_MIN_MS * 1000;
assert_eq!(s.pkt_tsbpd_time(ts), expected);
}
#[test]
fn tlpktdrop_disabled_never_drops() {
let mut s = TsbpdScheduler::new(0, TIME_BASE, DELAY_MS, 0, false, None);
let ts = 0u32;
let very_late_now =
Duration::from_micros(TIME_BASE + u64::from(ts) + DELAY_MS * 1000 + 1_000_000);
let outcome = s.feed_data(0, ts, very_late_now);
assert_eq!(outcome.delivered, vec![0]);
assert!(outcome.dropped.is_empty());
}
#[test]
fn gap_blocks_delivery() {
let mut s = sched();
s.feed_data(
1,
10_000,
Duration::from_micros(TIME_BASE + 10_000 + DELAY_MS * 1000),
);
s.feed_data(
2,
20_000,
Duration::from_micros(TIME_BASE + 20_000 + DELAY_MS * 1000),
);
assert!(s.has_gap());
assert_eq!(s.buffered_count(), 2);
let outcome = s.tick(Duration::from_micros(TIME_BASE + 100_000 + DELAY_MS * 1000));
assert!(outcome.delivered.is_empty());
assert!(s.has_gap());
}
#[test]
fn duplicate_arrival_does_not_advance_clock() {
let mut s = sched();
let ts = 0u32;
let now = Duration::from_micros(TIME_BASE + DELAY_MS * 1000);
s.feed_data(0, ts, now);
assert_eq!(s.buffered_count(), 0);
s.feed_data(0, ts, now);
assert_eq!(s.buffered_count(), 0, "duplicate must not re-buffer");
assert_eq!(s.next_release(), 1);
}
}