gemini_adk_fluent_rs/telephony/
rtp.rs

1//! RTP (RFC 3550) — the media framing of every SIP call.
2//!
3//! A raw SIP/PSTN leg carries audio as RTP packets over UDP: a 12-byte fixed
4//! header (sequence number, timestamp, SSRC) followed by the codec payload —
5//! for telephone audio, G.711 μ-law (payload type 0, `PCMU`) or A-law
6//! (payload type 8, `PCMA`) at 8 kHz, conventionally 20 ms (160 samples) per
7//! packet.
8//!
9//! This module is the pure layer: [`build`] and [`parse`] move between packet
10//! bytes and structured form (tolerating padding, CSRCs, and header
11//! extensions on the way in), and [`RtpSender`] carries the tiny amount of
12//! state a sender needs (sequence, timestamp, SSRC). No sockets — the `sip`
13//! feature's media loop drives it over UDP, and tests drive it with byte
14//! arrays.
15
16/// RTP payload type for G.711 μ-law at 8 kHz (RFC 3551 static assignment).
17pub const PT_PCMU: u8 = 0;
18/// RTP payload type for G.711 A-law at 8 kHz (RFC 3551 static assignment).
19pub const PT_PCMA: u8 = 8;
20
21/// Samples per packet at the conventional 20 ms packetisation (8 kHz mono).
22pub const SAMPLES_PER_PACKET: usize = 160;
23
24/// One parsed RTP packet (header fields we care about + payload bytes).
25#[derive(Debug, Clone, PartialEq)]
26pub struct RtpPacket {
27    /// Payload type (e.g. [`PT_PCMU`], [`PT_PCMA`]).
28    pub payload_type: u8,
29    /// Marker bit — set on the first packet after silence (talkspurt start).
30    pub marker: bool,
31    /// Sequence number, increments by one per packet.
32    pub sequence: u16,
33    /// Media timestamp in samples (8 kHz clock for G.711).
34    pub timestamp: u32,
35    /// Synchronisation source identifier.
36    pub ssrc: u32,
37    /// Codec payload bytes.
38    pub payload: Vec<u8>,
39}
40
41/// Build an RTP packet (version 2, no padding/extension/CSRC).
42pub fn build(packet: &RtpPacket) -> Vec<u8> {
43    let mut out = Vec::with_capacity(12 + packet.payload.len());
44    out.push(0x80); // V=2, P=0, X=0, CC=0
45    out.push((packet.payload_type & 0x7F) | if packet.marker { 0x80 } else { 0 });
46    out.extend_from_slice(&packet.sequence.to_be_bytes());
47    out.extend_from_slice(&packet.timestamp.to_be_bytes());
48    out.extend_from_slice(&packet.ssrc.to_be_bytes());
49    out.extend_from_slice(&packet.payload);
50    out
51}
52
53/// Parse an RTP packet, tolerating padding, CSRC entries, and a header
54/// extension. Returns `None` for datagrams that are not well-formed RTP v2 —
55/// a media port sees stray traffic (STUN probes, scans); dropping quietly is
56/// the correct posture.
57pub fn parse(datagram: &[u8]) -> Option<RtpPacket> {
58    if datagram.len() < 12 {
59        return None;
60    }
61    let b0 = datagram[0];
62    if b0 >> 6 != 2 {
63        return None; // not RTP version 2
64    }
65    let has_padding = b0 & 0x20 != 0;
66    let has_extension = b0 & 0x10 != 0;
67    let csrc_count = (b0 & 0x0F) as usize;
68    let b1 = datagram[1];
69
70    let mut offset = 12 + csrc_count * 4;
71    if datagram.len() < offset {
72        return None;
73    }
74    if has_extension {
75        if datagram.len() < offset + 4 {
76            return None;
77        }
78        let ext_words = u16::from_be_bytes([datagram[offset + 2], datagram[offset + 3]]) as usize;
79        offset += 4 + ext_words * 4;
80        if datagram.len() < offset {
81            return None;
82        }
83    }
84    let mut end = datagram.len();
85    if has_padding {
86        let pad = *datagram.last()? as usize;
87        if pad == 0 || offset + pad > end {
88            return None;
89        }
90        end -= pad;
91    }
92
93    Some(RtpPacket {
94        payload_type: b1 & 0x7F,
95        marker: b1 & 0x80 != 0,
96        sequence: u16::from_be_bytes([datagram[2], datagram[3]]),
97        timestamp: u32::from_be_bytes([datagram[4], datagram[5], datagram[6], datagram[7]]),
98        ssrc: u32::from_be_bytes([datagram[8], datagram[9], datagram[10], datagram[11]]),
99        payload: datagram[offset..end].to_vec(),
100    })
101}
102
103// ── Telephone events (RFC 4733 DTMF) ────────────────────────────────────────
104
105/// One parsed telephone-event (RFC 4733) payload — a DTMF keypress carried
106/// as RTP instead of audio tones.
107#[derive(Debug, Clone, Copy, PartialEq, Eq)]
108pub struct TelephoneEvent {
109    /// Event code: 0–9 are the digits, 10 is `*`, 11 is `#`, 12–15 are A–D.
110    pub event: u8,
111    /// End bit — set on the final packet(s) of the keypress. End packets are
112    /// conventionally retransmitted three times; deduplicate on
113    /// [`RtpPacket::timestamp`], which stays constant for one keypress.
114    pub end: bool,
115    /// Cumulative duration of the event so far, in timestamp units.
116    pub duration: u16,
117}
118
119impl TelephoneEvent {
120    /// The event as the character a keypad prints, `None` for codes > 15.
121    pub fn digit(&self) -> Option<char> {
122        Some(match self.event {
123            0..=9 => (b'0' + self.event) as char,
124            10 => '*',
125            11 => '#',
126            12..=15 => (b'A' + self.event - 12) as char,
127            _ => return None,
128        })
129    }
130}
131
132/// Parse an RFC 4733 telephone-event payload (the 4-byte named-event form).
133///
134/// The caller decides *whether* a packet is a telephone event from the
135/// payload type negotiated in SDP — this parses the payload of one that is.
136pub fn parse_telephone_event(payload: &[u8]) -> Option<TelephoneEvent> {
137    if payload.len() < 4 {
138        return None;
139    }
140    Some(TelephoneEvent {
141        event: payload[0],
142        end: payload[1] & 0x80 != 0,
143        duration: u16::from_be_bytes([payload[2], payload[3]]),
144    })
145}
146
147/// Sender-side RTP state: sequence, timestamp, and SSRC advance per packet.
148#[derive(Debug)]
149pub struct RtpSender {
150    payload_type: u8,
151    sequence: u16,
152    timestamp: u32,
153    ssrc: u32,
154    /// Set the marker bit on the next packet (start of a talkspurt).
155    mark_next: bool,
156}
157
158impl RtpSender {
159    /// Create a sender for the given payload type and SSRC.
160    ///
161    /// Callers should pick a random-ish SSRC and initial sequence/timestamp;
162    /// determinism here keeps the function pure — randomness is the caller's.
163    pub fn new(payload_type: u8, ssrc: u32, initial_sequence: u16, initial_timestamp: u32) -> Self {
164        Self {
165            payload_type,
166            sequence: initial_sequence,
167            timestamp: initial_timestamp,
168            ssrc,
169            mark_next: true,
170        }
171    }
172
173    /// Frame one payload as an RTP packet and advance sequence/timestamp.
174    ///
175    /// `samples` is the number of media samples the payload covers (for
176    /// G.711, one byte per sample — 160 for a 20 ms packet).
177    pub fn packetize(&mut self, payload: &[u8], samples: u32) -> Vec<u8> {
178        let packet = build(&RtpPacket {
179            payload_type: self.payload_type,
180            marker: self.mark_next,
181            sequence: self.sequence,
182            timestamp: self.timestamp,
183            ssrc: self.ssrc,
184            payload: payload.to_vec(),
185        });
186        self.mark_next = false;
187        self.sequence = self.sequence.wrapping_add(1);
188        self.timestamp = self.timestamp.wrapping_add(samples);
189        packet
190    }
191
192    /// Advance the media clock over a silent gap and mark the next packet as
193    /// the start of a new talkspurt.
194    pub fn skip_silence(&mut self, samples: u32) {
195        self.timestamp = self.timestamp.wrapping_add(samples);
196        self.mark_next = true;
197    }
198}
199
200#[cfg(test)]
201mod tests {
202    use super::*;
203
204    #[test]
205    fn builds_the_canonical_wire_layout() {
206        let bytes = build(&RtpPacket {
207            payload_type: PT_PCMU,
208            marker: true,
209            sequence: 0x0102,
210            timestamp: 0x03040506,
211            ssrc: 0x0708090A,
212            payload: vec![0xFF, 0xFE],
213        });
214        assert_eq!(
215            bytes,
216            vec![
217                0x80, 0x80, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0A, 0xFF, 0xFE
218            ]
219        );
220    }
221
222    #[test]
223    fn parse_round_trips_build() {
224        let packet = RtpPacket {
225            payload_type: PT_PCMA,
226            marker: false,
227            sequence: 65_535,
228            timestamp: u32::MAX - 1,
229            ssrc: 42,
230            payload: vec![1, 2, 3, 4],
231        };
232        assert_eq!(parse(&build(&packet)), Some(packet));
233    }
234
235    #[test]
236    fn parse_skips_csrc_extension_and_padding() {
237        // V=2, P=1, X=1, CC=1 · PT=0 · seq 1 · ts 2 · ssrc 3
238        let mut bytes = vec![0xB1, 0x00, 0x00, 0x01, 0, 0, 0, 2, 0, 0, 0, 3];
239        bytes.extend_from_slice(&[9, 9, 9, 9]); // one CSRC
240        bytes.extend_from_slice(&[0xBE, 0xDE, 0x00, 0x01, 0, 0, 0, 0]); // ext: 1 word
241        bytes.extend_from_slice(&[0xAA, 0xBB]); // payload
242        bytes.extend_from_slice(&[0, 0, 3]); // 3 bytes padding (last byte = count)
243        let packet = parse(&bytes).expect("valid despite extras");
244        assert_eq!(packet.payload, vec![0xAA, 0xBB]);
245        assert_eq!(packet.sequence, 1);
246    }
247
248    #[test]
249    fn parse_rejects_garbage() {
250        assert_eq!(parse(&[]), None);
251        assert_eq!(parse(&[0x80; 5]), None); // too short
252        assert_eq!(parse(&[0x00; 20]), None); // version 0 (e.g. STUN)
253    }
254
255    #[test]
256    fn telephone_events_parse_digits_and_end_bits() {
257        // '5' pressed, not yet released, 160 timestamp units in.
258        let event = parse_telephone_event(&[5, 0x0A, 0x00, 0xA0]).unwrap();
259        assert_eq!(event.digit(), Some('5'));
260        assert!(!event.end);
261        assert_eq!(event.duration, 160);
262
263        // '#' released (end bit set).
264        let end = parse_telephone_event(&[11, 0x8A, 0x03, 0x20]).unwrap();
265        assert_eq!(end.digit(), Some('#'));
266        assert!(end.end);
267
268        assert_eq!(
269            parse_telephone_event(&[12, 0x80, 0, 60]).unwrap().digit(),
270            Some('A')
271        );
272        assert_eq!(
273            parse_telephone_event(&[10, 0x80, 0, 60]).unwrap().digit(),
274            Some('*')
275        );
276        // Flash-hook (16) and other extended events carry no keypad digit.
277        assert_eq!(
278            parse_telephone_event(&[16, 0x80, 0, 60]).unwrap().digit(),
279            None
280        );
281        // Truncated payload.
282        assert_eq!(parse_telephone_event(&[5, 0x80]), None);
283    }
284
285    #[test]
286    fn sender_advances_and_marks_talkspurts() {
287        let mut sender = RtpSender::new(PT_PCMU, 7, 100, 1000);
288        let first = parse(&sender.packetize(&[0u8; 160], 160)).unwrap();
289        let second = parse(&sender.packetize(&[0u8; 160], 160)).unwrap();
290        assert!(first.marker, "first packet starts a talkspurt");
291        assert!(!second.marker);
292        assert_eq!(second.sequence, 101);
293        assert_eq!(second.timestamp, 1160);
294
295        sender.skip_silence(800); // 100 ms of silence
296        let resumed = parse(&sender.packetize(&[0u8; 160], 160)).unwrap();
297        assert!(resumed.marker, "resuming after silence re-marks");
298        // 1000 + 2×160 (sent) + 800 (skipped) = 2120.
299        assert_eq!(resumed.timestamp, 2120);
300    }
301
302    #[test]
303    fn rtp_packet_with_no_payload() {
304        // Valid RTP packet with zero-length payload
305        let packet = RtpPacket {
306            payload_type: PT_PCMU,
307            marker: false,
308            sequence: 100,
309            timestamp: 1000,
310            ssrc: 42,
311            payload: Vec::new(),
312        };
313        let built = build(&packet);
314        let parsed = parse(&built).expect("should parse empty payload");
315        assert_eq!(parsed.payload, Vec::<u8>::new());
316    }
317
318    #[test]
319    fn rtp_sequence_wrapping() {
320        // Sequence numbers should wrap at u16::MAX
321        let mut sender = RtpSender::new(PT_PCMU, 1, u16::MAX - 1, 0);
322        let p1 = parse(&sender.packetize(&[0u8; 160], 160)).unwrap();
323        let p2 = parse(&sender.packetize(&[0u8; 160], 160)).unwrap();
324        assert_eq!(p1.sequence, u16::MAX - 1);
325        assert_eq!(p2.sequence, u16::MAX);
326        let p3 = parse(&sender.packetize(&[0u8; 160], 160)).unwrap();
327        assert_eq!(p3.sequence, 0, "sequence should wrap to 0");
328    }
329
330    #[test]
331    fn rtp_timestamp_wrapping() {
332        // Timestamps should wrap at u32::MAX
333        let mut sender = RtpSender::new(PT_PCMU, 1, 0, u32::MAX - 100);
334        let p1 = parse(&sender.packetize(&[0u8; 160], 160)).unwrap();
335        assert_eq!(p1.timestamp, u32::MAX - 100);
336        let p2 = parse(&sender.packetize(&[0u8; 160], 160)).unwrap();
337        let expected = (u32::MAX as u64 - 100 + 160) as u32;
338        assert_eq!(p2.timestamp, expected, "timestamp should wrap correctly");
339    }
340
341    #[test]
342    fn rtp_csrc_count_limits() {
343        // Max CSRC count is 15 (4-bit field)
344        let mut bytes = vec![0x8F, 0x00, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0]; // V=2, CC=15, no P/X
345        for _ in 0..15 {
346            bytes.extend_from_slice(&[0, 0, 0, 0]); // 15 CSRCs
347        }
348        bytes.extend_from_slice(&[1, 2]); // payload
349        let packet = parse(&bytes).expect("should parse max CSRCs");
350        assert_eq!(packet.payload, vec![1, 2]);
351    }
352
353    #[test]
354    fn rtp_extension_header_handling() {
355        // RTP extension with multiple 4-byte words
356        let mut bytes = vec![0x90, 0x00, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0]; // no CSRC, has ext
357        bytes.extend_from_slice(&[0xAB, 0xCD, 0x00, 0x04]); // ext: profile, 4 words
358        for _ in 0..4 {
359            bytes.extend_from_slice(&[0xFF, 0xEE, 0xDD, 0xCC]);
360        }
361        bytes.extend_from_slice(&[0x12, 0x34]); // payload
362        let packet = parse(&bytes).expect("should parse extension");
363        assert_eq!(packet.payload, vec![0x12, 0x34]);
364    }
365
366    #[test]
367    fn rtp_padding_edge_cases() {
368        // Padding: payload [1, 2, 3] followed by 3-byte padding (includes count)
369        let mut bytes = vec![0xA0, 0x00, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0]; // P=1, X=0, CC=0
370        bytes.extend_from_slice(&[1, 2]); // payload
371        bytes.extend_from_slice(&[0, 0, 3]); // 3 bytes padding (0, 0, and count byte 3)
372        let packet = parse(&bytes).expect("should parse when pad equals end");
373        assert_eq!(packet.payload, vec![1, 2]); // last 3 bytes removed as padding
374    }
375
376    #[test]
377    fn rtp_parse_rejects_invalid_padding() {
378        // Padding count of 0 is invalid
379        let mut bytes = vec![0xA0, 0x00, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0]; // P=1
380        bytes.extend_from_slice(&[1, 2]);
381        bytes.push(0); // invalid: padding count must be >= 1
382        assert_eq!(parse(&bytes), None, "should reject zero padding count");
383
384        // Padding extends beyond datagram
385        let mut bytes = vec![0xA0, 0x00, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0]; // P=1
386        bytes.extend_from_slice(&[1, 2]);
387        bytes.push(10); // padding count = 10, but only 3 bytes follow the header
388        assert_eq!(parse(&bytes), None, "should reject oversized padding");
389    }
390
391    #[test]
392    fn telephone_event_all_digits() {
393        // Test all standard DTMF codes (0-15)
394        for event_code in 0u8..=15 {
395            let payload = [event_code, 0x00, 0x00, 0xA0];
396            let event = parse_telephone_event(&payload).unwrap();
397            assert_eq!(event.event, event_code);
398            assert!(!event.end);
399            // Verify digit() returns appropriate char for 0-11, None for 12-15
400            match event_code {
401                0..=9 => {
402                    assert!(event.digit().is_some());
403                }
404                10..=11 => {
405                    assert!(event.digit().is_some());
406                }
407                12..=15 => {
408                    assert!(event.digit().is_some());
409                }
410                _ => unreachable!(),
411            }
412        }
413    }
414
415    #[test]
416    fn telephone_event_duration_limits() {
417        // Duration is u16, test edge values
418        let payload = [5, 0x80, 0xFF, 0xFF]; // max duration
419        let event = parse_telephone_event(&payload).unwrap();
420        assert_eq!(event.duration, u16::MAX);
421
422        let payload = [5, 0x80, 0x00, 0x01]; // min non-zero duration
423        let event = parse_telephone_event(&payload).unwrap();
424        assert_eq!(event.duration, 1);
425    }
426
427    #[test]
428    fn payload_type_only_uses_7_bits() {
429        // Payload type is 7 bits; marker is bit 7
430        let packet = RtpPacket {
431            payload_type: 0x7F, // max 7-bit value
432            marker: true,
433            sequence: 0,
434            timestamp: 0,
435            ssrc: 0,
436            payload: vec![1, 2],
437        };
438        let built = build(&packet);
439        let parsed = parse(&built).unwrap();
440        assert_eq!(parsed.payload_type, 0x7F);
441        assert!(parsed.marker);
442    }
443}