rtc_interceptor/report/
receiver.rs1use 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
13pub struct ReceiverReportBuilder {
32 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 pub fn new() -> Self {
49 Self::default()
50 }
51
52 pub fn with_interval(mut self, interval: Duration) -> Self {
68 self.interval = interval;
69 self
70 }
71
72 pub fn build(self) -> ReceiverReportInterceptor {
85 ReceiverReportInterceptor::new(self.interval)
86 }
87}
88
89pub 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 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 fn process_rtp(&mut self, now: Instant, ssrc: u32, seq: u16, timestamp: u32) {
133 let stream = self.streams.entry(ssrc).or_insert_with(|| {
135 ReceiverStream::new(ssrc, 90000)
137 });
138
139 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 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 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 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 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 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}