arcly-stream 0.7.1

An open-extensible live-media streaming kernel: lock-free zero-copy frame fan-out, instant-start GOP cache, a pluggable multi-protocol ingestion layer (RTMP, RTSP, SRT, WHIP/WHEP shipped), and a feature-gated pure-Rust media plane (MPEG-TS/HLS/fMP4) — runtime, config, and metrics free.
Documentation
//! SRT ARQ (Automatic Repeat reQuest) loss-recovery window.
//!
//! The receiver ([`Receiver`]) reorders data packets by sequence number,
//! delivers them in order, reports gaps with **NAK** loss messages, and
//! periodically **ACK**s its in-order progress. A gap that never fills is
//! force-skipped once the reorder buffer exceeds its window, so a permanently
//! lost packet can't wedge the stream.
//!
//! The sender ([`SendBuffer`]) keeps recently transmitted packets in a bounded
//! window, drops them as ACKs advance, and replays the ones a NAK reports lost
//! (with the data header's retransmit `R` bit set).
//!
//! Sequence numbers are 31-bit and wrap; comparisons below are modulo 2³¹. (A
//! gap straddling the single wrap point self-heals via the buffer-relief skip.)

use std::collections::BTreeMap;

use super::packet::ControlType;

/// 31-bit sequence-number modulus.
const SEQ_MOD: u32 = 0x8000_0000;
const SEQ_MASK: u32 = 0x7FFF_FFFF;

/// Next sequence number after `s` (wrapping in 31-bit space).
pub(crate) fn seq_inc(s: u32) -> u32 {
    s.wrapping_add(1) & SEQ_MASK
}

fn seq_dec(s: u32) -> u32 {
    s.wrapping_sub(1) & SEQ_MASK
}

/// `a < b` in 31-bit modular sequence space.
fn seq_lt(a: u32, b: u32) -> bool {
    let d = (b.wrapping_sub(a)) & SEQ_MASK;
    d != 0 && d < SEQ_MOD / 2
}

fn seq_le(a: u32, b: u32) -> bool {
    a == b || seq_lt(a, b)
}

// ---------------------------------------------------------------------------
// Control packet (NAK / ACK / ACKACK) build + parse.
// ---------------------------------------------------------------------------

const CTRL_ACK: u16 = 0x0002;
const CTRL_NAK: u16 = 0x0003;
const CTRL_ACKACK: u16 = 0x0006;

fn build_control(ctrl_type: u16, type_info: u32, timestamp: u32, dest: u32, cif: &[u8]) -> Vec<u8> {
    let mut out = Vec::with_capacity(16 + cif.len());
    let word0 = 0x8000_0000 | ((ctrl_type as u32) << 16);
    out.extend_from_slice(&word0.to_be_bytes());
    out.extend_from_slice(&type_info.to_be_bytes());
    out.extend_from_slice(&timestamp.to_be_bytes());
    out.extend_from_slice(&dest.to_be_bytes());
    out.extend_from_slice(cif);
    out
}

/// Build a NAK (loss report) for the given inclusive `ranges`. A single lost
/// packet is one word (high bit 0); a range is two words (high bit set on the
/// low bound, followed by the inclusive high bound).
pub(crate) fn build_nak(ranges: &[(u32, u32)], timestamp: u32, dest: u32) -> Vec<u8> {
    let mut cif = Vec::with_capacity(ranges.len() * 8);
    for &(lo, hi) in ranges {
        if lo == hi {
            cif.extend_from_slice(&(lo & SEQ_MASK).to_be_bytes());
        } else {
            cif.extend_from_slice(&((lo & SEQ_MASK) | SEQ_MOD).to_be_bytes());
            cif.extend_from_slice(&(hi & SEQ_MASK).to_be_bytes());
        }
    }
    build_control(CTRL_NAK, 0, timestamp, dest, &cif)
}

/// Parse a NAK's loss list into inclusive ranges. `datagram` is the full packet.
pub(crate) fn parse_nak(datagram: &[u8]) -> Vec<(u32, u32)> {
    let mut ranges = Vec::new();
    let cif = &datagram[datagram.len().min(16)..];
    let words: Vec<u32> = cif
        .chunks_exact(4)
        .map(|c| u32::from_be_bytes([c[0], c[1], c[2], c[3]]))
        .collect();
    let mut i = 0;
    while i < words.len() {
        let w = words[i];
        if w & SEQ_MOD != 0 {
            // Range start; the next word is the inclusive end.
            let lo = w & SEQ_MASK;
            if let Some(&hi) = words.get(i + 1) {
                ranges.push((lo, hi & SEQ_MASK));
                i += 2;
            } else {
                break;
            }
        } else {
            ranges.push((w & SEQ_MASK, w & SEQ_MASK));
            i += 1;
        }
    }
    ranges
}

/// Build an ACK acknowledging every packet below `ack_seq` (the next-expected
/// sequence number). `ack_no` is the ACK's own counter, echoed in the ACKACK.
pub(crate) fn build_ack(ack_no: u32, ack_seq: u32, timestamp: u32, dest: u32) -> Vec<u8> {
    build_control(
        CTRL_ACK,
        ack_no,
        timestamp,
        dest,
        &(ack_seq & SEQ_MASK).to_be_bytes(),
    )
}

/// Acknowledge an ACK (echoing its `ack_no` in the type-specific field).
pub(crate) fn build_ackack(ack_no: u32, timestamp: u32, dest: u32) -> Vec<u8> {
    build_control(CTRL_ACKACK, ack_no, timestamp, dest, &[])
}

/// The next-expected sequence number an ACK acknowledges, or `None` if the
/// datagram isn't a parseable ACK.
pub(crate) fn parse_ack(datagram: &[u8]) -> Option<u32> {
    let cif = datagram.get(16..20)?;
    Some(u32::from_be_bytes([cif[0], cif[1], cif[2], cif[3]]) & SEQ_MASK)
}

/// The ACK's own counter (its type-specific field), echoed back in the ACKACK.
pub(crate) fn ack_seqno(datagram: &[u8]) -> Option<u32> {
    let w = datagram.get(4..8)?;
    Some(u32::from_be_bytes([w[0], w[1], w[2], w[3]]))
}

/// Classify a control datagram (thin wrapper over the packet header type).
pub(crate) fn control_type(datagram: &[u8]) -> Option<ControlType> {
    match super::packet::SrtPacket::parse(datagram)? {
        super::packet::SrtPacket::Control { control_type, .. } => Some(control_type),
        _ => None,
    }
}

// ---------------------------------------------------------------------------
// Receiver: reorder, deliver in order, report loss.
// ---------------------------------------------------------------------------

/// Reorders incoming data packets and tracks loss for NAK/ACK generation.
pub(crate) struct Receiver {
    next: Option<u32>,
    buf: BTreeMap<u32, Vec<u8>>,
    window: usize,
}

impl Receiver {
    /// A receiver that force-skips a stuck gap once more than `window` packets
    /// are buffered out of order.
    pub(crate) fn new(window: usize) -> Receiver {
        Receiver {
            next: None,
            buf: BTreeMap::new(),
            window,
        }
    }

    /// Accept a data packet; return any payloads that are now in order.
    pub(crate) fn push(&mut self, seq: u32, payload: Vec<u8>) -> Vec<Vec<u8>> {
        let next = match self.next {
            None => {
                self.next = Some(seq);
                seq
            }
            Some(n) => n,
        };
        if seq_lt(seq, next) || self.buf.contains_key(&seq) {
            return Vec::new(); // late duplicate or already buffered
        }
        self.buf.insert(seq, payload);
        self.drain()
    }

    fn drain(&mut self) -> Vec<Vec<u8>> {
        let mut out = Vec::new();
        while let Some(n) = self.next {
            match self.buf.remove(&n) {
                Some(p) => {
                    out.push(p);
                    self.next = Some(seq_inc(n));
                }
                None => break,
            }
        }
        out
    }

    /// Inclusive ranges of missing sequence numbers between the next-expected
    /// packet and the highest one buffered — the NAK loss list.
    pub(crate) fn missing(&self) -> Vec<(u32, u32)> {
        let (Some(next), Some(&max)) = (self.next, self.buf.keys().next_back()) else {
            return Vec::new();
        };
        let mut ranges = Vec::new();
        let mut gap_start: Option<u32> = None;
        let mut s = next;
        while seq_le(s, max) {
            if self.buf.contains_key(&s) {
                if let Some(g) = gap_start.take() {
                    ranges.push((g, seq_dec(s)));
                }
            } else if gap_start.is_none() {
                gap_start = Some(s);
            }
            if ranges.len() >= 256 {
                break; // bound a pathological loss list
            }
            s = seq_inc(s);
        }
        ranges
    }

    /// If the reorder buffer has overflowed its window, abandon the stuck gap by
    /// jumping `next` to the lowest buffered packet; return what that releases.
    pub(crate) fn relieve(&mut self) -> Vec<Vec<u8>> {
        if self.buf.len() <= self.window {
            return Vec::new();
        }
        if let Some(&first) = self.buf.keys().next() {
            self.next = Some(first);
        }
        self.drain()
    }

    /// The sequence number an ACK should carry (everything below it is in hand).
    pub(crate) fn ack_seq(&self) -> Option<u32> {
        self.next
    }
}

// ---------------------------------------------------------------------------
// Sender: retransmission window.
// ---------------------------------------------------------------------------

/// A bounded buffer of recently sent data packets, for NAK-driven retransmit.
pub(crate) struct SendBuffer {
    packets: BTreeMap<u32, Vec<u8>>,
    window: usize,
}

impl SendBuffer {
    /// A window holding at most `window` packets (older ones are evicted).
    pub(crate) fn new(window: usize) -> SendBuffer {
        SendBuffer {
            packets: BTreeMap::new(),
            window,
        }
    }

    /// Record a freshly sent `datagram` under its sequence number.
    pub(crate) fn record(&mut self, seq: u32, datagram: Vec<u8>) {
        self.packets.insert(seq, datagram);
        while self.packets.len() > self.window {
            let oldest = *self.packets.keys().next().unwrap();
            self.packets.remove(&oldest);
        }
    }

    /// Drop everything an ACK has acknowledged (sequence `< ack_seq`).
    pub(crate) fn acknowledge(&mut self, ack_seq: u32) {
        self.packets.retain(|&seq, _| !seq_lt(seq, ack_seq));
    }

    /// Datagrams to retransmit for the NAK'd `ranges`, each with the retransmit
    /// `R` bit set. Packets no longer in the window are silently skipped.
    pub(crate) fn retransmit(&self, ranges: &[(u32, u32)]) -> Vec<Vec<u8>> {
        let mut out = Vec::new();
        for &(lo, hi) in ranges {
            let mut s = lo;
            loop {
                if let Some(pkt) = self.packets.get(&s) {
                    let mut dg = pkt.clone();
                    if dg.len() >= 5 {
                        dg[4] |= 0x04; // set R (retransmit) bit, word1 bit 26
                    }
                    out.push(dg);
                }
                if s == hi {
                    break;
                }
                s = seq_inc(s);
            }
        }
        out
    }

    /// Whether the window currently holds packet `seq`.
    #[cfg(test)]
    pub(crate) fn holds(&self, seq: u32) -> bool {
        self.packets.contains_key(&seq)
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn delivers_in_order_immediately() {
        let mut rx = Receiver::new(64);
        assert_eq!(rx.push(10, vec![1]), vec![vec![1]]);
        assert_eq!(rx.push(11, vec![2]), vec![vec![2]]);
        assert!(rx.missing().is_empty());
        assert_eq!(rx.ack_seq(), Some(12));
    }

    #[test]
    fn reorders_and_reports_gap() {
        let mut rx = Receiver::new(64);
        assert_eq!(rx.push(10, vec![10]), vec![vec![10]]);
        // 11 lost; 12,13 arrive and must be held.
        assert!(rx.push(12, vec![12]).is_empty());
        assert!(rx.push(13, vec![13]).is_empty());
        assert_eq!(rx.missing(), vec![(11, 11)]);
        // The retransmit of 11 unblocks 11,12,13 in order.
        assert_eq!(rx.push(11, vec![11]), vec![vec![11], vec![12], vec![13]]);
        assert!(rx.missing().is_empty());
    }

    #[test]
    fn reports_multi_packet_range() {
        let mut rx = Receiver::new(64);
        rx.push(5, vec![5]);
        rx.push(9, vec![9]); // 6,7,8 missing
        assert_eq!(rx.missing(), vec![(6, 8)]);
    }

    #[test]
    fn drops_duplicates_and_old() {
        let mut rx = Receiver::new(64);
        rx.push(20, vec![20]);
        rx.push(21, vec![21]);
        assert!(rx.push(20, vec![20]).is_empty(), "duplicate ignored");
        assert!(rx.push(5, vec![5]).is_empty(), "stale ignored");
    }

    #[test]
    fn relieves_a_stuck_gap_past_the_window() {
        let mut rx = Receiver::new(4);
        rx.push(0, vec![0]); // delivered, next=1
                             // 1 lost forever; pile up 2..=7 (6 buffered > window 4)
        for s in 2..=7u32 {
            rx.push(s, vec![s as u8]);
        }
        let released = rx.relieve();
        assert_eq!(
            released,
            (2..=7u32).map(|s| vec![s as u8]).collect::<Vec<_>>()
        );
        assert_eq!(rx.ack_seq(), Some(8));
    }

    #[test]
    fn nak_round_trips_ranges() {
        let ranges = vec![(11, 11), (20, 24)];
        let nak = build_nak(&ranges, 0, 7);
        assert_eq!(control_type(&nak), Some(ControlType::Nak));
        assert_eq!(parse_nak(&nak), ranges);
    }

    #[test]
    fn ack_round_trips() {
        let ack = build_ack(3, 100, 0, 9);
        assert_eq!(control_type(&ack), Some(ControlType::Ack));
        assert_eq!(parse_ack(&ack), Some(100));
    }

    #[test]
    fn send_buffer_retransmits_with_r_bit() {
        let mut tx = SendBuffer::new(8);
        // Minimal data datagrams: word1 (bytes 4..8) starts clear.
        for seq in 0..4u32 {
            let mut dg = vec![0u8; 16];
            dg[..4].copy_from_slice(&seq.to_be_bytes());
            tx.record(seq, dg);
        }
        let re = tx.retransmit(&[(1, 2)]);
        assert_eq!(re.len(), 2);
        assert_eq!(re[0][4] & 0x04, 0x04, "retransmit bit set");
    }

    #[test]
    fn send_buffer_evicts_and_acks() {
        let mut tx = SendBuffer::new(3);
        for seq in 0..5u32 {
            tx.record(seq, vec![0u8; 16]);
        }
        assert!(!tx.holds(0), "oldest evicted past window");
        assert!(tx.holds(4));
        tx.acknowledge(4); // drop everything < 4
        assert!(!tx.holds(3));
        assert!(tx.holds(4));
    }
}