use std::time::{Duration, Instant};
pub mod constants {
use std::time::Duration;
pub const MIN_FRAME_INTERVAL_FLOOR: Duration = Duration::from_millis(20);
pub const COLLECTION_INTERVAL: Duration = Duration::from_millis(8);
pub const DELAYED_ACK_TIMEOUT: Duration = Duration::from_millis(100);
pub const MAX_FRAME_RATE_HZ: u32 = 50;
pub const KEEPALIVE_INTERVAL: Duration = Duration::from_secs(25);
pub const DEAD_INTERVAL: Duration = Duration::from_secs(60);
pub const MAX_RETRANSMITS: u32 = 10;
pub const RETRANSMIT_BACKOFF: u32 = 2;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SendReason {
StateChange,
Ack,
Keepalive,
Retransmit,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PacerAction {
SendNow,
WaitUntil(Instant),
Idle,
}
#[derive(Debug, Clone)]
pub struct FramePacer {
last_frame_sent: Option<Instant>,
state_change_time: Option<Instant>,
ack_pending_since: Option<Instant>,
data_pending: bool,
srtt_ms: f64,
}
impl Default for FramePacer {
fn default() -> Self {
Self::new()
}
}
impl FramePacer {
pub fn new() -> Self {
Self {
last_frame_sent: None,
state_change_time: None,
ack_pending_since: None,
data_pending: false,
srtt_ms: 0.0,
}
}
pub fn set_srtt(&mut self, srtt: Duration) {
self.srtt_ms = srtt.as_secs_f64() * 1000.0;
}
pub fn on_state_change(&mut self) {
if self.state_change_time.is_none() {
self.state_change_time = Some(Instant::now());
}
self.data_pending = true;
}
pub fn on_ack_needed(&mut self) {
if self.ack_pending_since.is_none() {
self.ack_pending_since = Some(Instant::now());
}
}
pub fn on_frame_sent(&mut self) {
self.last_frame_sent = Some(Instant::now());
self.state_change_time = None;
self.ack_pending_since = None;
self.data_pending = false;
}
pub fn clear_pending(&mut self) {
self.data_pending = false;
self.state_change_time = None;
}
fn min_frame_interval(&self) -> Duration {
let srtt_half_ms = self.srtt_ms / 2.0;
let floor_ms = constants::MIN_FRAME_INTERVAL_FLOOR.as_millis() as f64;
let interval_ms = f64::max(srtt_half_ms, floor_ms);
let max_interval_ms = 1000.0 / constants::MAX_FRAME_RATE_HZ as f64;
let interval_ms = f64::max(interval_ms, max_interval_ms);
Duration::from_secs_f64(interval_ms / 1000.0)
}
pub fn poll(&self) -> PacerAction {
let now = Instant::now();
let needs_send = self.data_pending || self.ack_pending_since.is_some();
if !needs_send {
return PacerAction::Idle;
}
if let Some(last_sent) = self.last_frame_sent {
let min_interval = self.min_frame_interval();
let next_allowed = last_sent + min_interval;
if now < next_allowed {
return PacerAction::WaitUntil(next_allowed);
}
}
if let Some(state_time) = self.state_change_time {
let collection_end = state_time + constants::COLLECTION_INTERVAL;
if now < collection_end && self.ack_pending_since.is_none() {
return PacerAction::WaitUntil(collection_end);
}
}
if !self.data_pending
&& let Some(ack_time) = self.ack_pending_since
{
let ack_deadline = ack_time + constants::DELAYED_ACK_TIMEOUT;
if now < ack_deadline {
return PacerAction::WaitUntil(ack_deadline);
}
}
PacerAction::SendNow
}
pub fn needs_keepalive(&self, last_received: Instant) -> bool {
if let Some(last_sent) = self.last_frame_sent {
let now = Instant::now();
let since_sent = now.duration_since(last_sent);
let since_received = now.duration_since(last_received);
since_sent >= constants::KEEPALIVE_INTERVAL
&& since_received < constants::DEAD_INTERVAL
} else {
false
}
}
pub fn is_connection_dead(&self, last_received: Instant) -> bool {
Instant::now().duration_since(last_received) >= constants::DEAD_INTERVAL
}
}
#[derive(Debug, Clone)]
pub struct RetransmitController {
retransmit_count: u32,
last_retransmit: Option<Instant>,
current_timeout: Duration,
base_rto: Duration,
}
impl RetransmitController {
pub fn new(initial_rto: Duration) -> Self {
Self {
retransmit_count: 0,
last_retransmit: None,
current_timeout: initial_rto,
base_rto: initial_rto,
}
}
pub fn set_rto(&mut self, rto: Duration) {
self.base_rto = rto;
if self.retransmit_count == 0 {
self.current_timeout = rto;
}
}
pub fn should_retransmit(&self, unacked_data: bool) -> bool {
if !unacked_data {
return false;
}
if self.retransmit_count >= constants::MAX_RETRANSMITS {
return false; }
match self.last_retransmit {
Some(last) => Instant::now().duration_since(last) >= self.current_timeout,
None => true, }
}
pub fn on_retransmit(&mut self) {
self.retransmit_count += 1;
self.last_retransmit = Some(Instant::now());
let new_timeout = self.current_timeout * constants::RETRANSMIT_BACKOFF;
self.current_timeout = new_timeout.min(super::timing::constants::MAX_RTO);
}
pub fn on_ack(&mut self) {
self.retransmit_count = 0;
self.last_retransmit = None;
self.current_timeout = self.base_rto;
}
pub fn retransmit_count(&self) -> u32 {
self.retransmit_count
}
pub fn is_failed(&self) -> bool {
self.retransmit_count >= constants::MAX_RETRANSMITS
}
pub fn time_until_retransmit(&self) -> Option<Duration> {
self.last_retransmit.map(|last| {
let elapsed = Instant::now().duration_since(last);
self.current_timeout.saturating_sub(elapsed)
})
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_pacer_initial_state() {
let pacer = FramePacer::new();
assert_eq!(pacer.poll(), PacerAction::Idle);
}
#[test]
fn test_pacer_state_change() {
let mut pacer = FramePacer::new();
pacer.on_state_change();
match pacer.poll() {
PacerAction::WaitUntil(_) => {}
other => panic!("Expected WaitUntil, got {:?}", other),
}
std::thread::sleep(constants::COLLECTION_INTERVAL + Duration::from_millis(1));
assert_eq!(pacer.poll(), PacerAction::SendNow);
}
#[test]
fn test_pacer_ack_only() {
let mut pacer = FramePacer::new();
pacer.on_ack_needed();
match pacer.poll() {
PacerAction::WaitUntil(_) => {}
other => panic!("Expected WaitUntil, got {:?}", other),
}
}
#[test]
fn test_pacer_ack_with_data() {
let mut pacer = FramePacer::new();
pacer.on_ack_needed();
pacer.on_state_change();
std::thread::sleep(constants::COLLECTION_INTERVAL + Duration::from_millis(1));
assert_eq!(pacer.poll(), PacerAction::SendNow);
}
#[test]
fn test_pacer_min_interval() {
let mut pacer = FramePacer::new();
pacer.set_srtt(Duration::from_millis(100));
let min_interval = pacer.min_frame_interval();
assert!(min_interval >= Duration::from_millis(50));
}
#[test]
fn test_pacer_frame_sent_clears_state() {
let mut pacer = FramePacer::new();
pacer.on_state_change();
pacer.on_ack_needed();
pacer.on_frame_sent();
assert_eq!(pacer.poll(), PacerAction::Idle);
}
#[test]
fn test_retransmit_controller() {
let mut controller = RetransmitController::new(Duration::from_millis(100));
assert!(controller.should_retransmit(true));
assert!(!controller.should_retransmit(false));
controller.on_retransmit();
assert!(!controller.should_retransmit(true));
controller.on_ack();
assert_eq!(controller.retransmit_count(), 0);
}
#[test]
fn test_retransmit_max_attempts() {
let mut controller = RetransmitController::new(Duration::from_millis(1));
for _ in 0..constants::MAX_RETRANSMITS {
controller.on_retransmit();
}
assert!(controller.is_failed());
assert!(!controller.should_retransmit(true));
}
#[test]
fn test_keepalive_check() {
let pacer = FramePacer::new();
assert!(!pacer.needs_keepalive(Instant::now()));
}
#[test]
fn test_connection_dead() {
let pacer = FramePacer::new();
assert!(!pacer.is_connection_dead(Instant::now()));
}
}