use super::mode::AmrFrameType;
use super::payload::AmrPayloadFrame;
use crate::error::{CodecError, Result};
use std::collections::VecDeque;
pub const FRAME_BLOCK_MS: u16 = 20;
#[derive(Debug, Clone)]
pub struct RedundancyScheduler {
depth: usize,
history: VecDeque<AmrPayloadFrame>,
}
impl RedundancyScheduler {
pub fn new(max_red_ms: Option<u16>, requested_depth: usize) -> Result<Self> {
if requested_depth == 0 {
return Err(CodecError::invalid_config(
"redundancy depth is a frame count and must be at least 1",
));
}
if requested_depth > 32 {
return Err(CodecError::invalid_config(
"an AMR payload's table of contents addresses at most 32 frame-blocks",
));
}
let permitted = Self::permitted_depth(max_red_ms);
if requested_depth > permitted {
return Err(CodecError::invalid_config(format!(
"redundancy depth {requested_depth} exceeds the {permitted} frame-blocks \
the peer's max-red allows"
)));
}
Ok(Self {
depth: requested_depth,
history: VecDeque::new(),
})
}
#[must_use]
pub const fn permitted_depth(max_red_ms: Option<u16>) -> usize {
match max_red_ms {
None => 1,
Some(ms) => 1 + (ms / FRAME_BLOCK_MS) as usize,
}
}
#[must_use]
pub const fn depth(&self) -> usize {
self.depth
}
pub fn next_payload(&mut self, frame: AmrPayloadFrame) -> Vec<AmrPayloadFrame> {
let mut frames: Vec<AmrPayloadFrame> = self.history.iter().cloned().collect();
frames.push(frame.clone());
if matches!(
frame.frame_type,
AmrFrameType::Speech(_) | AmrFrameType::Sid(_)
) {
self.history.push_back(frame);
}
while self.history.len() >= self.depth {
self.history.pop_front();
}
frames
}
#[must_use]
pub fn payload_timestamp(
newest_timestamp: u32,
frame_count: usize,
samples_per_frame: u32,
) -> u32 {
let older = u32::try_from(frame_count.saturating_sub(1)).unwrap_or(u32::MAX);
newest_timestamp.wrapping_sub(older.wrapping_mul(samples_per_frame))
}
}
#[derive(Debug, Clone, Default)]
pub struct RedundancyDedup {
newest: Option<u32>,
}
impl RedundancyDedup {
#[must_use]
pub const fn new() -> Self {
Self { newest: None }
}
pub fn accept(
&mut self,
packet_timestamp: u32,
frame_count: usize,
samples_per_frame: u32,
) -> Vec<bool> {
let mut flags = Vec::with_capacity(frame_count);
for index in 0..frame_count {
let offset = u32::try_from(index).unwrap_or(u32::MAX);
let timestamp = packet_timestamp.wrapping_add(offset.wrapping_mul(samples_per_frame));
let is_new = self.newest.is_none_or(|newest| {
let ahead = timestamp.wrapping_sub(newest);
ahead != 0 && ahead < u32::MAX / 2
});
if is_new {
self.newest = Some(timestamp);
}
flags.push(is_new);
}
flags
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::codecs::amr::mode::{AmrMode, AmrVariant};
const NB: AmrVariant = AmrVariant::NarrowBand;
const NB_SAMPLES: u32 = 160;
fn speech(index: u8) -> AmrPayloadFrame {
let mode = AmrMode::new(NB, 7).expect("12.2 is a narrowband mode");
AmrPayloadFrame::new(
AmrFrameType::Speech(mode),
true,
vec![index; mode.octet_aligned_bytes()],
)
.expect("a full-length 12.2 frame")
}
#[test]
fn max_red_bounds_the_depth() {
assert_eq!(RedundancyScheduler::permitted_depth(Some(0)), 1);
assert_eq!(RedundancyScheduler::permitted_depth(Some(20)), 2);
assert_eq!(RedundancyScheduler::permitted_depth(Some(220)), 12);
assert_eq!(RedundancyScheduler::permitted_depth(None), 1);
assert!(RedundancyScheduler::new(Some(0), 2).is_err());
assert!(RedundancyScheduler::new(Some(20), 2).is_ok());
assert!(RedundancyScheduler::new(Some(20), 3).is_err());
assert!(RedundancyScheduler::new(Some(0), 0).is_err());
assert!(RedundancyScheduler::new(Some(10_000), 33).is_err());
}
#[test]
fn depth_one_sends_each_frame_exactly_once() {
let mut scheduler = RedundancyScheduler::new(Some(0), 1).expect("depth 1");
for index in 0..4u8 {
let payload = scheduler.next_payload(speech(index));
assert_eq!(payload.len(), 1, "depth 1 must never bundle");
assert_eq!(payload[0].data[0], index);
}
}
#[test]
fn deeper_payloads_carry_the_previous_frames_oldest_first() {
let mut scheduler = RedundancyScheduler::new(Some(40), 3).expect("depth 3");
let first = scheduler.next_payload(speech(0));
assert_eq!(first.len(), 1);
let second = scheduler.next_payload(speech(1));
assert_eq!(second.len(), 2);
assert_eq!(second.iter().map(|f| f.data[0]).collect::<Vec<_>>(), [0, 1]);
let third = scheduler.next_payload(speech(2));
assert_eq!(
third.iter().map(|f| f.data[0]).collect::<Vec<_>>(),
[0, 1, 2]
);
let fourth = scheduler.next_payload(speech(3));
assert_eq!(
fourth.iter().map(|f| f.data[0]).collect::<Vec<_>>(),
[1, 2, 3],
"the window must slide, not grow"
);
}
#[test]
fn the_payload_timestamp_names_the_oldest_frame() {
assert_eq!(
RedundancyScheduler::payload_timestamp(1_000, 3, NB_SAMPLES),
1_000 - 2 * NB_SAMPLES
);
assert_eq!(
RedundancyScheduler::payload_timestamp(1_000, 1, NB_SAMPLES),
1_000
);
}
#[test]
fn repeats_are_dropped_when_nothing_is_lost() {
let mut dedup = RedundancyDedup::new();
assert_eq!(dedup.accept(0, 1, NB_SAMPLES), [true]);
assert_eq!(dedup.accept(0, 2, NB_SAMPLES), [false, true]);
assert_eq!(dedup.accept(0, 3, NB_SAMPLES), [false, false, true]);
assert_eq!(
dedup.accept(NB_SAMPLES, 3, NB_SAMPLES),
[false, false, true]
);
}
#[test]
fn a_lost_packet_is_recovered_from_the_next_one() {
let mut dedup = RedundancyDedup::new();
assert_eq!(dedup.accept(0, 1, NB_SAMPLES), [true]);
let flags = dedup.accept(0, 3, NB_SAMPLES);
assert_eq!(flags, [false, true, true]);
}
#[test]
fn every_frame_survives_a_one_in_two_loss_pattern_at_depth_two() {
let mut dedup = RedundancyDedup::new();
let mut recovered = Vec::new();
for index in 0..20u32 {
let newest = index * NB_SAMPLES;
let frame_count = if index == 0 { 1 } else { 2 };
let packet_timestamp =
RedundancyScheduler::payload_timestamp(newest, frame_count, NB_SAMPLES);
if index % 2 == 1 {
continue; }
for (slot, is_new) in dedup
.accept(packet_timestamp, frame_count, NB_SAMPLES)
.into_iter()
.enumerate()
{
if is_new {
recovered
.push(packet_timestamp + u32::try_from(slot).unwrap_or(0) * NB_SAMPLES);
}
}
}
let expected: Vec<u32> = (0..19u32).map(|index| index * NB_SAMPLES).collect();
assert_eq!(
recovered, expected,
"depth 2 must survive alternate-packet loss with no gaps and no repeats"
);
}
#[test]
fn the_timestamp_space_wraps_without_swallowing_an_epoch() {
let mut dedup = RedundancyDedup::new();
let before = u32::MAX - NB_SAMPLES;
assert_eq!(dedup.accept(before, 1, NB_SAMPLES), [true]);
assert_eq!(
dedup.accept(before.wrapping_add(NB_SAMPLES), 1, NB_SAMPLES),
[true],
"a wrapped timestamp is newer, not older"
);
}
#[test]
fn no_data_frames_are_not_worth_repeating() {
let mut scheduler = RedundancyScheduler::new(Some(40), 3).expect("depth 3");
scheduler.next_payload(speech(0));
let gap =
AmrPayloadFrame::new(AmrFrameType::NoData, true, Vec::new()).expect("a NO_DATA frame");
scheduler.next_payload(gap);
let next = scheduler.next_payload(speech(2));
assert!(
next.iter()
.all(|frame| !matches!(frame.frame_type, AmrFrameType::NoData)),
"a NO_DATA frame was repeated as redundancy"
);
}
}