#![allow(dead_code)]
use rtc_interceptor::{AttributedPacket, Packet, TaggedPacket};
use shared::TransportContext;
use std::collections::VecDeque;
use std::time::{Duration, Instant};
#[derive(Debug, Clone, Copy)]
pub struct PathProfile {
pub propagation: Duration,
pub capacity_bits_per_second: f64,
pub queue_capacity_bits: f64,
pub drop_one_in: Option<u32>,
}
impl PathProfile {
pub fn steady() -> Self {
Self {
propagation: Duration::from_millis(20),
capacity_bits_per_second: 3_000_000.0,
queue_capacity_bits: 3_000_000.0,
drop_one_in: None,
}
}
pub fn queue_building() -> Self {
Self {
propagation: Duration::from_millis(20),
capacity_bits_per_second: 300_000.0,
queue_capacity_bits: 6_000_000.0,
drop_one_in: None,
}
}
pub fn lossy_without_queueing() -> Self {
Self {
propagation: Duration::from_millis(20),
capacity_bits_per_second: 3_000_000.0,
queue_capacity_bits: 3_000_000.0,
drop_one_in: Some(5),
}
}
pub fn mildly_lossy() -> Self {
Self {
drop_one_in: Some(20),
..Self::lossy_without_queueing()
}
}
pub fn recovering() -> Self {
Self {
propagation: Duration::from_millis(20),
capacity_bits_per_second: 600_000.0,
queue_capacity_bits: 6_000_000.0,
drop_one_in: None,
}
}
fn capacity_at(&self, elapsed: Duration, widens_after: Option<Duration>) -> f64 {
match widens_after {
Some(after) if elapsed >= after => self.capacity_bits_per_second * 5.0,
_ => self.capacity_bits_per_second,
}
}
}
#[derive(Debug, Clone, Copy)]
struct InFlight {
twcc_sequence_number: u16,
arrival: Instant,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Arrival {
pub twcc_sequence_number: u16,
pub at: Option<Instant>,
}
pub struct Path {
profile: PathProfile,
widens_after: Option<Duration>,
epoch: Instant,
queued_bits: f64,
drains_at: Instant,
in_flight: VecDeque<InFlight>,
offered: u32,
arrived: Vec<Arrival>,
}
impl Path {
pub fn new(profile: PathProfile, epoch: Instant) -> Self {
Self {
profile,
widens_after: None,
epoch,
queued_bits: 0.0,
drains_at: epoch,
in_flight: VecDeque::new(),
offered: 0,
arrived: Vec::new(),
}
}
pub fn widening_after(mut self, after: Duration) -> Self {
self.widens_after = Some(after);
self
}
pub fn offer(&mut self, now: Instant, twcc_sequence_number: u16, size_bits: f64) -> bool {
self.drain_to(now);
self.offered += 1;
if let Some(one_in) = self.profile.drop_one_in
&& self.offered.is_multiple_of(one_in)
{
self.arrived.push(Arrival {
twcc_sequence_number,
at: None,
});
return false;
}
if self.queued_bits + size_bits > self.profile.queue_capacity_bits {
self.arrived.push(Arrival {
twcc_sequence_number,
at: None,
});
return false;
}
let capacity = self
.profile
.capacity_at(now.saturating_duration_since(self.epoch), self.widens_after);
let service = Duration::from_secs_f64(size_bits / capacity);
let departs = self.drains_at.max(now) + service;
self.drains_at = departs;
self.queued_bits += size_bits;
self.in_flight.push_back(InFlight {
twcc_sequence_number,
arrival: departs + self.profile.propagation,
});
true
}
pub fn drain_to(&mut self, now: Instant) {
while let Some(head) = self.in_flight.front().copied() {
if head.arrival > now {
break;
}
self.in_flight.pop_front();
self.arrived.push(Arrival {
twcc_sequence_number: head.twcc_sequence_number,
at: Some(head.arrival),
});
}
if now >= self.drains_at {
self.queued_bits = 0.0;
}
}
pub fn take_arrivals(&mut self) -> Vec<Arrival> {
std::mem::take(&mut self.arrived)
}
pub fn queueing_delay(&self, sent_at: Instant, arrived_at: Instant) -> Duration {
arrived_at
.saturating_duration_since(sent_at)
.saturating_sub(self.profile.propagation)
}
}
pub fn twcc_feedback_for(
now: Instant,
epoch: Instant,
media_ssrc: u32,
arrivals: &[Arrival],
) -> Option<TaggedPacket> {
use rtcp::transport_feedbacks::transport_layer_cc::{
PacketStatusChunk, RecvDelta, RunLengthChunk, StatusChunkTypeTcc, SymbolTypeTcc,
TransportLayerCc,
};
let base = arrivals.first()?.twcc_sequence_number;
let mut packet_chunks = Vec::new();
let mut recv_deltas = Vec::new();
let mut previous: Option<Instant> = None;
for arrival in arrivals {
let symbol = match arrival.at {
Some(at) => {
let since = previous.map_or_else(
|| at.saturating_duration_since(epoch),
|last| at.saturating_duration_since(last),
);
previous = Some(at);
let ticks = (since.as_micros() / 250).min(255) as u16;
recv_deltas.push(RecvDelta {
type_tcc_packet: SymbolTypeTcc::PacketReceivedSmallDelta,
delta: i64::from(ticks) * 250,
});
SymbolTypeTcc::PacketReceivedSmallDelta
}
None => SymbolTypeTcc::PacketNotReceived,
};
packet_chunks.push(PacketStatusChunk::RunLengthChunk(RunLengthChunk {
type_tcc: StatusChunkTypeTcc::RunLengthChunk,
packet_status_symbol: symbol,
run_length: 1,
}));
}
let feedback = TransportLayerCc {
sender_ssrc: 0,
media_ssrc,
base_sequence_number: base,
packet_status_count: arrivals.len() as u16,
reference_time: (now.saturating_duration_since(epoch).as_millis() / 64) as u32,
fb_pkt_count: 0,
packet_chunks,
recv_deltas,
};
Some(TaggedPacket {
now,
transport: TransportContext::default(),
message: AttributedPacket::new(Packet::Rtcp(vec![Box::new(feedback)])),
})
}