Skip to main content

rtc_interceptor/report/
receiver.rs

1//! Receiver Report Interceptor - Generates RTCP Receiver Reports.
2
3use super::receiver_stream::ReceiverStream;
4use crate::Interceptor;
5use crate::stream_info::StreamInfo;
6use crate::{AttributedPacket, Packet, TaggedPacket};
7use sansio::Protocol;
8use shared::TransportContext;
9use shared::error::Error;
10use std::collections::{HashMap, VecDeque};
11use std::time::{Duration, Instant};
12
13/// Builder for the ReceiverReportInterceptor.
14///
15/// # Example
16///
17/// ```
18/// use rtc_interceptor::{Slot, Registry, ReceiverReportBuilder};
19/// use std::time::Duration;
20///
21/// // With default interval (1 second)
22/// let chain = Registry::new()
23///     .with(Slot::ReceiverReport, ReceiverReportBuilder::new().build())
24///     .build();
25///
26/// // With custom interval
27/// let chain = Registry::new()
28///     .with(Slot::ReceiverReport, ReceiverReportBuilder::new().with_interval(Duration::from_millis(500)).build())
29///     .build();
30/// ```
31pub struct ReceiverReportBuilder {
32    /// Interval between receiver reports.
33    interval: Duration,
34}
35
36impl Default for ReceiverReportBuilder {
37    fn default() -> Self {
38        Self {
39            interval: Duration::from_secs(1),
40        }
41    }
42}
43
44impl ReceiverReportBuilder {
45    /// Create a new builder with default settings.
46    ///
47    /// Default interval is 1 second.
48    pub fn new() -> Self {
49        Self::default()
50    }
51
52    /// Set a custom interval between receiver reports.
53    ///
54    /// # Example
55    ///
56    /// ```
57    /// use rtc_interceptor::{ReceiverReportBuilder, Registry, Slot};
58    /// use std::time::Duration;
59    ///
60    /// let registry = Registry::new().with(
61    ///     Slot::ReceiverReport,
62    ///     ReceiverReportBuilder::new()
63    ///         .with_interval(Duration::from_millis(500))
64    ///         .build(),
65    /// );
66    /// ```
67    pub fn with_interval(mut self, interval: Duration) -> Self {
68        self.interval = interval;
69        self
70    }
71
72    /// Create a builder function for use with Registry.
73    ///
74    /// This returns a closure that can be passed to `Registry::with()`.
75    ///
76    /// # Example
77    ///
78    /// ```
79    /// use rtc_interceptor::{Slot, Registry, ReceiverReportBuilder};
80    ///
81    /// let registry = Registry::new()
82    ///     .with(Slot::ReceiverReport, ReceiverReportBuilder::new().build());
83    /// ```
84    pub fn build(self) -> ReceiverReportInterceptor {
85        ReceiverReportInterceptor::new(self.interval)
86    }
87}
88
89/// Interceptor that generates RTCP Receiver Reports.
90///
91/// This interceptor monitors incoming RTP packets, tracks statistics per stream,
92/// and periodically generates RTCP Receiver Reports.
93///
94/// # Type Parameters
95///
96/// - `P`: The inner protocol being wrapped
97///
98/// # Example
99///
100/// ```
101/// use rtc_interceptor::{Slot, Registry, ReceiverReportBuilder};
102///
103/// let chain = Registry::new()
104///     .with(Slot::ReceiverReport, ReceiverReportBuilder::new().build())
105///     .build();
106/// ```
107pub struct ReceiverReportInterceptor {
108    interval: Duration,
109    next_timeout: Option<Instant>,
110
111    streams: HashMap<u32, ReceiverStream>,
112
113    read_queue: VecDeque<TaggedPacket>,
114    write_queue: VecDeque<TaggedPacket>,
115}
116
117impl ReceiverReportInterceptor {
118    /// Create a new ReceiverReportInterceptor with default configuration.
119    fn new(interval: Duration) -> Self {
120        Self {
121            interval,
122            next_timeout: None,
123
124            streams: HashMap::new(),
125
126            read_queue: VecDeque::new(),
127            write_queue: VecDeque::new(),
128        }
129    }
130
131    /// Process an incoming RTP packet for statistics.
132    fn process_rtp(&mut self, now: Instant, ssrc: u32, seq: u16, timestamp: u32) {
133        // Create stream if it doesn't exist
134        let stream = self.streams.entry(ssrc).or_insert_with(|| {
135            // Default clock rate, should be configured per stream in real usage
136            ReceiverStream::new(ssrc, 90000)
137        });
138
139        // Create a minimal RTP packet for processing
140        let pkt = rtp::packet::Packet {
141            header: rtp::header::Header {
142                ssrc,
143                sequence_number: seq,
144                timestamp,
145                ..Default::default()
146            },
147            ..Default::default()
148        };
149
150        stream.process_rtp(now, &pkt);
151    }
152
153    /// Process an incoming RTCP Sender Report.
154    fn process_sender_report(&mut self, now: Instant, sr: &rtcp::sender_report::SenderReport) {
155        if let Some(stream) = self.streams.get_mut(&sr.ssrc) {
156            stream.process_sender_report(now, sr);
157        }
158    }
159
160    /// Generate receiver reports for all tracked streams.
161    fn generate_reports(&mut self, now: Instant) -> Vec<rtcp::receiver_report::ReceiverReport> {
162        self.streams
163            .values_mut()
164            .map(|stream| stream.generate_report(now))
165            .collect()
166    }
167
168    /// Register a new stream with its clock rate.
169    fn register_stream(&mut self, ssrc: u32, clock_rate: u32) {
170        self.streams
171            .entry(ssrc)
172            .or_insert_with(|| ReceiverStream::new(ssrc, clock_rate));
173    }
174}
175
176impl Protocol<TaggedPacket, TaggedPacket, ()> for ReceiverReportInterceptor {
177    type Rout = TaggedPacket;
178    type Wout = TaggedPacket;
179    type Eout = ();
180    type Error = Error;
181    type Time = Instant;
182
183    fn handle_read(&mut self, msg: TaggedPacket) -> Result<(), Self::Error> {
184        if let Packet::Rtcp(rtcp_packets) = &msg.message.packet {
185            for rtcp_packet in rtcp_packets {
186                if let Some(sr) = rtcp_packet
187                    .as_any()
188                    .downcast_ref::<rtcp::sender_report::SenderReport>()
189                    && let Some(stream) = self.streams.get_mut(&sr.ssrc)
190                {
191                    stream.process_sender_report(msg.now, sr);
192                }
193            }
194        } else if let Packet::Rtp(rtp_packet) = &msg.message.packet
195            && let Some(stream) = self.streams.get_mut(&rtp_packet.header.ssrc)
196        {
197            stream.process_rtp(msg.now, rtp_packet);
198
199            // Arm the report timer from the first packet's instant (see nack::generator).
200            if self.next_timeout.is_none() {
201                self.next_timeout = Some(msg.now + self.interval);
202            }
203        }
204
205        self.read_queue.push_back(msg);
206
207        Ok(())
208    }
209
210    fn poll_read(&mut self) -> Option<Self::Rout> {
211        self.read_queue.pop_front()
212    }
213
214    fn handle_write(&mut self, msg: TaggedPacket) -> Result<(), Self::Error> {
215        self.write_queue.push_back(msg);
216        Ok(())
217    }
218
219    fn poll_write(&mut self) -> Option<TaggedPacket> {
220        // First drain generated RTCP reports
221        if let Some(pkt) = self.write_queue.pop_front() {
222            return Some(pkt);
223        }
224        None
225    }
226
227    fn handle_timeout(&mut self, now: Instant) -> Result<(), Error> {
228        if let Some(next_timeout) = self.next_timeout
229            && now >= next_timeout
230        {
231            self.next_timeout = Some(now + self.interval);
232
233            for stream in self.streams.values_mut() {
234                let rr = stream.generate_report(now);
235                self.write_queue.push_back(TaggedPacket {
236                    now,
237                    transport: TransportContext::default(),
238                    message: AttributedPacket::new(Packet::Rtcp(vec![Box::new(rr)])),
239                });
240            }
241        }
242        Ok(())
243    }
244
245    fn poll_timeout(&mut self) -> Option<Instant> {
246        self.next_timeout
247    }
248}
249
250impl Interceptor for ReceiverReportInterceptor {
251    fn bind_remote_stream(&mut self, info: &StreamInfo) {
252        let stream = ReceiverStream::new(info.ssrc, info.clock_rate);
253        self.streams.insert(info.ssrc, stream);
254    }
255
256    fn unbind_remote_stream(&mut self, info: &StreamInfo) {
257        self.streams.remove(&info.ssrc);
258    }
259
260    fn bind_local_stream(&mut self, _info: &StreamInfo) {}
261
262    fn unbind_local_stream(&mut self, _info: &StreamInfo) {}
263}