Skip to content
File

Blob: firmware/vendor/str0m/src/streams/send_queue.rs

rust347 lines
1use std::collections::VecDeque;
2use std::time::{Duration, Instant};
3 
4use crate::pacer::{QueuePriority, QueueSnapshot};
5use crate::util::not_happening;
6 
7use super::RtpPacket;
8 
9#[derive(Debug)]
10pub(crate) struct SendQueue {
11 queue: VecDeque<RtpPacket>,
12 total: TotalQueue,
13 last_emitted: Option<Instant>,
14}
15 
16impl SendQueue {
17 pub fn new() -> Self {
18 Self {
19 queue: VecDeque::new(),
20 total: TotalQueue::default(),
21 last_emitted: None,
22 }
23 }
24 
25 pub fn push(&mut self, mut packet: RtpPacket) {
26 // Every incoming packet must be timestamped withe a handle_timeout.
27 // This sentinel value indicates it is needed.
28 packet.timestamp = not_happening();
29 
30 self.queue.push_back(packet);
31 }
32 
33 pub fn handle_timeout(&mut self, now: Instant) {
34 for pkt in self.queue.iter_mut().rev() {
35 if pkt.timestamp != not_happening() {
36 // all enqueued packets are timestamped.
37 break;
38 } else {
39 pkt.timestamp = now;
40 self.total.increase(now, pkt.payload.len());
41 }
42 }
43 }
44 
45 pub fn need_timeout(&self) -> bool {
46 // Packets are timestamped contiguously from the front,
47 // so checking the tail summarizes the entire queue's state.
48 self.queue
49 .back()
50 .map(|p| p.timestamp == not_happening())
51 .unwrap_or_default()
52 }
53 
54 pub fn peek(&mut self) -> Option<&mut RtpPacket> {
55 let peeked = self.queue.front_mut()?;
56 if peeked.timestamp == not_happening() {
57 None
58 } else {
59 Some(peeked)
60 }
61 }
62 
63 pub fn pop(&mut self, now: Instant) -> Option<RtpPacket> {
64 // Don't release packets without a timestamp.
65 self.peek()?;
66 
67 // Unwrap is OK, because peek() above must have returned a value
68 // for us to be here.
69 let packet = self.queue.pop_front().unwrap();
70 
71 // Must be timestamped
72 assert!(packet.timestamp != not_happening());
73 
74 let queue_time = now - packet.timestamp;
75 self.total.decrease(now, packet.payload.len(), queue_time);
76 self.last_emitted = Some(now);
77 
78 Some(packet)
79 }
80 
81 pub fn is_empty(&self) -> bool {
82 self.queue.is_empty()
83 }
84 
85 pub fn last(&self) -> Option<&RtpPacket> {
86 self.queue.back()
87 }
88 
89 pub(crate) fn snapshot(&mut self, now: Instant) -> QueueSnapshot {
90 self.total.move_time_forward(now);
91 
92 QueueSnapshot {
93 created_at: now,
94 byte_size: self.total.unsent_size,
95 packet_count: self.total.unsent_count as u32,
96 total_queue_time_origin: self.total.queue_time,
97 last_emitted: self.last_emitted,
98 first_unsent: self
99 .queue
100 .iter()
101 .find(|p| p.timestamp != not_happening())
102 .map(|p| p.timestamp),
103 priority: if self.total.unsent_count > 0 {
104 QueuePriority::Media
105 } else {
106 QueuePriority::Empty
107 },
108 }
109 }
110 
111 pub(crate) fn clear(&mut self) {
112 self.queue.clear();
113 self.total.clear();
114 self.last_emitted = None;
115 }
116}
117 
118// Total queue time in buffer. This lovely drawing explains how to add more time.
119//
120// -time--------------------------------------------------------->
121//
122// +--------------+
123// | |
124// +--------------+
125// +---------+ Already
126// | | queued
127// +---------+ durations
128// +-----+
129// | |
130// +-----+
131// +-+
132// | | <----- Add next
133// +-+ packet
134//
135//
136//
137// +--------------+--------+
138// | |@@@@@@@@|
139// +--------------+--------+
140// +---------+--------+ The @ is
141// | |@@@@@@@@| what's
142// +---------+--------+ added
143// +-----+--------+
144// | |@@@@@@@@|
145// +-----+--------+
146// +-+
147// |@|
148// +-+
149#[derive(Debug, Default)]
150struct TotalQueue {
151 /// Number of unsent packets.
152 unsent_count: usize,
153 /// The data size (bytes) of the unsent packets.
154 unsent_size: usize,
155 // /// When we last added some value to `queue_time`.
156 // last: Option<Instant>,
157 /// The total queue time of all the unsent packets.
158 queue_time: Duration,
159 // Timestamp of the last added packet.
160 last: Option<Instant>,
161}
162 
163impl TotalQueue {
164 fn move_time_forward(&mut self, now: Instant) {
165 if let Some(last) = self.last {
166 assert!(self.unsent_count > 0);
167 let from_last = now - last;
168 self.queue_time += from_last * (self.unsent_count as u32);
169 self.last = Some(now);
170 } else {
171 assert!(self.unsent_count == 0);
172 assert!(self.unsent_size == 0);
173 assert!(self.queue_time == Duration::ZERO);
174 }
175 }
176 
177 fn increase(&mut self, now: Instant, size: usize) {
178 self.move_time_forward(now);
179 self.unsent_count += 1;
180 self.unsent_size += size;
181 self.last = Some(now);
182 }
183 
184 fn decrease(&mut self, now: Instant, size: usize, queue_time: Duration) {
185 self.move_time_forward(now);
186 
187 self.unsent_count -= 1;
188 self.unsent_size -= size;
189 
190 self.queue_time -= queue_time;
191 
192 if self.unsent_count == 0 {
193 assert!(self.unsent_size == 0);
194 self.queue_time = Duration::ZERO;
195 self.last = None;
196 }
197 }
198 
199 fn clear(&mut self) {
200 *self = Self::default();
201 }
202}
203 
204#[cfg(test)]
205mod test {
206 use crate::rtp_::MediaTime;
207 use crate::rtp_::RtpHeader;
208 
209 use super::*;
210 
211 #[test]
212 fn peek_pop_no_timestamp() {
213 let mut queue = SendQueue::new();
214 
215 queue.push(RtpPacket {
216 seq_no: 0.into(),
217 time: MediaTime::from_90khz(10),
218 header: RtpHeader::default(),
219 payload: [].into(),
220 vp8_patch: None,
221 timestamp: Instant::now(),
222 last_sender_info: None,
223 nackable: true,
224 });
225 
226 assert!(queue.peek().is_none());
227 assert!(queue.pop(Instant::now()).is_none());
228 assert!(queue.need_timeout());
229 
230 let snapshot_at = Instant::now() + Duration::from_secs(3);
231 assert_eq!(
232 queue.snapshot(snapshot_at),
233 QueueSnapshot {
234 created_at: snapshot_at,
235 packet_count: 0,
236 byte_size: 0,
237 total_queue_time_origin: Duration::ZERO,
238 first_unsent: None,
239 priority: QueuePriority::Empty,
240 ..Default::default()
241 }
242 );
243 }
244 
245 #[test]
246 fn peek_pop_after_timestamp() {
247 let mut queue = SendQueue::new();
248 
249 let start = Instant::now();
250 
251 queue.push(RtpPacket {
252 seq_no: 0.into(),
253 time: MediaTime::from_90khz(10),
254 header: RtpHeader::default(),
255 payload: [42, 42].into(),
256 vp8_patch: None,
257 timestamp: start,
258 last_sender_info: None,
259 nackable: true,
260 });
261 
262 queue.handle_timeout(start);
263 
264 assert!(queue.peek().is_some());
265 assert!(!queue.need_timeout());
266 
267 let snapshot_at = start + Duration::from_secs(3);
268 assert_eq!(
269 queue.snapshot(snapshot_at),
270 QueueSnapshot {
271 created_at: snapshot_at,
272 packet_count: 1,
273 byte_size: 2,
274 total_queue_time_origin: Duration::from_secs(3),
275 first_unsent: Some(start),
276 priority: QueuePriority::Media,
277 ..Default::default()
278 }
279 );
280 
281 assert!(queue.pop(Instant::now()).is_some());
282 }
283 
284 #[test]
285 fn untimed_packets_are_contiguous_at_tail() {
286 let mut queue = SendQueue::new();
287 let start = Instant::now();
288 
289 queue.push(RtpPacket {
290 seq_no: 0.into(),
291 time: MediaTime::from_90khz(10),
292 header: RtpHeader::default(),
293 payload: [0; 10].into(),
294 vp8_patch: None,
295 timestamp: start,
296 last_sender_info: None,
297 nackable: true,
298 });
299 queue.push(RtpPacket {
300 seq_no: 1.into(),
301 time: MediaTime::from_90khz(20),
302 header: RtpHeader::default(),
303 payload: [1; 10].into(),
304 vp8_patch: None,
305 timestamp: start,
306 last_sender_info: None,
307 nackable: true,
308 });
309 queue.push(RtpPacket {
310 seq_no: 2.into(),
311 time: MediaTime::from_90khz(20),
312 header: RtpHeader::default(),
313 payload: [1; 10].into(),
314 vp8_patch: None,
315 timestamp: start,
316 last_sender_info: None,
317 nackable: true,
318 });
319 
320 assert!(
321 queue.queue.iter().all(|q| q.timestamp == not_happening()),
322 "expect every new packet's timestamp to start with the sentinel value"
323 );
324 assert!(queue.need_timeout());
325 
326 let now = start + Duration::from_millis(10);
327 queue.handle_timeout(now);
328 
329 assert!(
330 queue.queue.iter().all(|q| q.timestamp != not_happening()),
331 "expect handle_timeout to have timestamped every packet"
332 );
333 assert!(!queue.need_timeout());
334 }
335 
336 #[test]
337 fn total_queue() {
338 let mut total_queue = TotalQueue::default();
339 let now = Instant::now();
340 total_queue.increase(now, 0);
341 total_queue.increase(now, 1);
342 total_queue.decrease(now, 1, Duration::ZERO);
343 // Doesn't panic
344 total_queue.move_time_forward(now + Duration::from_millis(1));
345 }
346}