Skip to main content

rtc_interceptor/nack/
responder.rs

1//! NACK Responder Interceptor - Responds to NACK requests by retransmitting packets.
2
3use super::send_buffer::SendBuffer;
4use super::stream_supports_nack;
5use crate::Interceptor;
6use crate::stream_info::StreamInfo;
7use crate::{Attribute, AttributedPacket, Packet, TaggedPacket};
8use sansio::Protocol;
9use shared::TransportContext;
10use shared::error::Error;
11use std::collections::{HashMap, VecDeque};
12use std::time::Instant;
13
14/// Builder for the NackResponderInterceptor.
15///
16/// # Example
17///
18/// ```
19/// use rtc_interceptor::{Slot, Registry, NackResponderBuilder};
20///
21/// let chain = Registry::new()
22///     .with(Slot::NackResponder, NackResponderBuilder::new()
23///         .with_size(1024)
24///         .build())
25///     .build();
26/// ```
27pub struct NackResponderBuilder {
28    /// Size of the send buffer (must be power of 2: 1, 2, 4, ..., 32768).
29    size: u16,
30}
31
32impl Default for NackResponderBuilder {
33    fn default() -> Self {
34        Self { size: 1024 }
35    }
36}
37
38impl NackResponderBuilder {
39    /// Create a new builder with default settings.
40    pub fn new() -> Self {
41        Self::default()
42    }
43
44    /// Set the size of the send buffer.
45    ///
46    /// Size must be a power of 2 between 1 and 32768 (inclusive).
47    /// Larger buffers can retransmit older packets but use more memory.
48    pub fn with_size(mut self, size: u16) -> Self {
49        self.size = size;
50        self
51    }
52
53    /// Build the interceptor.
54    pub fn build(self) -> NackResponderInterceptor {
55        NackResponderInterceptor::new(self.size)
56    }
57}
58
59/// Per-stream state for the responder.
60struct LocalStream {
61    /// Buffer of sent packets for retransmission.
62    send_buffer: SendBuffer,
63    /// RTX SSRC for RFC4588 retransmission (if configured).
64    ssrc_rtx: Option<u32>,
65    /// RTX payload type for RFC4588 retransmission (if configured).
66    payload_type_rtx: Option<u8>,
67    /// Sequence number counter for RTX packets.
68    rtx_sequence_number: u16,
69}
70
71/// Interceptor that responds to NACK requests by retransmitting packets.
72///
73/// This interceptor buffers outgoing RTP packets on local streams and
74/// retransmits them when RTCP TransportLayerNack packets are received.
75pub struct NackResponderInterceptor {
76    /// Configuration
77    size: u16,
78
79    /// Send buffers per local stream SSRC
80    streams: HashMap<u32, LocalStream>,
81
82    /// Queue for retransmitted packets
83    write_queue: VecDeque<TaggedPacket>,
84    /// Inbound packets ready for the next interceptor.
85    read_queue: VecDeque<TaggedPacket>,
86}
87
88impl NackResponderInterceptor {
89    fn new(size: u16) -> Self {
90        Self {
91            read_queue: VecDeque::new(),
92            size,
93            streams: HashMap::new(),
94            write_queue: VecDeque::new(),
95        }
96    }
97
98    /// Handle a NACK request by queuing retransmissions.
99    fn handle_nack(
100        &mut self,
101        now: Instant,
102        nack: &rtcp::transport_feedbacks::transport_layer_nack::TransportLayerNack,
103    ) {
104        // Collect sequence numbers to retransmit
105        let mut seqs_to_retransmit = Vec::new();
106
107        for nack_pair in &nack.nacks {
108            // Check the base packet ID
109            seqs_to_retransmit.push(nack_pair.packet_id);
110
111            // Check each bit in lost_packets bitmap
112            for i in 0..16 {
113                if nack_pair.lost_packets & (1 << i) != 0 {
114                    let seq = nack_pair.packet_id.wrapping_add(i + 1);
115                    seqs_to_retransmit.push(seq);
116                }
117            }
118        }
119
120        let Some(stream) = self.streams.get_mut(&nack.media_ssrc) else {
121            return;
122        };
123
124        // Queue retransmissions
125        for seq in seqs_to_retransmit {
126            let Some(original_packet) = stream.send_buffer.get(seq) else {
127                continue;
128            };
129
130            let packet = if let (Some(ssrc_rtx), Some(pt_rtx)) =
131                (stream.ssrc_rtx, stream.payload_type_rtx)
132            {
133                // RFC4588: Create RTX packet
134                // - Use RTX SSRC and payload type
135                // - Prepend original sequence number (2 bytes big-endian) to payload
136                // - Use separate RTX sequence number counter
137                let original_seq = original_packet.header.sequence_number;
138                let mut rtx_payload = Vec::with_capacity(2 + original_packet.payload.len());
139                rtx_payload.extend_from_slice(&original_seq.to_be_bytes());
140                rtx_payload.extend_from_slice(&original_packet.payload);
141
142                let rtx_seq = stream.rtx_sequence_number;
143                stream.rtx_sequence_number = stream.rtx_sequence_number.wrapping_add(1);
144
145                rtp::Packet {
146                    header: rtp::header::Header {
147                        // Not left to `..Default::default()`: the default
148                        // header is version 0, and receivers discard
149                        // version != 2 before examining anything else, so a
150                        // defaulted RTX packet is dropped on arrival and the
151                        // NACKed gap never repairs.
152                        version: 2,
153                        ssrc: ssrc_rtx,
154                        payload_type: pt_rtx,
155                        sequence_number: rtx_seq,
156                        timestamp: original_packet.header.timestamp,
157                        marker: original_packet.header.marker,
158                        ..Default::default()
159                    },
160                    payload: rtx_payload.into(),
161                }
162            } else {
163                // No RTX: retransmit original packet as-is
164                original_packet.clone()
165            };
166
167            // Tagged so a send history downstream counts it as new bytes on the wire rather
168            // than as the original transmission. A bandwidth estimator that cannot tell the two
169            // apart under-counts exactly when the path is lossy — it sees fewer bytes than are
170            // really being sent, infers headroom, and raises the target during loss.
171            self.write_queue.push_back(TaggedPacket {
172                now,
173                transport: TransportContext::default(),
174                message: AttributedPacket::new(Packet::Rtp(packet)).with(Attribute::Retransmission),
175            });
176        }
177    }
178}
179
180impl Protocol<TaggedPacket, TaggedPacket, ()> for NackResponderInterceptor {
181    type Rout = TaggedPacket;
182    type Wout = TaggedPacket;
183    type Eout = ();
184    type Error = Error;
185    type Time = Instant;
186
187    fn handle_read(&mut self, msg: TaggedPacket) -> Result<(), Self::Error> {
188        // Process NACK packets
189        if let Packet::Rtcp(ref rtcp_packets) = msg.message.packet {
190            for rtcp_packet in rtcp_packets {
191                if let Some(nack) = rtcp_packet
192                    .as_any()
193                    .downcast_ref::<rtcp::transport_feedbacks::transport_layer_nack::TransportLayerNack>()
194                {
195                    self.handle_nack(msg.now, nack);
196                }
197            }
198        }
199
200        self.read_queue.push_back(msg);
201
202        Ok(())
203    }
204
205    fn poll_read(&mut self) -> Option<Self::Rout> {
206        self.read_queue.pop_front()
207    }
208
209    fn handle_write(&mut self, msg: TaggedPacket) -> Result<(), Self::Error> {
210        // Buffer outgoing RTP packets
211        if let Packet::Rtp(ref rtp_packet) = msg.message.packet
212            && let Some(stream) = self.streams.get_mut(&rtp_packet.header.ssrc)
213        {
214            stream.send_buffer.add(rtp_packet.clone());
215        }
216
217        self.write_queue.push_back(msg);
218
219        Ok(())
220    }
221
222    fn poll_write(&mut self) -> Option<TaggedPacket> {
223        // First drain retransmitted packets
224        self.write_queue.pop_front()
225    }
226
227    fn handle_timeout(&mut self, _now: Instant) -> Result<(), Self::Error> {
228        Ok(())
229    }
230
231    fn poll_timeout(&mut self) -> Option<Self::Time> {
232        None
233    }
234}
235
236impl Interceptor for NackResponderInterceptor {
237    fn bind_local_stream(&mut self, info: &StreamInfo) {
238        if stream_supports_nack(info)
239            && let Some(send_buffer) = SendBuffer::new(self.size)
240        {
241            self.streams.insert(
242                info.ssrc,
243                LocalStream {
244                    send_buffer,
245                    ssrc_rtx: info.ssrc_rtx,
246                    payload_type_rtx: info.payload_type_rtx,
247                    rtx_sequence_number: 0,
248                },
249            );
250        }
251    }
252
253    fn unbind_local_stream(&mut self, info: &StreamInfo) {
254        self.streams.remove(&info.ssrc);
255    }
256
257    fn bind_remote_stream(&mut self, _info: &StreamInfo) {}
258
259    fn unbind_remote_stream(&mut self, _info: &StreamInfo) {}
260}