pub mod ice;
pub mod room;
pub mod rtcp;
pub mod sdp;
pub use ice::{parse_trickle, IceCandidate};
pub use room::{DominantSpeaker, Room};
pub use sdp::{MediaDirection, SdpAnswerParams, SdpOffer};
use crate::bus::PlaybackRegistry;
use crate::inbound::{IngestContext, PublishSession};
#[cfg(feature = "codec-av1")]
use crate::protocol::rtp::Av1Packetizer;
use crate::protocol::rtp::{
H264Depacketizer, OpusPacketizer, RtpHeader, RtpPacketizer, Vp9Packetizer,
};
use crate::{CodecId, MediaFrame, Result, StreamKey};
use async_trait::async_trait;
use std::sync::Arc;
#[derive(Debug, Clone, Default, PartialEq)]
pub struct PeerStats {
pub estimated_bitrate_bps: Option<u64>,
pub rtt_ms: Option<f32>,
pub egress_loss: Option<f32>,
}
#[async_trait]
pub trait DtlsSrtpTransport: Send + Sync {
fn fingerprint(&self) -> String;
fn ice_credentials(&self) -> (String, String);
async fn recv_rtp(&self) -> Option<Vec<u8>> {
None
}
async fn send_rtp(&self, _packet: &[u8]) -> Result<()> {
Ok(())
}
async fn send_rtcp(&self, packet: &[u8]) -> Result<()>;
async fn recv_rtcp(&self) -> Option<Vec<u8>> {
None
}
fn estimated_bitrate(&self) -> Option<u64> {
None
}
fn peer_stats(&self) -> Option<PeerStats> {
None
}
async fn add_remote_candidate(&self, _candidate: &str) -> Result<()> {
Ok(())
}
async fn recv_data(&self) -> Option<(String, Vec<u8>)> {
None
}
async fn send_data(&self, _label: &str, _data: &[u8]) -> Result<()> {
Ok(())
}
fn answer(&self, offer_sdp: &str, direction: MediaDirection) -> String {
let Some(offer) = SdpOffer::parse(offer_sdp) else {
return String::new();
};
let (ice_ufrag, ice_pwd) = self.ice_credentials();
sdp::build_answer_directed(
&offer,
&SdpAnswerParams {
fingerprint: self.fingerprint(),
ice_ufrag,
ice_pwd,
},
direction,
)
}
}
#[derive(Clone)]
pub struct WhipEndpoint {
ctx: IngestContext,
}
impl WhipEndpoint {
pub fn new(ctx: IngestContext) -> Self {
Self { ctx }
}
pub fn accept_offer(
&self,
offer_sdp: &str,
key: StreamKey,
transport: std::sync::Arc<dyn DtlsSrtpTransport>,
) -> Result<(WhipResource, String)> {
let offer = SdpOffer::parse(offer_sdp)
.ok_or_else(|| crate::StreamError::protocol("malformed SDP offer"))?;
let answer = transport.answer(offer_sdp, MediaDirection::RecvOnly);
let resource = WhipResource {
ctx: self.ctx.clone(),
key,
transport,
video_pt: offer.payload_type,
audio_pt: offer.audio_payload_type,
rid_ext_id: offer.rid_ext_id,
simulcast_rids: offer.simulcast_rids,
};
Ok((resource, answer))
}
}
pub struct WhipResource {
ctx: IngestContext,
key: StreamKey,
transport: std::sync::Arc<dyn DtlsSrtpTransport>,
video_pt: u8,
audio_pt: Option<u8>,
rid_ext_id: Option<u8>,
simulcast_rids: Vec<String>,
}
impl WhipResource {
pub async fn pump(self) -> Result<()> {
match self.rid_ext_id {
Some(ext) if self.simulcast_rids.len() > 1 => self.pump_simulcast(ext).await,
_ => self.pump_single().await,
}
}
async fn pump_single(self) -> Result<()> {
let session: PublishSession = self.ctx.open_publish(self.key.clone()).await?;
let handle = session.handle().clone();
let mut depack = H264Depacketizer::new();
let mut needs_keyframe = true;
let mut last_ssrc = 0u32;
loop {
let pkt = tokio::select! {
pkt = self.transport.recv_rtp() => match pkt {
Some(p) => p,
None => break,
},
_ = handle.keyframe_requested() => {
let pli = rtcp::build_pli(0, last_ssrc);
let _ = self.transport.send_rtcp(&pli).await;
continue;
}
};
let Some(header) = RtpHeader::parse(&pkt) else {
continue;
};
last_ssrc = header.ssrc;
let payload = &pkt[header.payload_offset..];
if self.audio_pt == Some(header.payload_type) {
if !payload.is_empty() {
let pts = (header.timestamp / 48) as i64;
let data = bytes::Bytes::copy_from_slice(payload);
let frame = MediaFrame::new_audio(pts, data, CodecId::Opus);
let _ = session.publish_frame(frame)?;
}
continue;
}
let _ = self.video_pt; match depack.push(payload, header.marker, header.timestamp, header.sequence) {
Ok(Some(au)) => {
needs_keyframe = false;
let pts = (au.timestamp / 90) as i64;
let frame =
MediaFrame::new_video(pts, pts, au.data, CodecId::H264, au.keyframe);
let _ = session.publish_frame(frame)?;
}
Ok(None) => {}
Err(_) => {
needs_keyframe = true;
}
}
if needs_keyframe {
let pli = rtcp::build_pli(0, header.ssrc);
let _ = self.transport.send_rtcp(&pli).await;
}
}
session.finish().await
}
async fn pump_simulcast(self, rid_ext: u8) -> Result<()> {
use std::collections::HashMap;
struct Layer {
session: PublishSession,
depack: H264Depacketizer,
needs_keyframe: bool,
}
let base = self.simulcast_rids[0].clone();
let mut layers: HashMap<String, Layer> = HashMap::new();
while let Some(pkt) = self.transport.recv_rtp().await {
let Some(header) = RtpHeader::parse(&pkt) else {
continue;
};
let rid = crate::protocol::rtp::rtp_extension_value(&pkt, rid_ext)
.and_then(|b| std::str::from_utf8(b).ok())
.map(str::to_owned)
.unwrap_or_else(|| base.clone());
if !self.simulcast_rids.contains(&rid) {
continue; }
if !layers.contains_key(&rid) {
let key = self.layer_key(&rid, &base);
let session = self.ctx.open_publish(key).await?;
layers.insert(
rid.clone(),
Layer {
session,
depack: H264Depacketizer::new(),
needs_keyframe: true,
},
);
}
let layer = layers.get_mut(&rid).unwrap();
let payload = &pkt[header.payload_offset..];
match layer
.depack
.push(payload, header.marker, header.timestamp, header.sequence)
{
Ok(Some(au)) => {
layer.needs_keyframe = false;
let pts = (au.timestamp / 90) as i64;
let frame =
MediaFrame::new_video(pts, pts, au.data, CodecId::H264, au.keyframe);
let _ = layer.session.publish_frame(frame)?;
}
Ok(None) => {}
Err(_) => layer.needs_keyframe = true,
}
if layer.needs_keyframe {
let pli = rtcp::build_pli(0, header.ssrc);
let _ = self.transport.send_rtcp(&pli).await;
}
}
for (_, layer) in layers {
layer.session.finish().await?;
}
Ok(())
}
fn layer_key(&self, rid: &str, base: &str) -> StreamKey {
if rid == base {
self.key.clone()
} else {
self.key.layer(rid)
}
}
pub async fn close(self) -> Result<()> {
Ok(())
}
}
#[derive(Clone)]
pub struct WhepEndpoint {
playback: Arc<dyn PlaybackRegistry>,
gate: Option<crate::auth::EgressGate>,
}
impl WhepEndpoint {
pub fn new(playback: Arc<dyn PlaybackRegistry>) -> Self {
Self {
playback,
gate: None,
}
}
pub fn with_gate(mut self, gate: crate::auth::EgressGate) -> Self {
self.gate = Some(gate);
self
}
pub async fn accept_offer_gated(
&self,
offer_sdp: &str,
key: StreamKey,
token: Option<String>,
peer: Option<std::net::SocketAddr>,
transport: Arc<dyn DtlsSrtpTransport>,
) -> Result<(WhepResource, String)> {
if let Some(gate) = self.gate.as_ref() {
if !gate(key.clone(), token, peer).await {
return Err(crate::StreamError::Unauthorized(
"whep egress denied by gate".into(),
));
}
}
self.accept_offer(offer_sdp, key, transport)
}
pub fn accept_offer(
&self,
offer_sdp: &str,
key: StreamKey,
transport: Arc<dyn DtlsSrtpTransport>,
) -> Result<(WhepResource, String)> {
let offer = SdpOffer::parse(offer_sdp)
.ok_or_else(|| crate::StreamError::protocol("malformed SDP offer"))?;
let answer = transport.answer(offer_sdp, MediaDirection::SendOnly);
let resource = WhepResource {
playback: Arc::clone(&self.playback),
key,
transport,
payload_type: offer.payload_type,
audio_payload_type: offer.audio_payload_type,
warned_unsupported: std::sync::atomic::AtomicBool::new(false),
};
Ok((resource, answer))
}
}
pub struct WhepResource {
playback: Arc<dyn PlaybackRegistry>,
key: StreamKey,
transport: Arc<dyn DtlsSrtpTransport>,
payload_type: u8,
audio_payload_type: Option<u8>,
warned_unsupported: std::sync::atomic::AtomicBool,
}
enum EgressPacketizer {
Nal { p: RtpPacketizer, codec: CodecId },
Vp9(Vp9Packetizer),
#[cfg(feature = "codec-av1")]
Av1(Av1Packetizer),
}
impl EgressPacketizer {
fn for_codec(payload_type: u8, ssrc: u32, mtu: usize, codec: CodecId) -> Self {
match codec {
CodecId::H265 => EgressPacketizer::Nal {
p: RtpPacketizer::new_h265(payload_type, ssrc, mtu),
codec: CodecId::H265,
},
CodecId::VP9 => EgressPacketizer::Vp9(Vp9Packetizer::new(payload_type, ssrc, mtu)),
#[cfg(feature = "codec-av1")]
CodecId::AV1 => EgressPacketizer::Av1(Av1Packetizer::new(payload_type, ssrc, mtu)),
_ => EgressPacketizer::Nal {
p: RtpPacketizer::new(payload_type, ssrc, mtu),
codec: CodecId::H264,
},
}
}
fn packetize_into(&mut self, frame: &MediaFrame, ts_ms: i64, out: &mut Vec<Vec<u8>>) -> bool {
let ts = (ts_ms.max(0) as u64).wrapping_mul(90) as u32; match self {
EgressPacketizer::Nal { p, codec } if frame.codec == *codec => {
p.packetize_into(&frame.data, ts, out);
true
}
EgressPacketizer::Vp9(p) if frame.codec == CodecId::VP9 => {
p.packetize_into(&frame.data, ts, frame.is_keyframe(), out);
true
}
#[cfg(feature = "codec-av1")]
EgressPacketizer::Av1(p) if frame.codec == CodecId::AV1 => {
p.packetize_into(&frame.data, ts, out);
true
}
_ => false,
}
}
}
struct MonoClock {
started: bool,
last_in: i64,
out: i64,
nominal: i64,
}
impl MonoClock {
fn new() -> Self {
Self {
started: false,
last_in: 0,
out: 0,
nominal: 33, }
}
fn map(&mut self, in_ms: i64) -> i64 {
if !self.started {
self.started = true;
self.last_in = in_ms;
self.out = in_ms.max(0);
return self.out;
}
let delta = in_ms - self.last_in;
self.last_in = in_ms;
let step = if delta <= 0 {
1 } else if delta > 1_000 {
self.nominal } else {
self.nominal = delta; delta
};
self.out += step;
self.out
}
}
fn select_layer(
layers: &[(StreamKey, u64)],
estimate: Option<u64>,
current: &StreamKey,
) -> StreamKey {
if layers.is_empty() {
return current.clone();
}
let floor = layers
.iter()
.min_by_key(|(_, bps)| *bps)
.map(|(k, _)| k.clone())
.unwrap();
let current_bps = layers.iter().find(|(k, _)| k == current).map(|(_, b)| *b);
let Some(estimate) = estimate else {
return if current_bps.is_some() {
current.clone()
} else {
floor
};
};
let desired = layers
.iter()
.filter(|(_, bps)| *bps > 0 && *bps <= estimate)
.max_by_key(|(_, bps)| *bps);
let Some((desired_key, desired_bps)) = desired else {
return floor; };
let current_bps = match current_bps {
Some(b) => b,
None => return desired_key.clone(), };
if *desired_bps > current_bps {
if estimate >= desired_bps.saturating_mul(5) / 4 {
return desired_key.clone();
}
} else if *desired_bps < current_bps {
if estimate < current_bps.saturating_mul(19) / 20 {
return desired_key.clone();
}
}
current.clone()
}
impl WhepResource {
pub async fn pump(self) -> Result<()> {
let handle = self.playback.get_stream(&self.key)?;
let ssrc = 0x5745_4850; let mut sub = handle.subscribe_resilient();
let (mut vcfg, _) = handle.cached_configs();
let replay = handle.replay_buffer();
let mut kf_handle = handle.clone();
drop(handle);
let mut current_key = self.key.clone();
let video_codec = vcfg
.as_ref()
.map(|c| c.codec)
.or_else(|| replay.iter().find(|f| f.is_video()).map(|f| f.codec))
.unwrap_or(CodecId::H264);
let mut packetizer =
EgressPacketizer::for_codec(self.payload_type, ssrc, 1200, video_codec);
let mut audio = self
.audio_payload_type
.map(|pt| OpusPacketizer::new(pt, 0x5745_4151));
let mut pkts: Vec<Vec<u8>> = Vec::new();
let mut vclock = MonoClock::new();
let mut aclock = MonoClock::new();
if let Some(cfg) = vcfg.as_ref() {
self.send_frame(
cfg,
&mut packetizer,
&mut audio,
&mut pkts,
&mut vclock,
&mut aclock,
)
.await?;
}
let mut last_keyframe: Option<Arc<MediaFrame>> = None;
for frame in replay {
if frame.is_video() && frame.is_keyframe() {
last_keyframe = Some(frame.clone());
}
self.send_frame(
&frame,
&mut packetizer,
&mut audio,
&mut pkts,
&mut vclock,
&mut aclock,
)
.await?;
}
let mut rtcp_open = true;
let mut abr_tick = tokio::time::interval(std::time::Duration::from_secs(1));
abr_tick.tick().await;
const NO_RTCP_GRACE: std::time::Duration = std::time::Duration::from_secs(2);
let pump_start = std::time::Instant::now();
let mut last_dropped = sub.dropped();
let mut awaiting_keyframe = false;
loop {
let feedback = async {
if rtcp_open {
self.transport.recv_rtcp().await
} else {
std::future::pending().await
}
};
tokio::select! {
frame = sub.recv() => {
let Some(frame) = frame else { break };
let dropped = sub.dropped();
if dropped > last_dropped {
last_dropped = dropped;
awaiting_keyframe = true;
kf_handle.request_keyframe();
}
if frame.is_video() {
if frame.is_keyframe() {
last_keyframe = Some(frame.clone());
awaiting_keyframe = false;
} else if awaiting_keyframe {
continue;
}
} else if awaiting_keyframe {
continue;
}
self.send_frame(&frame, &mut packetizer, &mut audio, &mut pkts, &mut vclock, &mut aclock)
.await?;
}
rtcp = feedback => {
match rtcp {
Some(buf) => {
self.handle_feedback(
&buf,
&kf_handle,
vcfg.as_ref(),
last_keyframe.as_ref(),
&mut packetizer,
&mut audio,
&mut pkts,
&mut vclock,
&mut aclock,
)
.await?;
}
None => {
if pump_start.elapsed() > NO_RTCP_GRACE {
tracing::debug!(stream = %self.key, "WHEP egress: peer transport closed, ending pump");
break;
}
rtcp_open = false;
}
}
}
_ = abr_tick.tick() => {
let layers = self.discover_layers();
let estimate = self.transport.estimated_bitrate();
let target = select_layer(&layers, estimate, ¤t_key);
if target != current_key {
if let Ok(next) = self.playback.get_stream(&target) {
tracing::debug!(
stream = %self.key, from = %current_key, to = %target,
estimate_bps = estimate.unwrap_or(0),
"WHEP egress: adaptive-bitrate layer switch",
);
sub = next.subscribe_resilient();
vcfg = next.cached_configs().0;
kf_handle = next.clone();
current_key = target;
if let Some(cfg) = vcfg.as_ref() {
self.send_frame(cfg, &mut packetizer, &mut audio, &mut pkts, &mut vclock, &mut aclock)
.await?;
}
last_keyframe = None;
awaiting_keyframe = true;
kf_handle.request_keyframe();
}
}
}
}
}
Ok(())
}
fn discover_layers(&self) -> Vec<(StreamKey, u64)> {
let app = &self.key.app;
let base = self.key.stream_id.as_str();
let prefix = format!("{base}~");
let mut out = Vec::new();
let ids = self.playback.list_streams(app).unwrap_or_default();
for id in ids {
let s = id.as_str();
if s == base || s.starts_with(&prefix) {
let key = StreamKey::new(app.as_str(), s);
let bitrate = self
.playback
.get_stream(&key)
.map(|h| h.qos().video_bitrate_bps)
.unwrap_or(0);
out.push((key, bitrate));
}
}
if !out.iter().any(|(k, _)| k == &self.key) {
out.push((self.key.clone(), 0));
}
out
}
#[allow(clippy::too_many_arguments)]
async fn handle_feedback(
&self,
rtcp: &[u8],
kf_handle: &crate::bus::StreamHandle,
vcfg: Option<&Arc<MediaFrame>>,
last_keyframe: Option<&Arc<MediaFrame>>,
packetizer: &mut EgressPacketizer,
audio: &mut Option<OpusPacketizer>,
pkts: &mut Vec<Vec<u8>>,
vclock: &mut MonoClock,
aclock: &mut MonoClock,
) -> Result<()> {
let mut refresh = false;
for fb in rtcp::parse_compound(rtcp) {
match fb {
rtcp::RtcpFeedback::Pli { .. } | rtcp::RtcpFeedback::Fir { .. } => refresh = true,
rtcp::RtcpFeedback::ReceiverReport {
fraction_lost,
cumulative_lost,
jitter,
..
} => {
tracing::trace!(
stream = %self.key,
fraction_lost,
cumulative_lost,
jitter,
"WHEP egress: viewer receiver report",
);
}
rtcp::RtcpFeedback::Nack { lost, .. } => {
tracing::trace!(
stream = %self.key,
lost = lost.len(),
"WHEP egress: viewer NACK (retransmission not yet implemented)",
);
}
rtcp::RtcpFeedback::Remb { bitrate_bps, .. } => {
tracing::trace!(
stream = %self.key,
bitrate_bps,
"WHEP egress: viewer REMB bandwidth estimate",
);
}
}
}
if refresh {
if let Some(cfg) = vcfg {
self.send_frame(cfg, packetizer, audio, pkts, vclock, aclock)
.await?;
}
if let Some(kf) = last_keyframe {
self.send_frame(kf, packetizer, audio, pkts, vclock, aclock)
.await?;
}
kf_handle.request_keyframe();
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
async fn send_frame(
&self,
frame: &MediaFrame,
packetizer: &mut EgressPacketizer,
audio: &mut Option<OpusPacketizer>,
pkts: &mut Vec<Vec<u8>>,
vclock: &mut MonoClock,
aclock: &mut MonoClock,
) -> Result<()> {
if frame.is_audio() {
if let Some(ap) = audio.as_mut() {
if frame.codec == CodecId::Opus {
let ts = (aclock.map(frame.dts).max(0) as u64).wrapping_mul(48) as u32; ap.packetize_into(&frame.data, ts, pkts);
for packet in pkts.iter() {
self.transport.send_rtp(packet).await?;
}
}
}
return Ok(());
}
if !frame.is_video() {
return Ok(());
}
let ts_ms = vclock.map(frame.dts);
if packetizer.packetize_into(frame, ts_ms, pkts) {
for packet in pkts.iter() {
self.transport.send_rtp(packet).await?;
}
} else {
use std::sync::atomic::Ordering;
if !self.warned_unsupported.swap(true, Ordering::Relaxed) {
tracing::warn!(
stream = %self.key,
codec = ?frame.codec,
"WHEP egress: unsupported video codec; frames skipped",
);
}
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::bus::PlaybackRegistry;
use std::sync::Arc;
use tokio::sync::Mutex;
#[test]
fn mono_clock_is_strictly_increasing_through_glitches() {
let mut c = MonoClock::new();
let inputs = [
1000, 1033, 1066, 1050, 1100, 1100, 1133, 50_000, 50_033, 50_066,
];
let mut prev = i64::MIN;
for ms in inputs {
let out = c.map(ms);
assert!(out > prev, "output regressed: {out} after {prev}");
prev = out;
}
}
#[test]
fn mono_clock_passes_through_steady_cadence() {
let mut c = MonoClock::new();
assert_eq!(c.map(0), 0);
assert_eq!(c.map(33), 33);
assert_eq!(c.map(66), 66);
}
struct FakeTransport {
packets: Mutex<std::collections::VecDeque<Vec<u8>>>,
rtcp: Mutex<Vec<Vec<u8>>>,
sent_rtp: Mutex<Vec<Vec<u8>>>,
rtcp_in: Mutex<std::collections::VecDeque<Vec<u8>>>,
keep_open: bool,
}
impl FakeTransport {
fn with_packets(packets: std::collections::VecDeque<Vec<u8>>) -> Self {
Self {
packets: Mutex::new(packets),
rtcp: Mutex::new(Vec::new()),
sent_rtp: Mutex::new(Vec::new()),
rtcp_in: Mutex::new(Default::default()),
keep_open: false,
}
}
fn with_packets_keep_open(packets: std::collections::VecDeque<Vec<u8>>) -> Self {
Self {
keep_open: true,
..Self::with_packets(packets)
}
}
fn with_inbound_rtcp(rtcp: std::collections::VecDeque<Vec<u8>>) -> Self {
Self {
keep_open: true,
rtcp_in: Mutex::new(rtcp),
..Self::with_packets(Default::default())
}
}
}
#[async_trait]
impl DtlsSrtpTransport for FakeTransport {
fn fingerprint(&self) -> String {
"sha-256 AA:BB".into()
}
fn ice_credentials(&self) -> (String, String) {
("ufrag".into(), "pwd".into())
}
async fn recv_rtp(&self) -> Option<Vec<u8>> {
match self.packets.lock().await.pop_front() {
Some(p) => Some(p),
None if self.keep_open => std::future::pending().await,
None => None,
}
}
async fn send_rtp(&self, packet: &[u8]) -> Result<()> {
self.sent_rtp.lock().await.push(packet.to_vec());
Ok(())
}
async fn send_rtcp(&self, packet: &[u8]) -> Result<()> {
self.rtcp.lock().await.push(packet.to_vec());
Ok(())
}
async fn recv_rtcp(&self) -> Option<Vec<u8>> {
match self.rtcp_in.lock().await.pop_front() {
Some(p) => Some(p),
None if self.keep_open => std::future::pending().await,
None => None,
}
}
}
fn rtp_packet(seq: u16, ts: u32, marker: bool, payload: &[u8]) -> Vec<u8> {
rtp_packet_pt(96, seq, ts, marker, payload)
}
fn rtp_packet_pt(pt: u8, seq: u16, ts: u32, marker: bool, payload: &[u8]) -> Vec<u8> {
let mut p = vec![0x80, if marker { 0x80 | pt } else { pt & 0x7F }];
p.extend_from_slice(&seq.to_be_bytes());
p.extend_from_slice(&ts.to_be_bytes());
p.extend_from_slice(&[0, 0, 0, 7]);
p.extend_from_slice(payload);
p
}
fn rtp_with_rid(ext_id: u8, rid: &str, seq: u16, marker: bool, payload: &[u8]) -> Vec<u8> {
let mut p = vec![0x90, if marker { 0x80 | 96 } else { 96 }]; p.extend_from_slice(&seq.to_be_bytes());
p.extend_from_slice(&0u32.to_be_bytes()); p.extend_from_slice(&[0, 0, 0, 7]); p.extend_from_slice(&0xBEDEu16.to_be_bytes()); let mut ext = vec![(ext_id << 4) | (rid.len() as u8 - 1)];
ext.extend_from_slice(rid.as_bytes());
while ext.len() % 4 != 0 {
ext.push(0);
}
p.extend_from_slice(&((ext.len() / 4) as u16).to_be_bytes());
p.extend_from_slice(&ext);
p.extend_from_slice(payload);
p
}
#[tokio::test]
async fn pump_routes_simulcast_layers_to_per_layer_streams() {
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live").gop_cache(4))
.build();
let ctx = IngestContext::new(engine.clone());
let offer = "v=0\r\n\
o=- 0 0 IN IP4 0.0.0.0\r\n\
m=video 9 UDP/TLS/RTP/SAVPF 96\r\n\
a=mid:0\r\n\
a=sendonly\r\n\
a=rtpmap:96 H264/90000\r\n\
a=extmap:4 urn:ietf:params:rtp-hdrext:sdes:rid\r\n\
a=rid:q send\r\n\
a=rid:h send\r\n\
a=simulcast:send q;h\r\n";
let mut q = std::collections::VecDeque::new();
q.push_back(rtp_with_rid(4, "q", 1, true, &[0x65, 0x11]));
q.push_back(rtp_with_rid(4, "h", 2, true, &[0x65, 0x22]));
let transport = Arc::new(FakeTransport::with_packets_keep_open(q));
let endpoint = WhipEndpoint::new(ctx);
let (resource, _answer) = endpoint
.accept_offer(offer, StreamKey::new("live", "cam"), transport)
.unwrap();
let pump = tokio::spawn(resource.pump());
let base = wait_for_stream(&engine, &StreamKey::new("live", "cam")).await;
let high = wait_for_stream(&engine, &StreamKey::new("live", "cam~h")).await;
assert!(base, "base simulcast layer published to the requested key");
assert!(high, "second simulcast layer published to a per-rid key");
pump.abort();
}
async fn wait_for_stream(engine: &Arc<crate::Engine>, key: &StreamKey) -> bool {
for _ in 0..200 {
if engine.get_stream(key).is_ok() {
return true;
}
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
}
false
}
#[tokio::test]
async fn accept_offer_builds_answer_with_transport_credentials() {
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live"))
.build();
let endpoint = WhipEndpoint::new(IngestContext::new(engine));
let transport = Arc::new(FakeTransport::with_packets(Default::default()));
let offer = "v=0\r\no=- 0 0 IN IP4 0.0.0.0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n";
let (_res, answer) = endpoint
.accept_offer(offer, StreamKey::new("live", "web"), transport)
.unwrap();
assert!(answer.contains("a=ice-ufrag:ufrag"));
assert!(answer.contains("a=fingerprint:sha-256 AA:BB"));
assert!(answer.contains("a=setup:passive"));
}
#[tokio::test]
async fn pump_publishes_idr_then_releases_slot() {
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live").gop_cache(4))
.build();
let key = StreamKey::new("live", "web");
let ctx = IngestContext::new(engine.clone());
let mut q = std::collections::VecDeque::new();
q.push_back(rtp_packet(1, 0, true, &[0x65, 0x11])); let transport = Arc::new(FakeTransport::with_packets(q));
let resource = WhipResource {
ctx,
key: key.clone(),
transport,
video_pt: 96,
audio_pt: None,
rid_ext_id: None,
simulcast_rids: Vec::new(),
};
resource.pump().await.unwrap();
assert!(engine.get_stream(&key).is_err());
}
#[tokio::test]
async fn pump_requests_keyframe_on_a_depacketize_gap() {
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live").gop_cache(4))
.build();
let ctx = IngestContext::new(engine);
let mut q = std::collections::VecDeque::new();
q.push_back(rtp_packet(1, 0, false, &[0x7C, 0x05, 0x11])); let transport = Arc::new(FakeTransport::with_packets(q));
let resource = WhipResource {
ctx,
key: StreamKey::new("live", "web2"),
transport: transport.clone(),
video_pt: 96,
audio_pt: None,
rid_ext_id: None,
simulcast_rids: Vec::new(),
};
resource.pump().await.unwrap();
assert!(
!transport.rtcp.lock().await.is_empty(),
"a PLI was sent after the depacketize gap"
);
}
#[tokio::test]
async fn pump_routes_opus_audio_onto_the_bus() {
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live").gop_cache(8))
.build();
let key = StreamKey::new("live", "av");
let ctx = IngestContext::new(engine.clone());
let handle = engine.get_stream(&key);
assert!(handle.is_err(), "stream not live until pump opens publish");
let mut q = std::collections::VecDeque::new();
q.push_back(rtp_packet_pt(111, 7, 4800, true, &[0xAA, 0xBB, 0xCC]));
let transport = Arc::new(FakeTransport::with_packets(q));
let resource = WhipResource {
ctx,
key: key.clone(),
transport,
video_pt: 96,
audio_pt: Some(111),
rid_ext_id: None,
simulcast_rids: Vec::new(),
};
let pump = tokio::spawn(async move { resource.pump().await });
let _ = pump.await.unwrap();
assert!(engine.get_stream(&key).is_err());
}
#[tokio::test]
async fn whep_egress_packetizes_published_frames_as_rtp() {
use crate::FrameFlags;
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live").gop_cache(8))
.build();
let key = StreamKey::new("live", "show");
let ctx = IngestContext::new(engine.clone());
let session = ctx.open_publish(key.clone()).await.unwrap();
let mut cfg = MediaFrame::new_video(
0,
0,
bytes::Bytes::from_static(&[0, 0, 0, 1, 0x67, 0x42]),
CodecId::H264,
false,
);
cfg.flags |= FrameFlags::CONFIG;
session.publish_frame(cfg).unwrap();
session
.publish_frame(MediaFrame::new_video(
10,
10,
bytes::Bytes::from_static(&[0, 0, 0, 1, 0x65, 0x88, 0x99]),
CodecId::H264,
true,
))
.unwrap();
let whep = WhepEndpoint::new(engine.clone());
let transport = Arc::new(FakeTransport::with_packets(Default::default()));
let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n";
let (resource, answer) = whep
.accept_offer(offer, key.clone(), transport.clone())
.unwrap();
assert!(answer.contains("a=sendonly"), "WHEP answer is sendonly");
let pump = tokio::spawn(resource.pump());
for _ in 0..32 {
if !transport.sent_rtp.lock().await.is_empty() {
break;
}
tokio::task::yield_now().await;
}
session.finish().await.unwrap();
let _ = pump.await.unwrap();
let sent = transport.sent_rtp.lock().await;
assert!(!sent.is_empty(), "egress sent RTP packets");
let h = RtpHeader::parse(&sent[0]).unwrap();
assert_eq!(h.payload_type, 96);
}
#[tokio::test]
async fn whep_gate_denies_and_permits() {
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live"))
.build();
let key = StreamKey::new("live", "show");
let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n";
let transport = || Arc::new(FakeTransport::with_packets(Default::default()));
let gate: crate::auth::EgressGate = Arc::new(|_key, token, _peer| {
Box::pin(async move { token.as_deref() == Some("good") })
});
let whep = WhepEndpoint::new(engine.clone()).with_gate(gate);
assert!(whep
.accept_offer_gated(offer, key.clone(), None, None, transport())
.await
.is_err());
assert!(whep
.accept_offer_gated(offer, key.clone(), Some("good".into()), None, transport())
.await
.is_ok());
let open = WhepEndpoint::new(engine.clone());
assert!(open
.accept_offer_gated(offer, key.clone(), None, None, transport())
.await
.is_ok());
}
#[tokio::test]
async fn whep_egress_resends_keyframe_on_viewer_pli() {
use crate::FrameFlags;
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live").gop_cache(8))
.build();
let key = StreamKey::new("live", "fb");
let ctx = IngestContext::new(engine.clone());
let session = ctx.open_publish(key.clone()).await.unwrap();
let mut cfg = MediaFrame::new_video(
0,
0,
bytes::Bytes::from_static(&[0, 0, 0, 1, 0x67, 0x42]),
CodecId::H264,
false,
);
cfg.flags |= FrameFlags::CONFIG;
session.publish_frame(cfg).unwrap();
session
.publish_frame(MediaFrame::new_video(
10,
10,
bytes::Bytes::from_static(&[0, 0, 0, 1, 0x65, 0x88, 0x99]),
CodecId::H264,
true,
))
.unwrap();
let mut script = std::collections::VecDeque::new();
script.push_back(rtcp::build_pli(0, 0));
let transport = Arc::new(FakeTransport::with_inbound_rtcp(script));
let whep = WhepEndpoint::new(engine.clone());
let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n";
let (resource, _) = whep
.accept_offer(offer, key.clone(), transport.clone())
.unwrap();
let pump = tokio::spawn(resource.pump());
let count_keyframes = |pkts: &[Vec<u8>]| {
pkts.iter()
.filter(|p| {
RtpHeader::parse(p)
.map(|h| p[h.payload_offset..].windows(1).any(|w| w[0] == 0x65))
.unwrap_or(false)
})
.count()
};
let mut refreshed = false;
for _ in 0..500 {
if count_keyframes(&transport.sent_rtp.lock().await) >= 2 {
refreshed = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(2)).await;
}
session.finish().await.unwrap();
let _ = pump.await.unwrap();
assert!(refreshed, "viewer PLI re-sent the keyframe");
}
#[test]
fn select_layer_picks_fit_with_hysteresis() {
let k = |s: &str| StreamKey::new("live", s);
let layers = vec![
(k("show"), 300_000u64),
(k("show~h"), 800_000),
(k("show~f"), 2_500_000),
];
assert_eq!(
select_layer(&layers, Some(4_000_000), &k("show")),
k("show~f")
);
assert_eq!(select_layer(&layers, Some(820_000), &k("show")), k("show"));
assert_eq!(
select_layer(&layers, Some(900_000), &k("show~f")),
k("show~h")
);
assert_eq!(select_layer(&layers, None, &k("show~h")), k("show~h"));
assert_eq!(
select_layer(&layers, Some(100_000), &k("show~f")),
k("show")
);
let one = vec![(k("solo"), 0u64)];
assert_eq!(select_layer(&one, Some(5_000_000), &k("solo")), k("solo"));
}
#[tokio::test]
async fn discover_layers_lists_base_and_simulcast_siblings() {
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live").gop_cache(4))
.build();
let ctx = IngestContext::new(engine.clone());
for id in ["show", "show~h", "show~f", "other"] {
let s = ctx.open_publish(StreamKey::new("live", id)).await.unwrap();
s.publish_frame(MediaFrame::new_video(
0,
0,
bytes::Bytes::from_static(&[0, 0, 0, 1, 0x65]),
CodecId::H264,
true,
))
.unwrap();
std::mem::forget(s); }
let whep = WhepEndpoint::new(engine.clone());
let transport = Arc::new(FakeTransport::with_packets(Default::default()));
let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n";
let (resource, _) = whep
.accept_offer(offer, StreamKey::new("live", "show"), transport)
.unwrap();
let mut ids: Vec<String> = resource
.discover_layers()
.into_iter()
.map(|(k, _)| k.stream_id.as_str().to_string())
.collect();
ids.sort();
assert_eq!(ids, vec!["show", "show~f", "show~h"]);
}
#[tokio::test]
async fn request_keyframe_signals_the_publishers_handle() {
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live").gop_cache(8))
.build();
let key = StreamKey::new("live", "kf");
let ctx = IngestContext::new(engine.clone());
let session = ctx.open_publish(key.clone()).await.unwrap();
let pub_handle = session.handle().clone();
let waiter = tokio::spawn(async move {
tokio::time::timeout(
std::time::Duration::from_secs(2),
pub_handle.keyframe_requested(),
)
.await
});
let view_handle = engine.get_stream(&key).unwrap();
tokio::task::yield_now().await;
view_handle.request_keyframe();
assert!(waiter.await.unwrap().is_ok(), "publisher saw the request");
}
#[tokio::test]
async fn whep_egress_packetizes_opus_audio() {
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live").gop_cache(8))
.build();
let key = StreamKey::new("live", "aud");
let ctx = IngestContext::new(engine.clone());
let session = ctx.open_publish(key.clone()).await.unwrap();
session
.publish_frame(MediaFrame::new_video(
0,
0,
bytes::Bytes::from_static(&[0, 0, 0, 1, 0x65, 0x88]),
CodecId::H264,
true,
))
.unwrap();
session
.publish_frame(MediaFrame::new_audio(
20,
bytes::Bytes::from_static(&[0xDE, 0xAD, 0xBE, 0xEF]),
CodecId::Opus,
))
.unwrap();
let whep = WhepEndpoint::new(engine.clone());
let transport = Arc::new(FakeTransport::with_packets(Default::default()));
let offer = "v=0\r\n\
m=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n\
m=audio 9 UDP/TLS/RTP/SAVPF 111\r\na=rtpmap:111 opus/48000/2\r\n";
let (resource, _answer) = whep
.accept_offer(offer, key.clone(), transport.clone())
.unwrap();
let pump = tokio::spawn(resource.pump());
for _ in 0..64 {
if transport
.sent_rtp
.lock()
.await
.iter()
.any(|p| RtpHeader::parse(p).is_some_and(|h| h.payload_type == 111))
{
break;
}
tokio::task::yield_now().await;
}
session.finish().await.unwrap();
let _ = pump.await.unwrap();
let sent = transport.sent_rtp.lock().await;
assert!(
sent.iter()
.any(|p| RtpHeader::parse(p).is_some_and(|h| h.payload_type == 111)),
"egress sent an Opus audio RTP packet on PT 111"
);
}
#[tokio::test]
async fn whep_egress_packetizes_vp9_frames() {
let engine = crate::Engine::builder()
.application(crate::AppSpec::new("live").gop_cache(8))
.build();
let key = StreamKey::new("live", "vp9");
let ctx = IngestContext::new(engine.clone());
let session = ctx.open_publish(key.clone()).await.unwrap();
let frame_data = bytes::Bytes::from_static(&[0xAA, 0xBB, 0xCC, 0xDD, 0xEE]);
session
.publish_frame(MediaFrame::new_video(
0,
0,
frame_data.clone(),
CodecId::VP9,
true,
))
.unwrap();
let whep = WhepEndpoint::new(engine.clone());
let transport = Arc::new(FakeTransport::with_packets(Default::default()));
let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 VP9/90000\r\n";
let (resource, _answer) = whep
.accept_offer(offer, key.clone(), transport.clone())
.unwrap();
let pump = tokio::spawn(resource.pump());
for _ in 0..32 {
if !transport.sent_rtp.lock().await.is_empty() {
break;
}
tokio::task::yield_now().await;
}
session.finish().await.unwrap();
let _ = pump.await.unwrap();
let sent = transport.sent_rtp.lock().await;
assert!(!sent.is_empty(), "VP9 egress sent RTP packets");
let mut depack = crate::protocol::rtp::Vp9Depacketizer::new();
let mut out = None;
for p in sent.iter() {
let h = RtpHeader::parse(p).unwrap();
if let Some(f) = depack
.push(&p[h.payload_offset..], h.marker, h.timestamp)
.unwrap()
{
out = Some(f);
}
}
let out = out.expect("VP9 frame completed");
assert_eq!(&out.data[..], &frame_data[..], "VP9 frame reconstructed");
assert!(out.keyframe);
}
}