Skip to main content

rtc_rtp/packetizer/
mod.rs

1#[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
16/// Payloader payloads a byte array for use as rtp.Packet payloads
17pub trait Payloader: Send + Sync + fmt::Debug {
18    /// Splits one encoded frame into payloads no larger than `mtu`.
19    ///
20    /// # Errors
21    ///
22    /// Fails if the frame is malformed for this codec, or `mtu` is too small to make progress.
23    fn payload(&mut self, mtu: usize, b: &Bytes) -> Result<Vec<Bytes>>;
24    /// Clones this payloader behind a trait object.
25    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
34/// Packetizer packetizes a payload
35pub trait Packetizer: Send + Sync + fmt::Debug {
36    /// Attaches the absolute-send-time header extension under id `value`.
37    fn enable_abs_send_time(&mut self, value: u8);
38    /// Packetizes one frame, advancing the timestamp by `samples`.
39    ///
40    /// Assigns sequence numbers, sets the marker bit on the final packet, and applies any
41    /// enabled header extensions.
42    ///
43    /// # Errors
44    ///
45    /// Propagates payloader failures.
46    fn packetize(&mut self, payload: &Bytes, samples: u32) -> Result<Vec<Packet>>;
47    /// Advances the timestamp without sending anything, for dropped or silent frames.
48    fn skip_samples(&mut self, skipped_samples: u32);
49    /// Clones this packetizer behind a trait object.
50    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
59/// Depacketizer depacketizes a RTP payload, removing any RTP specific data from the payload
60pub trait Depacketizer {
61    /// Reassembles a frame from one RTP payload, buffering fragments as needed.
62    ///
63    /// # Errors
64    ///
65    /// Fails if the payload is malformed for this codec.
66    fn depacketize(&mut self, b: &Bytes) -> Result<Bytes>;
67
68    /// Checks if the packet is at the beginning of a partition.  This
69    /// should return false if the result could not be determined, in
70    /// which case the caller will detect timestamp discontinuities.
71    fn is_partition_head(&self, payload: &Bytes) -> bool;
72
73    /// Checks if the packet is at the end of a partition.  This should
74    /// return false if the result could not be determined.
75    fn is_partition_tail(&self, marker: bool, payload: &Bytes) -> bool;
76}
77
78/// FnTimeGen provides current time (Instant)
79pub 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, //http://www.webrtc.org/experiments/rtp-hdrext/abs-send-time
91    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
108/// Builds a packetizer for one outbound stream.
109///
110/// Ties together the codec's payloader, a sequencer, and the SSRC, payload type, MTU and clock
111/// rate the stream was negotiated with.
112pub 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, //TODO: Figure out how to do timestamps
153                    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            //apply http://www.webrtc.org/experiments/rtp-hdrext/abs-send-time
170            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    /// skip_samples causes a gap in sample count between Packetize requests so the
182    /// RTP payloads produced have a gap in timestamps
183    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}