rtc_rtp/packetizer/
mod.rs1#[cfg(test)]
2mod packetizer_test;
3
4use crate::{extension::abs_send_time_extension::*, header::*, packet::*, sequence::*};
5use shared::{
6 error::Result,
7 marshal::{Marshal, MarshalSize},
8 time::SystemInstant,
9};
10
11use bytes::{Bytes, BytesMut};
12use std::fmt;
13use std::sync::Arc;
14use std::time::Instant;
15
16pub trait Payloader: Send + Sync + fmt::Debug {
18 fn payload(&mut self, mtu: usize, b: &Bytes) -> Result<Vec<Bytes>>;
24 fn clone_to(&self) -> Box<dyn Payloader>;
26}
27
28impl Clone for Box<dyn Payloader> {
29 fn clone(&self) -> Box<dyn Payloader> {
30 self.clone_to()
31 }
32}
33
34pub trait Packetizer: Send + Sync + fmt::Debug {
36 fn enable_abs_send_time(&mut self, value: u8);
38 fn packetize(&mut self, payload: &Bytes, samples: u32) -> Result<Vec<Packet>>;
47 fn skip_samples(&mut self, skipped_samples: u32);
49 fn clone_to(&self) -> Box<dyn Packetizer>;
51}
52
53impl Clone for Box<dyn Packetizer> {
54 fn clone(&self) -> Box<dyn Packetizer> {
55 self.clone_to()
56 }
57}
58
59pub trait Depacketizer {
61 fn depacketize(&mut self, b: &Bytes) -> Result<Bytes>;
67
68 fn is_partition_head(&self, payload: &Bytes) -> bool;
72
73 fn is_partition_tail(&self, marker: bool, payload: &Bytes) -> bool;
76}
77
78pub type FnTimeGen = Arc<dyn (Fn() -> Instant) + Send + Sync>;
80
81#[derive(Clone)]
82pub(crate) struct PacketizerImpl {
83 pub(crate) mtu: usize,
84 pub(crate) payload_type: u8,
85 pub(crate) ssrc: u32,
86 pub(crate) payloader: Box<dyn Payloader>,
87 pub(crate) sequencer: Box<dyn Sequencer>,
88 pub(crate) timestamp: u32,
89 pub(crate) clock_rate: u32,
90 pub(crate) abs_send_time_ext_id: u8, pub(crate) time_gen: Option<FnTimeGen>,
92 time_baseline: SystemInstant,
93}
94
95impl fmt::Debug for PacketizerImpl {
96 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
97 f.debug_struct("PacketizerImpl")
98 .field("mtu", &self.mtu)
99 .field("payload_type", &self.payload_type)
100 .field("ssrc", &self.ssrc)
101 .field("timestamp", &self.timestamp)
102 .field("clock_rate", &self.clock_rate)
103 .field("abs_send_time_ext_id", &self.abs_send_time_ext_id)
104 .finish()
105 }
106}
107
108pub fn new_packetizer(
113 mtu: usize,
114 payload_type: u8,
115 ssrc: u32,
116 payloader: Box<dyn Payloader>,
117 sequencer: Box<dyn Sequencer>,
118 clock_rate: u32,
119) -> impl Packetizer {
120 PacketizerImpl {
121 mtu,
122 payload_type,
123 ssrc,
124 payloader,
125 sequencer,
126 timestamp: rand::random::<u32>(),
127 clock_rate,
128 abs_send_time_ext_id: 0,
129 time_gen: None,
130 time_baseline: SystemInstant::now(),
131 }
132}
133
134impl Packetizer for PacketizerImpl {
135 fn enable_abs_send_time(&mut self, id: u8) {
136 self.abs_send_time_ext_id = id
137 }
138
139 fn packetize(&mut self, payload: &Bytes, samples: u32) -> Result<Vec<Packet>> {
140 let payloads = self.payloader.payload(self.mtu - 12, payload)?;
141 let payloads_len = payloads.len();
142 let mut packets = Vec::with_capacity(payloads_len);
143 for (i, payload) in payloads.into_iter().enumerate() {
144 packets.push(Packet {
145 header: Header {
146 version: 2,
147 padding: false,
148 extension: false,
149 marker: i == payloads_len - 1,
150 payload_type: self.payload_type,
151 sequence_number: self.sequencer.next_sequence_number(),
152 timestamp: self.timestamp, ssrc: self.ssrc,
154 ..Default::default()
155 },
156 payload,
157 });
158 }
159
160 self.timestamp = self.timestamp.wrapping_add(samples);
161
162 if payloads_len != 0 && self.abs_send_time_ext_id != 0 {
163 let now = if let Some(fn_time_gen) = &self.time_gen {
164 fn_time_gen()
165 } else {
166 Instant::now()
167 };
168 let send_time = AbsSendTimeExtension::new(self.time_baseline.ntp(now));
169 let mut raw = BytesMut::with_capacity(send_time.marshal_size());
171 raw.resize(send_time.marshal_size(), 0);
172 let _ = send_time.marshal_to(&mut raw)?;
173 packets[payloads_len - 1]
174 .header
175 .set_extension(self.abs_send_time_ext_id, raw.freeze())?;
176 }
177
178 Ok(packets)
179 }
180
181 fn skip_samples(&mut self, skipped_samples: u32) {
184 self.timestamp = self.timestamp.wrapping_add(skipped_samples);
185 }
186
187 fn clone_to(&self) -> Box<dyn Packetizer> {
188 Box::new(self.clone())
189 }
190}