File
Blob: firmware/vendor/str0m/src/pacer/queue.rs
| 1 | use std::time::{Duration, Instant}; |
| 2 | |
| 3 | use crate::rtp_::MidRid; |
| 4 | use crate::util::not_happening; |
| 5 | |
| 6 | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 7 | pub struct QueueSnapshot { |
| 8 | /// Time this snapshot was made |
| 9 | pub created_at: Instant, |
| 10 | /// The total byte size of the snapshot. |
| 11 | pub byte_size: usize, |
| 12 | /// The total number of packets in the queue. |
| 13 | /// NB: This is not a [`usize`] because it will later be used to divide a [`Duration`], for which |
| 14 | /// [`usize`] isn't implement. If the queues end up with 2^32 packets something has gone very wrong |
| 15 | /// in any case. |
| 16 | pub packet_count: u32, |
| 17 | /// Accumulation of all queue time at the time point `created_at`. To use this |
| 18 | /// Look at `total_queue_time(now)` which allows getting the queue time at a later Instant. |
| 19 | pub total_queue_time_origin: Duration, |
| 20 | /// Last time something was emitted from this queue. |
| 21 | pub last_emitted: Option<Instant>, |
| 22 | /// Time the first unsent packet has spent in the queue. |
| 23 | pub first_unsent: Option<Instant>, |
| 24 | /// The priority of the most important packet in the queue. |
| 25 | pub priority: QueuePriority, |
| 26 | } |
| 27 | |
| 28 | /// Priority for a given queue. |
| 29 | /// |
| 30 | /// When sorted, higher priority sorts first. |
| 31 | #[derive(Default, Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord)] |
| 32 | pub enum QueuePriority { |
| 33 | // Highest, priority for a queue that contains media. |
| 34 | Media = 0, |
| 35 | // Priority for a queue that only contains padding. |
| 36 | Padding = 1, |
| 37 | // Priority for an empty queue. |
| 38 | #[default] |
| 39 | Empty = 2, |
| 40 | } |
| 41 | |
| 42 | impl QueueSnapshot { |
| 43 | /// Update the priority if the snapshot is non-empty. |
| 44 | /// |
| 45 | /// Sets the priority to the provided priority if the queue is non-empty, otherwise empty. |
| 46 | pub fn update_priority(&mut self, priority: QueuePriority) { |
| 47 | if self.packet_count > 0 { |
| 48 | self.priority = priority; |
| 49 | } else { |
| 50 | self.priority = QueuePriority::Empty; |
| 51 | } |
| 52 | } |
| 53 | } |
| 54 | |
| 55 | impl Default for QueueSnapshot { |
| 56 | fn default() -> Self { |
| 57 | Self { |
| 58 | created_at: not_happening(), |
| 59 | byte_size: Default::default(), |
| 60 | packet_count: Default::default(), |
| 61 | total_queue_time_origin: Default::default(), |
| 62 | last_emitted: Default::default(), |
| 63 | first_unsent: Default::default(), |
| 64 | priority: QueuePriority::default(), |
| 65 | } |
| 66 | } |
| 67 | } |
| 68 | |
| 69 | /// The state of a single upstream queue. |
| 70 | /// The pacer manages packets across several upstream queues. |
| 71 | #[derive(Debug, Clone, Copy)] |
| 72 | pub struct QueueState { |
| 73 | pub midrid: MidRid, |
| 74 | pub unpaced: bool, |
| 75 | pub use_for_padding: bool, |
| 76 | pub snapshot: QueueSnapshot, |
| 77 | } |
| 78 | |
| 79 | /// A request to generate a specific amount of padding. |
| 80 | #[derive(Debug, Clone, Copy)] |
| 81 | pub struct PaddingRequest { |
| 82 | /// The Mid that should generate and queue the padding. |
| 83 | pub midrid: MidRid, |
| 84 | /// The amount of padding in bytes to generate. |
| 85 | pub padding: usize, |
| 86 | } |
| 87 | |
| 88 | impl QueueSnapshot { |
| 89 | /// Merge other into self. |
| 90 | pub fn merge(&mut self, other: &Self) { |
| 91 | self.created_at = self.created_at.min(other.created_at); |
| 92 | self.byte_size += other.byte_size; |
| 93 | self.packet_count += other.packet_count; |
| 94 | self.total_queue_time_origin += other.total_queue_time_origin; |
| 95 | self.last_emitted = self.last_emitted.max(other.last_emitted); |
| 96 | self.first_unsent = match (self.first_unsent, other.first_unsent) { |
| 97 | (None, None) => None, |
| 98 | (None, Some(v2)) => Some(v2), |
| 99 | (Some(v1), None) => Some(v1), |
| 100 | (Some(v1), Some(v2)) => Some(v1.min(v2)), |
| 101 | }; |
| 102 | self.priority = self.priority.min(other.priority); |
| 103 | } |
| 104 | |
| 105 | pub fn total_queue_time(&self, now: Instant) -> Duration { |
| 106 | self.total_queue_time_origin + self.packet_count * (now - self.created_at) |
| 107 | } |
| 108 | } |