use crate::peer_connection::configuration::media_engine::MediaEngine;
use crate::rtp_transceiver::rtp_sender::rtcp_parameters::{
TYPE_RTCP_FB_ACK, TYPE_RTCP_FB_NACK, TYPE_RTCP_FB_TRANSPORT_CC,
};
use crate::rtp_transceiver::rtp_sender::{
RTCPFeedback, RTCRtpCodec, RTCRtpHeaderExtensionCapability, RTCRtpHeaderExtensionParameters,
RtpCodecKind,
};
use crate::rtp_transceiver::{PayloadType, SSRC};
use interceptor::{
BandwidthEstimator, CongestionControlBuilder, NackGeneratorBuilder, NackResponderBuilder,
PacerBuilder, ReceiverReportBuilder, Rfc8888Builder, SenderReportBuilder, TwccReceiverBuilder,
TwccSenderBuilder,
};
pub use interceptor::{Registry, Slot};
use shared::error::Result;
pub fn register_default_interceptors(
registry: Registry,
media_engine: &mut MediaEngine,
) -> Result<Registry> {
let registry = configure_nack(registry, media_engine);
configure_simulcast_extension_headers(media_engine)?;
let registry = configure_twcc_receiver_only(registry, media_engine)?;
let registry = configure_rtcp_reports(registry);
Ok(registry)
}
pub fn configure_nack(registry: Registry, media_engine: &mut MediaEngine) -> Registry {
media_engine.register_feedback(
RTCPFeedback {
typ: TYPE_RTCP_FB_NACK.to_owned(),
parameter: "".to_owned(),
},
RtpCodecKind::Video,
);
media_engine.register_feedback(
RTCPFeedback {
typ: TYPE_RTCP_FB_NACK.to_owned(),
parameter: "pli".to_owned(),
},
RtpCodecKind::Video,
);
registry
.with(Slot::NackResponder, NackResponderBuilder::new().build())
.with(Slot::NackGenerator, NackGeneratorBuilder::new().build())
}
pub fn configure_rtcp_reports(registry: Registry) -> Registry {
registry
.with(Slot::ReceiverReport, ReceiverReportBuilder::new().build())
.with(Slot::SenderReport, SenderReportBuilder::new().build())
}
pub fn configure_simulcast_extension_headers(media_engine: &mut MediaEngine) -> Result<()> {
media_engine.register_header_extension(
RTCRtpHeaderExtensionCapability {
uri: ::sdp::extmap::SDES_MID_URI.to_owned(),
},
RtpCodecKind::Video,
None,
)?;
media_engine.register_header_extension(
RTCRtpHeaderExtensionCapability {
uri: ::sdp::extmap::SDES_RTP_STREAM_ID_URI.to_owned(),
},
RtpCodecKind::Video,
None,
)?;
media_engine.register_header_extension(
RTCRtpHeaderExtensionCapability {
uri: ::sdp::extmap::SDES_REPAIR_RTP_STREAM_ID_URI.to_owned(),
},
RtpCodecKind::Video,
None,
)?;
Ok(())
}
pub fn configure_twcc(registry: Registry, media_engine: &mut MediaEngine) -> Result<Registry> {
media_engine.register_feedback(
RTCPFeedback {
typ: TYPE_RTCP_FB_TRANSPORT_CC.to_owned(),
..Default::default()
},
RtpCodecKind::Video,
);
media_engine.register_header_extension(
RTCRtpHeaderExtensionCapability {
uri: sdp::extmap::TRANSPORT_CC_URI.to_owned(),
},
RtpCodecKind::Video,
None,
)?;
media_engine.register_feedback(
RTCPFeedback {
typ: TYPE_RTCP_FB_TRANSPORT_CC.to_owned(),
..Default::default()
},
RtpCodecKind::Audio,
);
media_engine.register_header_extension(
RTCRtpHeaderExtensionCapability {
uri: sdp::extmap::TRANSPORT_CC_URI.to_owned(),
},
RtpCodecKind::Audio,
None,
)?;
Ok(registry
.with(Slot::TwccSender, TwccSenderBuilder::new().build())
.with(Slot::TwccReceiver, TwccReceiverBuilder::new().build()))
}
pub fn configure_twcc_sender_only(
registry: Registry,
media_engine: &mut MediaEngine,
) -> Result<Registry> {
media_engine.register_feedback(
RTCPFeedback {
typ: TYPE_RTCP_FB_TRANSPORT_CC.to_owned(),
parameter: "".to_owned(),
},
RtpCodecKind::Video,
);
media_engine.register_header_extension(
RTCRtpHeaderExtensionCapability {
uri: sdp::extmap::TRANSPORT_CC_URI.to_owned(),
},
RtpCodecKind::Video,
None,
)?;
media_engine.register_feedback(
RTCPFeedback {
typ: TYPE_RTCP_FB_TRANSPORT_CC.to_owned(),
parameter: "".to_owned(),
},
RtpCodecKind::Audio,
);
media_engine.register_header_extension(
RTCRtpHeaderExtensionCapability {
uri: sdp::extmap::TRANSPORT_CC_URI.to_owned(),
},
RtpCodecKind::Audio,
None,
)?;
Ok(registry.with(Slot::TwccSender, TwccSenderBuilder::new().build()))
}
pub fn configure_twcc_receiver_only(
registry: Registry,
media_engine: &mut MediaEngine,
) -> Result<Registry> {
media_engine.register_feedback(
RTCPFeedback {
typ: TYPE_RTCP_FB_TRANSPORT_CC.to_owned(),
..Default::default()
},
RtpCodecKind::Video,
);
media_engine.register_header_extension(
RTCRtpHeaderExtensionCapability {
uri: sdp::extmap::TRANSPORT_CC_URI.to_owned(),
},
RtpCodecKind::Video,
None,
)?;
media_engine.register_feedback(
RTCPFeedback {
typ: TYPE_RTCP_FB_TRANSPORT_CC.to_owned(),
..Default::default()
},
RtpCodecKind::Audio,
);
media_engine.register_header_extension(
RTCRtpHeaderExtensionCapability {
uri: sdp::extmap::TRANSPORT_CC_URI.to_owned(),
},
RtpCodecKind::Audio,
None,
)?;
Ok(registry.with(Slot::TwccReceiver, TwccReceiverBuilder::new().build()))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[non_exhaustive]
pub enum CongestionFeedback {
#[default]
Twcc,
Rfc8888,
}
pub fn configure_congestion_control<E: BandwidthEstimator + 'static>(
registry: Registry,
estimator: E,
feedback: CongestionFeedback,
media_engine: &mut MediaEngine,
) -> Result<Registry> {
for kind in [RtpCodecKind::Video, RtpCodecKind::Audio] {
match feedback {
CongestionFeedback::Twcc => {
media_engine.register_feedback(
RTCPFeedback {
typ: TYPE_RTCP_FB_TRANSPORT_CC.to_owned(),
parameter: "".to_owned(),
},
kind,
);
media_engine.register_header_extension(
RTCRtpHeaderExtensionCapability {
uri: sdp::extmap::TRANSPORT_CC_URI.to_owned(),
},
kind,
None,
)?;
}
CongestionFeedback::Rfc8888 => {
media_engine.register_feedback(
RTCPFeedback {
typ: TYPE_RTCP_FB_ACK.to_owned(),
parameter: "ccfb".to_owned(),
},
kind,
);
}
}
}
let registry = registry
.with(
Slot::CongestionControl,
CongestionControlBuilder::new(estimator).build(),
)
.with(Slot::Pacer, PacerBuilder::new().build());
Ok(match feedback {
CongestionFeedback::Twcc => registry
.with(Slot::TwccSender, TwccSenderBuilder::new().build())
.with(Slot::TwccReceiver, TwccReceiverBuilder::new().build()),
CongestionFeedback::Rfc8888 => registry.with(Slot::Rfc8888, Rfc8888Builder::new().build()),
})
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn create_stream_info(
ssrc: SSRC,
ssrc_rtx: Option<SSRC>,
ssrc_fec: Option<SSRC>,
payload_type: PayloadType,
payload_type_rtx: Option<PayloadType>,
payload_type_fec: Option<PayloadType>,
codec: &RTCRtpCodec,
header_extensions: &[RTCRtpHeaderExtensionParameters],
) -> interceptor::StreamInfo {
let rtp_header_extensions: Vec<interceptor::RTPHeaderExtension> = header_extensions
.iter()
.map(|h| interceptor::RTPHeaderExtension {
id: h.id,
uri: h.uri.clone(),
})
.collect();
let feedbacks: Vec<_> = codec
.rtcp_feedback
.iter()
.map(|f| interceptor::RTCPFeedback {
typ: f.typ.clone(),
parameter: f.parameter.clone(),
})
.collect();
interceptor::StreamInfo {
ssrc,
ssrc_rtx,
ssrc_fec,
payload_type,
payload_type_rtx,
payload_type_fec,
rtp_header_extensions,
mime_type: codec.mime_type.clone(),
clock_rate: codec.clock_rate,
channels: codec.channels,
sdp_fmtp_line: codec.sdp_fmtp_line.clone(),
rtcp_feedback: feedbacks,
}
}
#[cfg(test)]
mod slot_order_tests {
use super::*;
fn slots(registry: &Registry) -> Vec<Slot> {
registry.slots().into_iter().map(|(slot, _)| slot).collect()
}
#[test]
fn helper_call_order_does_not_affect_the_chain() {
let mut media_engine = MediaEngine::default();
let nack_first = configure_twcc(
configure_nack(Registry::new(), &mut media_engine),
&mut media_engine,
)
.expect("twcc");
let mut media_engine = MediaEngine::default();
let twcc_first = configure_nack(
configure_twcc(Registry::new(), &mut media_engine).expect("twcc"),
&mut media_engine,
);
assert_eq!(slots(&nack_first), slots(&twcc_first));
}
#[test]
fn the_default_chain_is_emitted_in_slot_order() {
let mut media_engine = MediaEngine::default();
let registry = register_default_interceptors(Registry::new(), &mut media_engine)
.expect("default interceptors");
let slots = slots(®istry);
assert!(
slots.windows(2).all(|pair| pair[0] <= pair[1]),
"default chain is not wire-to-application: {slots:?}"
);
assert_eq!(
vec![
Slot::NackResponder,
Slot::NackGenerator,
Slot::TwccReceiver,
Slot::ReceiverReport,
Slot::SenderReport,
],
slots
);
}
#[test]
fn congestion_control_occupies_its_three_slots_in_order() {
use interceptor::Gcc;
let mut media_engine = MediaEngine::default();
let registry = configure_nack(Registry::new(), &mut media_engine);
let registry = configure_congestion_control(
registry,
Gcc::default(),
CongestionFeedback::Twcc,
&mut media_engine,
)
.expect("congestion control");
assert_eq!(
vec![
Slot::CongestionControl,
Slot::TwccSender,
Slot::Pacer,
Slot::NackResponder,
Slot::NackGenerator,
Slot::TwccReceiver,
],
slots(®istry)
);
}
#[test]
fn rfc8888_does_not_also_install_the_twcc_sender() {
use interceptor::Gcc;
let mut media_engine = MediaEngine::default();
let registry = configure_congestion_control(
Registry::new(),
Gcc::default(),
CongestionFeedback::Rfc8888,
&mut media_engine,
)
.expect("congestion control");
assert_eq!(
vec![Slot::CongestionControl, Slot::Pacer, Slot::Rfc8888],
slots(®istry),
"RFC 8888 needs no transport-wide sequence numbers"
);
}
#[test]
fn the_default_chain_has_no_congestion_control() {
let mut media_engine = MediaEngine::default();
let registry = register_default_interceptors(Registry::new(), &mut media_engine)
.expect("default interceptors");
let slots = slots(®istry);
assert!(
!slots.contains(&(Slot::CongestionControl)),
"no estimator by default: {slots:?}"
);
assert!(
!slots.contains(&(Slot::Pacer)),
"and no pacer by default: {slots:?}"
);
}
#[test]
fn a_custom_interceptor_fits_between_named_slots() {
let mut media_engine = MediaEngine::default();
let registry = configure_twcc(Registry::new(), &mut media_engine)
.expect("twcc")
.with(Slot::from(2_500), interceptor::NoopInterceptor::new());
assert_eq!(
vec![Slot::TwccSender, Slot::from(2_500), Slot::TwccReceiver],
slots(®istry)
);
}
#[test]
fn a_slot_holds_one_interceptor() {
let registry = Registry::new()
.with(Slot::Pacer, interceptor::NoopInterceptor::new())
.with(Slot::Pacer, interceptor::NoopInterceptor::new());
assert_eq!(vec![Slot::Pacer], slots(®istry));
}
#[test]
fn the_default_chain_names_what_it_registered() {
let mut media_engine = MediaEngine::default();
let registry = register_default_interceptors(Registry::new(), &mut media_engine)
.expect("default interceptors");
assert_eq!(
vec![
(Slot::NackResponder, "NackResponderInterceptor".to_owned()),
(Slot::NackGenerator, "NackGeneratorInterceptor".to_owned()),
(Slot::TwccReceiver, "TwccReceiverInterceptor".to_owned()),
(Slot::ReceiverReport, "ReceiverReportInterceptor".to_owned()),
(Slot::SenderReport, "SenderReportInterceptor".to_owned()),
],
registry.slots()
);
}
#[test]
fn the_twcc_recorder_is_not_registered_twice() {
use interceptor::Gcc;
let mut media_engine = MediaEngine::default();
let registry = register_default_interceptors(Registry::new(), &mut media_engine)
.expect("default interceptors");
let registry = configure_congestion_control(
registry,
Gcc::default(),
CongestionFeedback::Twcc,
&mut media_engine,
)
.expect("congestion control");
let recorders = registry
.slots()
.into_iter()
.filter(|(slot, _)| *slot == Slot::TwccReceiver)
.count();
assert_eq!(
1, recorders,
"one arrival recorder, however many helpers asked for one"
);
}
#[test]
fn rfc8888_alongside_the_defaults_leaves_two_recorders() {
use interceptor::Gcc;
let mut media_engine = MediaEngine::default();
let registry = register_default_interceptors(Registry::new(), &mut media_engine)
.expect("default interceptors");
let registry = configure_congestion_control(
registry,
Gcc::default(),
CongestionFeedback::Rfc8888,
&mut media_engine,
)
.expect("congestion control");
let recorders: Vec<Slot> = registry
.slots()
.into_iter()
.map(|(slot, _)| slot)
.filter(|slot| *slot == Slot::TwccReceiver || *slot == Slot::Rfc8888)
.collect();
assert_eq!(
vec![Slot::TwccReceiver, Slot::Rfc8888],
recorders,
"known gap: different slots, so nothing de-duplicates them — an RFC 8888 chain has to \
be built without `register_default_interceptors`, or `Registry` needs a way to drop a \
slot"
);
}
}