Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
123 changes: 123 additions & 0 deletions crates/asterisk-sip/src/channel_driver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -245,6 +245,44 @@ impl SipChannelDriver {
.map(|a| a.port())
}

/// Install (or clear) the telephone-event payload type from negotiated
/// remote SDP on an existing call.
pub(crate) async fn install_negotiated_dtmf_payload(
&self,
channel_name: &str,
sdp: &SessionDescription,
) -> AsteriskResult<()> {
let priv_data = self
.get_private(channel_name)
.ok_or_else(|| AsteriskError::NotFound(channel_name.to_string()))?;
let rtp = priv_data
.rtp
.lock()
.await
.clone()
.ok_or_else(|| AsteriskError::Internal("No RTP session".into()))?;

let negotiated = crate::sdp_rtp::negotiated_dtmf_payload_type(sdp, &self.codecs);
if let Some(payload_type) = negotiated {
rtp.set_dtmf_payload_type(payload_type);
} else {
rtp.clear_dtmf_payload_type();
}
debug!(
channel = channel_name,
payload_type = ?negotiated,
"Installed negotiated telephone-event payload type"
);
Ok(())
}

/// Return the channel's negotiated telephone-event payload type.
pub async fn channel_rtp_dtmf_payload_type(&self, channel_name: &str) -> Option<u8> {
let priv_data = self.get_private(channel_name)?;
let rtp = priv_data.rtp.lock().await.clone()?;
rtp.dtmf_payload_type()
}

fn get_transport(&self) -> AsteriskResult<Arc<dyn SipTransport>> {
self.transport.read().clone().ok_or_else(|| {
AsteriskError::Internal("SIP transport not initialized".into())
Expand Down Expand Up @@ -314,6 +352,13 @@ impl ChannelDriver for SipChannelDriver {

// Create RTP session
let rtp_session = self.allocate_rtp_session(self.local_addr.ip()).await?;
if let Some(codec) = self
.codecs
.iter()
.find(|codec| codec.name.eq_ignore_ascii_case("telephone-event"))
{
rtp_session.set_dtmf_payload_type(codec.payload_type);
}
let rtp_port = rtp_session.local_addr()?.port();

// Create SIP session
Expand Down Expand Up @@ -822,6 +867,84 @@ mod tests {
);
}

#[tokio::test]
async fn negotiated_dtmf_payload_controls_receiver_detection() {
let local: SocketAddr = "127.0.0.1:0".parse().unwrap();
let transport: Arc<dyn SipTransport> = Arc::new(
UdpTransport::bind(local).await.unwrap(),
);
let driver = SipChannelDriver::new(local);
driver.set_transport(transport);
let mut channel = driver
.request("sip:dtmf@127.0.0.1:9", None)
.await
.unwrap();
let rtp_port = driver
.channel_rtp_local_port(&channel.name)
.await
.unwrap();

let negotiated = SessionDescription::parse(
"v=0\r\n\
o=- 1 1 IN IP4 127.0.0.1\r\n\
s=Test\r\n\
c=IN IP4 127.0.0.1\r\n\
t=0 0\r\n\
m=audio 40000 RTP/AVP 0 110\r\n\
a=rtpmap:0 PCMU/8000\r\n\
a=rtpmap:110 telephone-event/8000\r\n",
)
.unwrap();
driver
.install_negotiated_dtmf_payload(&channel.name, &negotiated)
.await
.unwrap();
assert_eq!(
driver.channel_rtp_dtmf_payload_type(&channel.name).await,
Some(110)
);

let peer = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
let target: SocketAddr = format!("127.0.0.1:{}", rtp_port).parse().unwrap();
let event = crate::rtp::DtmfEvent {
event: 5,
end: true,
volume: 10,
duration: 800,
};
let packet = |payload_type| {
crate::rtp::build_rtp_packet(
&crate::rtp::RtpHeader {
version: 2,
padding: false,
extension: false,
csrc_count: 0,
marker: false,
payload_type,
sequence: 1,
timestamp: 160,
ssrc: 0x12345678,
},
&event.to_bytes(),
)
};

peer.send_to(&packet(101), target).await.unwrap();
assert!(matches!(
driver.read_frame(&mut channel).await.unwrap(),
Frame::Voice { .. }
));

peer.send_to(&packet(110), target).await.unwrap();
assert!(matches!(
driver.read_frame(&mut channel).await.unwrap(),
Frame::DtmfEnd {
digit: '5',
duration_ms: 100
}
));
}

/// Regression for the channel-name collision bug: the inbound INVITE path
/// previously derived its channel-name suffix from a truncated
/// `SystemTime::now()` nanosecond value ("rand_id"), which is not random and
Expand Down
27 changes: 27 additions & 0 deletions crates/asterisk-sip/src/event_handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -465,6 +465,12 @@ impl SipEventHandler {
if let Some(pt) = negotiated_payload_type(&remote_sdp, &self.supported_codecs) {
rtp.payload_type = pt;
}
if let Some(pt) = crate::sdp_rtp::negotiated_dtmf_payload_type(
&remote_sdp,
&self.supported_codecs,
) {
rtp.set_dtmf_payload_type(pt);
}
// Store the REAL inbound session (carries the INVITE,
// is_outbound = false) so driver.indicate()/hangup()
// work on this channel instead of silently no-opping on
Expand Down Expand Up @@ -926,8 +932,10 @@ impl SipEventHandler {
};
if let Some(cs_arc) = cs_arc {
let mut cs = cs_arc.lock().await;
let mut negotiated_sdp = None;
if cs.session.is_outbound {
cs.session.on_response(response);
negotiated_sdp = cs.session.remote_sdp.clone();
if let Some(ack) = cs.session.build_ack() {
if let Err(e) = self.transport.send(&ack, remote_addr).await {
warn!(call_id = %call_id, "Failed to send ACK: {}", e);
Expand All @@ -938,6 +946,22 @@ impl SipEventHandler {
eprintln!("[DEBUG] Failed to build ACK for call_id={}", call_id);
}
}
drop(cs);

if let (Some(driver), Some(sdp)) =
(self.channel_driver.get(), negotiated_sdp.as_ref())
{
if let Err(e) = driver
.install_negotiated_dtmf_payload(&channel_name, sdp)
.await
{
warn!(
call_id = %call_id,
"Failed to install negotiated telephone-event payload type: {}",
e
);
}
}
}
}
}
Expand Down Expand Up @@ -1537,6 +1561,9 @@ fn negotiated_payload_type(sdp: &SessionDescription, supported: &[Codec]) -> Opt
.iter()
.find(|m| m.media_type == "audio")?;
for oc in media.codecs() {
if oc.name.eq_ignore_ascii_case("telephone-event") {
continue;
}
for sc in supported {
if oc.name.eq_ignore_ascii_case(&sc.name) && oc.sample_rate == sc.sample_rate {
return Some(oc.payload_type);
Expand Down
51 changes: 40 additions & 11 deletions crates/asterisk-sip/src/rtp/mod.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
//! RTP/RTCP session management.
//!
//! Provides RTP send/receive with proper header handling, payload type
//! mapping, and RFC 2833 DTMF support. Also includes RTCP SR/RR.
//! mapping, and RFC 4733 DTMF support. Also includes RTCP SR/RR.
//!
//! Sub-modules:
//! - `jitter_buffer`: Fixed and adaptive jitter buffer implementations.
Expand All @@ -19,7 +19,7 @@ pub mod mos;

use std::net::SocketAddr;
use std::path::Path;
use std::sync::atomic::{AtomicU16, AtomicU32, Ordering};
use std::sync::atomic::{AtomicU8, AtomicU16, AtomicU32, Ordering};
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};

Expand All @@ -38,6 +38,8 @@ const RTP_MAX_MTU: usize = 1500;
const RTCP_PT_SR: u8 = 200;
/// RTCP receiver report type.
const RTCP_PT_RR: u8 = 201;
/// Sentinel outside RTP's 7-bit payload type space: no telephone-event map.
const NO_DTMF_PAYLOAD_TYPE: u8 = u8::MAX;

/// Default inclusive RTP range used when `rtp.conf` is absent.
pub const DEFAULT_RTP_PORT_START: u16 = 10000;
Expand Down Expand Up @@ -254,7 +256,7 @@ pub fn parse_rtp_header(data: &[u8]) -> Result<(RtpHeader, &[u8]), AsteriskError
Ok((header, &data[offset..]))
}

/// RFC 2833 DTMF event payload.
/// RFC 4733 DTMF event payload.
#[derive(Debug, Clone)]
pub struct DtmfEvent {
pub event: u8,
Expand Down Expand Up @@ -332,8 +334,8 @@ pub struct RtpSession {
timestamp: AtomicU32,
/// Payload type for outgoing packets.
pub payload_type: u8,
/// DTMF payload type (RFC 2833).
pub dtmf_payload_type: u8,
/// Negotiated telephone-event payload type (RFC 4733).
dtmf_payload_type: AtomicU8,
/// Samples per packet (for timestamp advancement).
pub samples_per_packet: u32,
/// Statistics.
Expand Down Expand Up @@ -365,7 +367,7 @@ impl RtpSession {
sequence: AtomicU16::new(0),
timestamp: AtomicU32::new(0),
payload_type: 0,
dtmf_payload_type: 101,
dtmf_payload_type: AtomicU8::new(NO_DTMF_PAYLOAD_TYPE),
samples_per_packet: 160,
stats: RtpStats::default(),
}
Expand All @@ -389,6 +391,30 @@ impl RtpSession {
*self.remote_addr.write() = Some(addr);
}

/// Return the negotiated RFC 4733 telephone-event payload type.
pub fn dtmf_payload_type(&self) -> Option<u8> {
match self.dtmf_payload_type.load(Ordering::Relaxed) {
NO_DTMF_PAYLOAD_TYPE => None,
payload_type => Some(payload_type),
}
}

/// Install the dynamic telephone-event payload type negotiated in SDP.
pub fn set_dtmf_payload_type(&self, payload_type: u8) {
let payload_type = if payload_type <= 0x7f {
payload_type
} else {
NO_DTMF_PAYLOAD_TYPE
};
self.dtmf_payload_type.store(payload_type, Ordering::Relaxed);
}

/// Disable RFC 4733 send/receive when SDP did not negotiate it.
pub fn clear_dtmf_payload_type(&self) {
self.dtmf_payload_type
.store(NO_DTMF_PAYLOAD_TYPE, Ordering::Relaxed);
}

/// Send an audio frame as RTP.
pub async fn send_frame(&self, frame: &Frame) -> AsteriskResult<()> {
let data = match frame {
Expand Down Expand Up @@ -473,8 +499,8 @@ impl RtpSession {
.octets_received
.fetch_add(payload.len() as u32, Ordering::Relaxed);

// Check for DTMF (RFC 2833)
if header.payload_type == self.dtmf_payload_type {
// Check for DTMF (RFC 4733)
if self.dtmf_payload_type() == Some(header.payload_type) {
if let Some(event) = DtmfEvent::from_bytes(payload) {
let digit = DtmfEvent::event_to_digit(event.event);
if event.end {
Expand All @@ -497,7 +523,7 @@ impl RtpSession {
))
}

/// Send a DTMF digit via RFC 2833.
/// Send a DTMF digit via RFC 4733.
pub async fn send_dtmf(
&self,
digit: char,
Expand All @@ -506,6 +532,9 @@ impl RtpSession {
let remote = self
.remote_addr()
.ok_or_else(|| AsteriskError::InvalidArgument("No remote address".into()))?;
let dtmf_payload_type = self.dtmf_payload_type().ok_or_else(|| {
AsteriskError::InvalidArgument("telephone-event was not negotiated".into())
})?;

let event_num = DtmfEvent::digit_to_event(digit);
let start_seq = self.sequence.fetch_add(1, Ordering::Relaxed);
Expand All @@ -525,7 +554,7 @@ impl RtpSession {
extension: false,
csrc_count: 0,
marker: i == 0,
payload_type: self.dtmf_payload_type,
payload_type: dtmf_payload_type,
sequence: start_seq.wrapping_add(i),
timestamp: start_ts,
ssrc: self.ssrc,
Expand All @@ -548,7 +577,7 @@ impl RtpSession {
extension: false,
csrc_count: 0,
marker: false,
payload_type: self.dtmf_payload_type,
payload_type: dtmf_payload_type,
sequence: start_seq.wrapping_add(3 + i),
timestamp: start_ts,
ssrc: self.ssrc,
Expand Down
Loading
Loading