From 50f5c0e10a141e6573b53e0f13ab67583349ad4b Mon Sep 17 00:00:00 2001 From: Repin Agent Date: Fri, 17 Jul 2026 21:55:35 -0600 Subject: [PATCH] fix(rtp): discard muxed RTCP before ingress hygiene --- crates/asterisk-sip/src/rtp/mod.rs | 189 +++++++++++++++++++++++++++++ crates/asterisk-sip/src/sdp/mod.rs | 32 +++++ 2 files changed, 221 insertions(+) diff --git a/crates/asterisk-sip/src/rtp/mod.rs b/crates/asterisk-sip/src/rtp/mod.rs index 9083d2a..45d290e 100644 --- a/crates/asterisk-sip/src/rtp/mod.rs +++ b/crates/asterisk-sip/src/rtp/mod.rs @@ -717,6 +717,14 @@ impl RtpSession { let (len, src) = self.socket.recv_from(&mut buf).await?; let packet = &buf[..len]; + // RTCP can share this socket when the peer uses rtcp-mux. Demux + // it before RTP source/framing/PT/SSRC hygiene so a report cannot + // be counted as bad RTP or disturb the accepted media stream. + if is_rtcp_packet(packet) { + trace!(source = %src, "Discarding RTCP datagram on RTP socket"); + continue; + } + if self.remote_addr().is_some_and(|remote| remote != src) { self.stats.discarded_wrong_source.fetch_add(1, Ordering::Relaxed); trace!(source = %src, expected = ?self.remote_addr(), @@ -1613,6 +1621,187 @@ mod tests { assert_eq!(session.stats.snapshot().packets_received, 2); } + #[tokio::test] + async fn chime_rtcp_does_not_disturb_rtp_ingress() { + // Reserve RTP+1 in the harness to model Chime's discrete RTCP port and + // prove the RTP session neither binds nor consumes it. + let mut sockets = None; + for _ in 0..32 { + let candidate = RtpSession::bind("127.0.0.1:0".parse().unwrap()) + .await + .unwrap(); + let rtp_addr = candidate.local_addr().unwrap(); + let Some(rtcp_port) = rtp_addr.port().checked_add(1) else { + continue; + }; + let rtcp_addr = SocketAddr::new(rtp_addr.ip(), rtcp_port); + if let Ok(rtcp_socket) = UdpSocket::bind(rtcp_addr).await { + sockets = Some((candidate, rtcp_socket)); + break; + } + } + let (session, separate_rtcp_socket) = + sockets.expect("could not reserve consecutive RTP/RTCP harness ports"); + session.set_payload_type(0); + + let target = session.local_addr().unwrap(); + let separate_target = separate_rtcp_socket.local_addr().unwrap(); + assert_eq!(separate_target.port(), target.port() + 1); + + let chime = UdpSocket::bind("127.0.0.1:0").await.unwrap(); + let chime_addr = chime.local_addr().unwrap(); + session.set_remote_addr(chime_addr); + + let first_rtp = RtpHeader { + version: 2, + padding: false, + extension: false, + csrc_count: 0, + marker: false, + payload_type: 0, + sequence: 1, + timestamp: 160, + ssrc: 0x12345678, + }; + chime + .send_to(&build_rtp_packet(&first_rtp, &[0x7f; 160]), target) + .await + .unwrap(); + assert!(matches!( + session.recv_frame().await.unwrap(), + Frame::Voice { + rtp_timing: Some(RtpTiming { + sequence: 1, + timestamp: 160, + }), + .. + } + )); + assert_eq!(*session.inbound_ssrc.lock(), Some(first_rtp.ssrc)); + + // A minimal valid Sender Report (PT 200, 28 bytes). + let sender_report = [ + 0x80, 200, 0x00, 0x06, // V=2, RC=0, PT=SR, length=6 + 0x11, 0x22, 0x33, 0x44, // sender SSRC + 0x00, 0x00, 0x00, 0x01, // NTP seconds + 0x00, 0x00, 0x00, 0x02, // NTP fraction + 0x00, 0x00, 0x00, 0xa0, // RTP timestamp + 0x00, 0x00, 0x00, 0x01, // sender packet count + 0x00, 0x00, 0x00, 0xa0, // sender octet count + ]; + // A minimal Receiver Report (PT 201, no report blocks). + let receiver_report = [ + 0x80, 201, 0x00, 0x01, // V=2, RC=0, PT=RR, length=1 + 0x11, 0x22, 0x33, 0x44, // reporter SSRC + ]; + + let before_separate = session.stats.snapshot(); + chime + .send_to(&sender_report, separate_target) + .await + .unwrap(); + let mut separate_packet = [0u8; 64]; + let (separate_len, separate_source) = tokio::time::timeout( + std::time::Duration::from_secs(2), + separate_rtcp_socket.recv_from(&mut separate_packet), + ) + .await + .expect("separate RTCP receive timed out") + .unwrap(); + assert_eq!(separate_source, chime_addr); + assert_eq!(&separate_packet[..separate_len], &sender_report); + assert_eq!(session.stats.snapshot(), before_separate); + + // Model Chime's rtcp-mux stream: SR and RR are interleaved with real + // PCMU packets on the RTP socket. Neither RTCP datagram may enter RTP + // ingress hygiene or disturb the accepted SSRC latch. + chime.send_to(&sender_report, target).await.unwrap(); + let resumed_rtp = RtpHeader { + sequence: 2, + timestamp: 320, + ..first_rtp + }; + chime + .send_to(&build_rtp_packet(&resumed_rtp, &[0x55; 160]), target) + .await + .unwrap(); + + let resumed = tokio::time::timeout( + std::time::Duration::from_secs(2), + session.recv_frame(), + ) + .await + .expect("RTP after muxed RTCP timed out") + .unwrap(); + assert!(matches!( + resumed, + Frame::Voice { + rtp_timing: Some(RtpTiming { + sequence: 2, + timestamp: 320, + }), + .. + } + )); + + chime.send_to(&receiver_report, target).await.unwrap(); + let third_rtp = RtpHeader { + sequence: 3, + timestamp: 480, + ..resumed_rtp + }; + chime + .send_to(&build_rtp_packet(&third_rtp, &[0x33; 160]), target) + .await + .unwrap(); + let third = tokio::time::timeout( + std::time::Duration::from_secs(2), + session.recv_frame(), + ) + .await + .expect("RTP after muxed RTCP RR timed out") + .unwrap(); + assert!(matches!( + third, + Frame::Voice { + rtp_timing: Some(RtpTiming { + sequence: 3, + timestamp: 480, + }), + .. + } + )); + + let after_mux = session.stats.snapshot(); + assert_eq!(after_mux.packets_received, before_separate.packets_received + 2); + assert_eq!( + after_mux.voice_frames_received, + before_separate.voice_frames_received + 2 + ); + assert_eq!( + after_mux.octets_received, + before_separate.octets_received + 320 + ); + assert_eq!( + after_mux.discarded_wrong_source, + before_separate.discarded_wrong_source + ); + assert_eq!( + after_mux.discarded_wrong_payload_type, + before_separate.discarded_wrong_payload_type + ); + assert_eq!( + after_mux.discarded_malformed, + before_separate.discarded_malformed + ); + assert_eq!( + after_mux.discarded_unstable_ssrc, + before_separate.discarded_unstable_ssrc + ); + assert_eq!(after_mux.remote_addr, Some(chime_addr)); + assert_eq!(*session.inbound_ssrc.lock(), Some(third_rtp.ssrc)); + } + #[tokio::test] async fn send_frame_rebases_source_timing_and_preserves_gaps() { let session = RtpSession::bind("127.0.0.1:0".parse().unwrap()).await.unwrap(); diff --git a/crates/asterisk-sip/src/sdp/mod.rs b/crates/asterisk-sip/src/sdp/mod.rs index 22e2459..cf3a1d7 100644 --- a/crates/asterisk-sip/src/sdp/mod.rs +++ b/crates/asterisk-sip/src/sdp/mod.rs @@ -1530,6 +1530,38 @@ a=sendrecv\r\n"; assert_eq!(parsed.media_descriptions[0].port, 20000); } + #[test] + fn chime_rtcp_attributes_do_not_reject_offer_or_answer() { + let offer = SessionDescription::parse( + "v=0\r\n\ + o=- 1 1 IN IP4 10.0.0.1\r\n\ + s=Chime\r\n\ + c=IN IP4 10.0.0.1\r\n\ + t=0 0\r\n\ + m=audio 40000 RTP/AVP 0\r\n\ + a=rtpmap:0 PCMU/8000\r\n\ + a=rtcp:40001 IN IP4 10.0.0.1\r\n\ + a=rtcp-mux\r\n\ + a=sendrecv\r\n", + ) + .expect("Chime-style SDP offer must parse"); + + let offered_audio = &offer.media_descriptions[0]; + assert!(offered_audio.has_rtcp_mux()); + assert!(offered_audio.attributes.iter().any(|(name, value)| { + name == "rtcp" && value.as_deref() == Some("40001 IN IP4 10.0.0.1") + })); + + let answer = + SessionDescription::create_answer(&offer, "127.0.0.1", 55555, &[codecs_pcmu()]); + assert_eq!(answer.media_descriptions.len(), 1); + assert_eq!(answer.media_descriptions[0].port, 55555); + + let reparsed_answer = SessionDescription::parse(&answer.to_string()) + .expect("answer to Chime-style offer must serialize and parse"); + assert_eq!(reparsed_answer.media_descriptions[0].port, 55555); + } + #[test] fn test_sdp_parse_dtls_attributes() { let sdp_text = "v=0\r\n\