1pub mod ice;
41pub mod room;
42pub mod rtcp;
43pub mod sdp;
44
45pub use ice::{parse_trickle, IceCandidate};
46pub use room::{DominantSpeaker, Room};
47pub use sdp::{MediaDirection, SdpAnswerParams, SdpOffer};
48
49use crate::bus::PlaybackRegistry;
50use crate::inbound::{IngestContext, PublishSession};
51#[cfg(feature = "codec-av1")]
52use crate::protocol::rtp::Av1Packetizer;
53use crate::protocol::rtp::{
54 H264Depacketizer, OpusPacketizer, RtpHeader, RtpPacketizer, Vp9Packetizer,
55};
56use crate::{CodecId, MediaFrame, Result, StreamKey};
57use async_trait::async_trait;
58use std::sync::Arc;
59
60#[derive(Debug, Clone, Default, PartialEq)]
66pub struct PeerStats {
67 pub estimated_bitrate_bps: Option<u64>,
70 pub rtt_ms: Option<f32>,
72 pub egress_loss: Option<f32>,
74}
75
76#[async_trait]
81pub trait DtlsSrtpTransport: Send + Sync {
82 fn fingerprint(&self) -> String;
85
86 fn ice_credentials(&self) -> (String, String);
93
94 async fn recv_rtp(&self) -> Option<Vec<u8>> {
98 None
99 }
100
101 async fn send_rtp(&self, _packet: &[u8]) -> Result<()> {
104 Ok(())
105 }
106
107 async fn send_rtcp(&self, packet: &[u8]) -> Result<()>;
109
110 async fn recv_rtcp(&self) -> Option<Vec<u8>> {
119 None
120 }
121
122 fn estimated_bitrate(&self) -> Option<u64> {
130 None
131 }
132
133 fn peer_stats(&self) -> Option<PeerStats> {
136 None
137 }
138
139 async fn add_remote_candidate(&self, _candidate: &str) -> Result<()> {
144 Ok(())
145 }
146
147 async fn recv_data(&self) -> Option<(String, Vec<u8>)> {
151 None
152 }
153
154 async fn send_data(&self, _label: &str, _data: &[u8]) -> Result<()> {
157 Ok(())
158 }
159
160 fn answer(&self, offer_sdp: &str, direction: MediaDirection) -> String {
170 let Some(offer) = SdpOffer::parse(offer_sdp) else {
171 return String::new();
172 };
173 let (ice_ufrag, ice_pwd) = self.ice_credentials();
174 sdp::build_answer_directed(
175 &offer,
176 &SdpAnswerParams {
177 fingerprint: self.fingerprint(),
178 ice_ufrag,
179 ice_pwd,
180 },
181 direction,
182 )
183 }
184}
185
186#[derive(Clone)]
193pub struct WhipEndpoint {
194 ctx: IngestContext,
195}
196
197impl WhipEndpoint {
198 pub fn new(ctx: IngestContext) -> Self {
200 Self { ctx }
201 }
202
203 pub fn accept_offer(
208 &self,
209 offer_sdp: &str,
210 key: StreamKey,
211 transport: std::sync::Arc<dyn DtlsSrtpTransport>,
212 ) -> Result<(WhipResource, String)> {
213 let offer = SdpOffer::parse(offer_sdp)
216 .ok_or_else(|| crate::StreamError::protocol("malformed SDP offer"))?;
217 let answer = transport.answer(offer_sdp, MediaDirection::RecvOnly);
219 let resource = WhipResource {
220 ctx: self.ctx.clone(),
221 key,
222 transport,
223 video_pt: offer.payload_type,
224 audio_pt: offer.audio_payload_type,
225 rid_ext_id: offer.rid_ext_id,
226 simulcast_rids: offer.simulcast_rids,
227 };
228 Ok((resource, answer))
229 }
230}
231
232pub struct WhipResource {
235 ctx: IngestContext,
236 key: StreamKey,
237 transport: std::sync::Arc<dyn DtlsSrtpTransport>,
238 video_pt: u8,
241 audio_pt: Option<u8>,
244 rid_ext_id: Option<u8>,
247 simulcast_rids: Vec<String>,
251}
252
253impl WhipResource {
254 pub async fn pump(self) -> Result<()> {
261 match self.rid_ext_id {
262 Some(ext) if self.simulcast_rids.len() > 1 => self.pump_simulcast(ext).await,
263 _ => self.pump_single().await,
264 }
265 }
266
267 async fn pump_single(self) -> Result<()> {
269 let session: PublishSession = self.ctx.open_publish(self.key.clone()).await?;
270 let handle = session.handle().clone();
271 let mut depack = H264Depacketizer::new();
272 let mut needs_keyframe = true;
273 let mut last_ssrc = 0u32;
276
277 loop {
278 let pkt = tokio::select! {
279 pkt = self.transport.recv_rtp() => match pkt {
280 Some(p) => p,
281 None => break,
282 },
283 _ = handle.keyframe_requested() => {
286 let pli = rtcp::build_pli(0, last_ssrc);
287 let _ = self.transport.send_rtcp(&pli).await;
288 continue;
289 }
290 };
291 let Some(header) = RtpHeader::parse(&pkt) else {
292 continue;
293 };
294 last_ssrc = header.ssrc;
295 let payload = &pkt[header.payload_offset..];
296
297 if self.audio_pt == Some(header.payload_type) {
300 if !payload.is_empty() {
301 let pts = (header.timestamp / 48) as i64;
302 let data = bytes::Bytes::copy_from_slice(payload);
303 let frame = MediaFrame::new_audio(pts, data, CodecId::Opus);
304 let _ = session.publish_frame(frame)?;
305 }
306 continue;
307 }
308
309 let _ = self.video_pt; match depack.push(payload, header.marker, header.timestamp, header.sequence) {
312 Ok(Some(au)) => {
313 needs_keyframe = false;
314 let pts = (au.timestamp / 90) as i64;
315 let frame =
316 MediaFrame::new_video(pts, pts, au.data, CodecId::H264, au.keyframe);
317 let _ = session.publish_frame(frame)?;
318 }
319 Ok(None) => {}
320 Err(_) => {
321 needs_keyframe = true;
323 }
324 }
325 if needs_keyframe {
326 let pli = rtcp::build_pli(0, header.ssrc);
327 let _ = self.transport.send_rtcp(&pli).await;
328 }
329 }
330
331 session.finish().await
332 }
333
334 async fn pump_simulcast(self, rid_ext: u8) -> Result<()> {
340 use std::collections::HashMap;
341 struct Layer {
342 session: PublishSession,
343 depack: H264Depacketizer,
344 needs_keyframe: bool,
345 }
346 let base = self.simulcast_rids[0].clone();
347 let mut layers: HashMap<String, Layer> = HashMap::new();
348
349 while let Some(pkt) = self.transport.recv_rtp().await {
350 let Some(header) = RtpHeader::parse(&pkt) else {
351 continue;
352 };
353 let rid = crate::protocol::rtp::rtp_extension_value(&pkt, rid_ext)
356 .and_then(|b| std::str::from_utf8(b).ok())
357 .map(str::to_owned)
358 .unwrap_or_else(|| base.clone());
359 if !self.simulcast_rids.contains(&rid) {
360 continue; }
362
363 if !layers.contains_key(&rid) {
364 let key = self.layer_key(&rid, &base);
365 let session = self.ctx.open_publish(key).await?;
366 layers.insert(
367 rid.clone(),
368 Layer {
369 session,
370 depack: H264Depacketizer::new(),
371 needs_keyframe: true,
372 },
373 );
374 }
375 let layer = layers.get_mut(&rid).unwrap();
376 let payload = &pkt[header.payload_offset..];
377 match layer
378 .depack
379 .push(payload, header.marker, header.timestamp, header.sequence)
380 {
381 Ok(Some(au)) => {
382 layer.needs_keyframe = false;
383 let pts = (au.timestamp / 90) as i64;
384 let frame =
385 MediaFrame::new_video(pts, pts, au.data, CodecId::H264, au.keyframe);
386 let _ = layer.session.publish_frame(frame)?;
387 }
388 Ok(None) => {}
389 Err(_) => layer.needs_keyframe = true,
390 }
391 if layer.needs_keyframe {
392 let pli = rtcp::build_pli(0, header.ssrc);
393 let _ = self.transport.send_rtcp(&pli).await;
394 }
395 }
396
397 for (_, layer) in layers {
398 layer.session.finish().await?;
399 }
400 Ok(())
401 }
402
403 fn layer_key(&self, rid: &str, base: &str) -> StreamKey {
406 if rid == base {
407 self.key.clone()
408 } else {
409 self.key.layer(rid)
410 }
411 }
412
413 pub async fn close(self) -> Result<()> {
415 Ok(())
416 }
417}
418
419#[derive(Clone)]
428pub struct WhepEndpoint {
429 playback: Arc<dyn PlaybackRegistry>,
430}
431
432impl WhepEndpoint {
433 pub fn new(playback: Arc<dyn PlaybackRegistry>) -> Self {
435 Self { playback }
436 }
437
438 pub fn accept_offer(
441 &self,
442 offer_sdp: &str,
443 key: StreamKey,
444 transport: Arc<dyn DtlsSrtpTransport>,
445 ) -> Result<(WhepResource, String)> {
446 let offer = SdpOffer::parse(offer_sdp)
447 .ok_or_else(|| crate::StreamError::protocol("malformed SDP offer"))?;
448 let answer = transport.answer(offer_sdp, MediaDirection::SendOnly);
450 let resource = WhepResource {
451 playback: Arc::clone(&self.playback),
452 key,
453 transport,
454 payload_type: offer.payload_type,
455 audio_payload_type: offer.audio_payload_type,
456 warned_unsupported: std::sync::atomic::AtomicBool::new(false),
457 };
458 Ok((resource, answer))
459 }
460}
461
462pub struct WhepResource {
465 playback: Arc<dyn PlaybackRegistry>,
466 key: StreamKey,
467 transport: Arc<dyn DtlsSrtpTransport>,
468 payload_type: u8,
469 audio_payload_type: Option<u8>,
472 warned_unsupported: std::sync::atomic::AtomicBool,
475}
476
477enum EgressPacketizer {
482 Nal { p: RtpPacketizer, codec: CodecId },
484 Vp9(Vp9Packetizer),
486 #[cfg(feature = "codec-av1")]
488 Av1(Av1Packetizer),
489}
490
491impl EgressPacketizer {
492 fn for_codec(payload_type: u8, ssrc: u32, mtu: usize, codec: CodecId) -> Self {
496 match codec {
497 CodecId::H265 => EgressPacketizer::Nal {
498 p: RtpPacketizer::new_h265(payload_type, ssrc, mtu),
499 codec: CodecId::H265,
500 },
501 CodecId::VP9 => EgressPacketizer::Vp9(Vp9Packetizer::new(payload_type, ssrc, mtu)),
502 #[cfg(feature = "codec-av1")]
503 CodecId::AV1 => EgressPacketizer::Av1(Av1Packetizer::new(payload_type, ssrc, mtu)),
504 _ => EgressPacketizer::Nal {
505 p: RtpPacketizer::new(payload_type, ssrc, mtu),
506 codec: CodecId::H264,
507 },
508 }
509 }
510
511 fn packetize_into(&mut self, frame: &MediaFrame, out: &mut Vec<Vec<u8>>) -> bool {
515 let ts = (frame.pts.max(0) as u64).wrapping_mul(90) as u32; match self {
519 EgressPacketizer::Nal { p, codec } if frame.codec == *codec => {
520 p.packetize_into(&frame.data, ts, out);
521 true
522 }
523 EgressPacketizer::Vp9(p) if frame.codec == CodecId::VP9 => {
524 p.packetize_into(&frame.data, ts, frame.is_keyframe(), out);
525 true
526 }
527 #[cfg(feature = "codec-av1")]
528 EgressPacketizer::Av1(p) if frame.codec == CodecId::AV1 => {
529 p.packetize_into(&frame.data, ts, out);
530 true
531 }
532 _ => false,
533 }
534 }
535}
536
537fn select_layer(
552 layers: &[(StreamKey, u64)],
553 estimate: Option<u64>,
554 current: &StreamKey,
555) -> StreamKey {
556 if layers.is_empty() {
557 return current.clone();
558 }
559 let floor = layers
562 .iter()
563 .min_by_key(|(_, bps)| *bps)
564 .map(|(k, _)| k.clone())
565 .unwrap();
566 let current_bps = layers.iter().find(|(k, _)| k == current).map(|(_, b)| *b);
567 let Some(estimate) = estimate else {
568 return if current_bps.is_some() {
570 current.clone()
571 } else {
572 floor
573 };
574 };
575
576 let desired = layers
578 .iter()
579 .filter(|(_, bps)| *bps > 0 && *bps <= estimate)
580 .max_by_key(|(_, bps)| *bps);
581 let Some((desired_key, desired_bps)) = desired else {
582 return floor; };
584 let current_bps = match current_bps {
585 Some(b) => b,
586 None => return desired_key.clone(), };
588 if *desired_bps > current_bps {
589 if estimate >= desired_bps.saturating_mul(5) / 4 {
591 return desired_key.clone();
592 }
593 } else if *desired_bps < current_bps {
594 if estimate < current_bps.saturating_mul(19) / 20 {
596 return desired_key.clone();
597 }
598 }
599 current.clone()
600}
601
602impl WhepResource {
603 pub async fn pump(self) -> Result<()> {
612 let handle = self.playback.get_stream(&self.key)?;
613 let ssrc = 0x5745_4850; let mut sub = handle.subscribe_resilient();
617
618 let (mut vcfg, _) = handle.cached_configs();
620 let replay = handle.replay_buffer();
621 let mut kf_handle = handle.clone();
627 drop(handle);
629
630 let mut current_key = self.key.clone();
634
635 let video_codec = vcfg
638 .as_ref()
639 .map(|c| c.codec)
640 .or_else(|| replay.iter().find(|f| f.is_video()).map(|f| f.codec))
641 .unwrap_or(CodecId::H264);
642 let mut packetizer =
643 EgressPacketizer::for_codec(self.payload_type, ssrc, 1200, video_codec);
644 let mut audio = self
648 .audio_payload_type
649 .map(|pt| OpusPacketizer::new(pt, 0x5745_4151)); let mut pkts: Vec<Vec<u8>> = Vec::new();
653
654 if let Some(cfg) = vcfg.as_ref() {
658 self.send_frame(cfg, &mut packetizer, &mut audio, &mut pkts)
659 .await?;
660 }
661 let mut last_keyframe: Option<Arc<MediaFrame>> = None;
662 for frame in replay {
663 if frame.is_video() && frame.is_keyframe() {
664 last_keyframe = Some(frame.clone());
665 }
666 self.send_frame(&frame, &mut packetizer, &mut audio, &mut pkts)
667 .await?;
668 }
669
670 let mut rtcp_open = true;
675 let mut abr_tick = tokio::time::interval(std::time::Duration::from_secs(1));
676 abr_tick.tick().await; loop {
678 let feedback = async {
679 if rtcp_open {
680 self.transport.recv_rtcp().await
681 } else {
682 std::future::pending().await
683 }
684 };
685 tokio::select! {
686 frame = sub.recv() => {
687 let Some(frame) = frame else { break };
688 if frame.is_video() && frame.is_keyframe() {
689 last_keyframe = Some(frame.clone());
690 }
691 self.send_frame(&frame, &mut packetizer, &mut audio, &mut pkts)
692 .await?;
693 }
694 rtcp = feedback => {
695 match rtcp {
696 Some(buf) => {
697 self.handle_feedback(
698 &buf,
699 &kf_handle,
700 vcfg.as_ref(),
701 last_keyframe.as_ref(),
702 &mut packetizer,
703 &mut audio,
704 &mut pkts,
705 )
706 .await?;
707 }
708 None => rtcp_open = false,
709 }
710 }
711 _ = abr_tick.tick() => {
712 let layers = self.discover_layers();
715 let estimate = self.transport.estimated_bitrate();
716 let target = select_layer(&layers, estimate, ¤t_key);
717 if target != current_key {
718 if let Ok(next) = self.playback.get_stream(&target) {
719 tracing::debug!(
720 stream = %self.key, from = %current_key, to = %target,
721 estimate_bps = estimate.unwrap_or(0),
722 "WHEP egress: adaptive-bitrate layer switch",
723 );
724 sub = next.subscribe_resilient();
725 vcfg = next.cached_configs().0;
726 kf_handle = next.clone();
727 current_key = target;
728 if let Some(cfg) = vcfg.as_ref() {
731 self.send_frame(cfg, &mut packetizer, &mut audio, &mut pkts)
732 .await?;
733 }
734 last_keyframe = None;
735 kf_handle.request_keyframe();
736 }
737 }
738 }
739 }
740 }
741 Ok(())
742 }
743
744 fn discover_layers(&self) -> Vec<(StreamKey, u64)> {
751 let app = &self.key.app;
752 let base = self.key.stream_id.as_str();
753 let prefix = format!("{base}~");
754 let mut out = Vec::new();
755 let ids = self.playback.list_streams(app).unwrap_or_default();
756 for id in ids {
757 let s = id.as_str();
758 if s == base || s.starts_with(&prefix) {
759 let key = StreamKey::new(app.as_str(), s);
760 let bitrate = self
761 .playback
762 .get_stream(&key)
763 .map(|h| h.qos().video_bitrate_bps)
764 .unwrap_or(0);
765 out.push((key, bitrate));
766 }
767 }
768 if !out.iter().any(|(k, _)| k == &self.key) {
770 out.push((self.key.clone(), 0));
771 }
772 out
773 }
774
775 #[allow(clippy::too_many_arguments)]
782 async fn handle_feedback(
783 &self,
784 rtcp: &[u8],
785 kf_handle: &crate::bus::StreamHandle,
786 vcfg: Option<&Arc<MediaFrame>>,
787 last_keyframe: Option<&Arc<MediaFrame>>,
788 packetizer: &mut EgressPacketizer,
789 audio: &mut Option<OpusPacketizer>,
790 pkts: &mut Vec<Vec<u8>>,
791 ) -> Result<()> {
792 let mut refresh = false;
793 for fb in rtcp::parse_compound(rtcp) {
794 match fb {
795 rtcp::RtcpFeedback::Pli { .. } | rtcp::RtcpFeedback::Fir { .. } => refresh = true,
796 rtcp::RtcpFeedback::ReceiverReport {
797 fraction_lost,
798 cumulative_lost,
799 jitter,
800 ..
801 } => {
802 tracing::trace!(
803 stream = %self.key,
804 fraction_lost,
805 cumulative_lost,
806 jitter,
807 "WHEP egress: viewer receiver report",
808 );
809 }
810 rtcp::RtcpFeedback::Nack { lost, .. } => {
811 tracing::trace!(
812 stream = %self.key,
813 lost = lost.len(),
814 "WHEP egress: viewer NACK (retransmission not yet implemented)",
815 );
816 }
817 rtcp::RtcpFeedback::Remb { bitrate_bps, .. } => {
818 tracing::trace!(
819 stream = %self.key,
820 bitrate_bps,
821 "WHEP egress: viewer REMB bandwidth estimate",
822 );
823 }
824 }
825 }
826 if refresh {
827 if let Some(cfg) = vcfg {
830 self.send_frame(cfg, packetizer, audio, pkts).await?;
831 }
832 if let Some(kf) = last_keyframe {
833 self.send_frame(kf, packetizer, audio, pkts).await?;
834 }
835 kf_handle.request_keyframe();
838 }
839 Ok(())
840 }
841
842 async fn send_frame(
850 &self,
851 frame: &MediaFrame,
852 packetizer: &mut EgressPacketizer,
853 audio: &mut Option<OpusPacketizer>,
854 pkts: &mut Vec<Vec<u8>>,
855 ) -> Result<()> {
856 if frame.is_audio() {
857 if let Some(ap) = audio.as_mut() {
858 if frame.codec == CodecId::Opus {
859 let ts = (frame.pts.max(0) as u64).wrapping_mul(48) as u32; ap.packetize_into(&frame.data, ts, pkts);
861 for packet in pkts.iter() {
862 self.transport.send_rtp(packet).await?;
863 }
864 }
865 }
866 return Ok(());
867 }
868 if !frame.is_video() {
869 return Ok(());
870 }
871 if packetizer.packetize_into(frame, pkts) {
872 for packet in pkts.iter() {
873 self.transport.send_rtp(packet).await?;
874 }
875 } else {
876 use std::sync::atomic::Ordering;
877 if !self.warned_unsupported.swap(true, Ordering::Relaxed) {
878 tracing::warn!(
879 stream = %self.key,
880 codec = ?frame.codec,
881 "WHEP egress: unsupported video codec; frames skipped",
882 );
883 }
884 }
885 Ok(())
886 }
887}
888
889#[cfg(test)]
890mod tests {
891 use super::*;
892 use crate::bus::PlaybackRegistry;
893 use std::sync::Arc;
894 use tokio::sync::Mutex;
895
896 struct FakeTransport {
898 packets: Mutex<std::collections::VecDeque<Vec<u8>>>,
899 rtcp: Mutex<Vec<Vec<u8>>>,
900 sent_rtp: Mutex<Vec<Vec<u8>>>,
901 rtcp_in: Mutex<std::collections::VecDeque<Vec<u8>>>,
903 keep_open: bool,
906 }
907
908 impl FakeTransport {
909 fn with_packets(packets: std::collections::VecDeque<Vec<u8>>) -> Self {
910 Self {
911 packets: Mutex::new(packets),
912 rtcp: Mutex::new(Vec::new()),
913 sent_rtp: Mutex::new(Vec::new()),
914 rtcp_in: Mutex::new(Default::default()),
915 keep_open: false,
916 }
917 }
918
919 fn with_packets_keep_open(packets: std::collections::VecDeque<Vec<u8>>) -> Self {
920 Self {
921 keep_open: true,
922 ..Self::with_packets(packets)
923 }
924 }
925
926 fn with_inbound_rtcp(rtcp: std::collections::VecDeque<Vec<u8>>) -> Self {
929 Self {
930 keep_open: true,
931 rtcp_in: Mutex::new(rtcp),
932 ..Self::with_packets(Default::default())
933 }
934 }
935 }
936
937 #[async_trait]
938 impl DtlsSrtpTransport for FakeTransport {
939 fn fingerprint(&self) -> String {
940 "sha-256 AA:BB".into()
941 }
942 fn ice_credentials(&self) -> (String, String) {
943 ("ufrag".into(), "pwd".into())
944 }
945 async fn recv_rtp(&self) -> Option<Vec<u8>> {
946 match self.packets.lock().await.pop_front() {
947 Some(p) => Some(p),
948 None if self.keep_open => std::future::pending().await,
949 None => None,
950 }
951 }
952 async fn send_rtp(&self, packet: &[u8]) -> Result<()> {
953 self.sent_rtp.lock().await.push(packet.to_vec());
954 Ok(())
955 }
956 async fn send_rtcp(&self, packet: &[u8]) -> Result<()> {
957 self.rtcp.lock().await.push(packet.to_vec());
958 Ok(())
959 }
960 async fn recv_rtcp(&self) -> Option<Vec<u8>> {
961 match self.rtcp_in.lock().await.pop_front() {
962 Some(p) => Some(p),
963 None if self.keep_open => std::future::pending().await,
964 None => None,
965 }
966 }
967 }
968
969 fn rtp_packet(seq: u16, ts: u32, marker: bool, payload: &[u8]) -> Vec<u8> {
970 rtp_packet_pt(96, seq, ts, marker, payload)
971 }
972
973 fn rtp_packet_pt(pt: u8, seq: u16, ts: u32, marker: bool, payload: &[u8]) -> Vec<u8> {
974 let mut p = vec![0x80, if marker { 0x80 | pt } else { pt & 0x7F }];
975 p.extend_from_slice(&seq.to_be_bytes());
976 p.extend_from_slice(&ts.to_be_bytes());
977 p.extend_from_slice(&[0, 0, 0, 7]);
978 p.extend_from_slice(payload);
979 p
980 }
981
982 fn rtp_with_rid(ext_id: u8, rid: &str, seq: u16, marker: bool, payload: &[u8]) -> Vec<u8> {
984 let mut p = vec![0x90, if marker { 0x80 | 96 } else { 96 }]; p.extend_from_slice(&seq.to_be_bytes());
986 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)];
990 ext.extend_from_slice(rid.as_bytes());
991 while ext.len() % 4 != 0 {
992 ext.push(0);
993 }
994 p.extend_from_slice(&((ext.len() / 4) as u16).to_be_bytes());
995 p.extend_from_slice(&ext);
996 p.extend_from_slice(payload);
997 p
998 }
999
1000 #[tokio::test]
1003 async fn pump_routes_simulcast_layers_to_per_layer_streams() {
1004 let engine = crate::Engine::builder()
1005 .application(crate::AppSpec::new("live").gop_cache(4))
1006 .build();
1007 let ctx = IngestContext::new(engine.clone());
1008 let offer = "v=0\r\n\
1009o=- 0 0 IN IP4 0.0.0.0\r\n\
1010m=video 9 UDP/TLS/RTP/SAVPF 96\r\n\
1011a=mid:0\r\n\
1012a=sendonly\r\n\
1013a=rtpmap:96 H264/90000\r\n\
1014a=extmap:4 urn:ietf:params:rtp-hdrext:sdes:rid\r\n\
1015a=rid:q send\r\n\
1016a=rid:h send\r\n\
1017a=simulcast:send q;h\r\n";
1018
1019 let mut q = std::collections::VecDeque::new();
1021 q.push_back(rtp_with_rid(4, "q", 1, true, &[0x65, 0x11]));
1022 q.push_back(rtp_with_rid(4, "h", 2, true, &[0x65, 0x22]));
1023 let transport = Arc::new(FakeTransport::with_packets_keep_open(q));
1024
1025 let endpoint = WhipEndpoint::new(ctx);
1026 let (resource, _answer) = endpoint
1027 .accept_offer(offer, StreamKey::new("live", "cam"), transport)
1028 .unwrap();
1029 let pump = tokio::spawn(resource.pump());
1030
1031 let base = wait_for_stream(&engine, &StreamKey::new("live", "cam")).await;
1033 let high = wait_for_stream(&engine, &StreamKey::new("live", "cam~h")).await;
1034 assert!(base, "base simulcast layer published to the requested key");
1035 assert!(high, "second simulcast layer published to a per-rid key");
1036
1037 pump.abort();
1038 }
1039
1040 async fn wait_for_stream(engine: &Arc<crate::Engine>, key: &StreamKey) -> bool {
1041 for _ in 0..200 {
1042 if engine.get_stream(key).is_ok() {
1043 return true;
1044 }
1045 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1046 }
1047 false
1048 }
1049
1050 #[tokio::test]
1051 async fn accept_offer_builds_answer_with_transport_credentials() {
1052 let engine = crate::Engine::builder()
1053 .application(crate::AppSpec::new("live"))
1054 .build();
1055 let endpoint = WhipEndpoint::new(IngestContext::new(engine));
1056 let transport = Arc::new(FakeTransport::with_packets(Default::default()));
1057 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";
1058 let (_res, answer) = endpoint
1059 .accept_offer(offer, StreamKey::new("live", "web"), transport)
1060 .unwrap();
1061 assert!(answer.contains("a=ice-ufrag:ufrag"));
1062 assert!(answer.contains("a=fingerprint:sha-256 AA:BB"));
1063 assert!(answer.contains("a=setup:passive"));
1064 }
1065
1066 #[tokio::test]
1067 async fn pump_publishes_idr_then_releases_slot() {
1068 let engine = crate::Engine::builder()
1069 .application(crate::AppSpec::new("live").gop_cache(4))
1070 .build();
1071 let key = StreamKey::new("live", "web");
1072 let ctx = IngestContext::new(engine.clone());
1073
1074 let mut q = std::collections::VecDeque::new();
1075 q.push_back(rtp_packet(1, 0, true, &[0x65, 0x11])); let transport = Arc::new(FakeTransport::with_packets(q));
1077
1078 let resource = WhipResource {
1079 ctx,
1080 key: key.clone(),
1081 transport,
1082 video_pt: 96,
1083 audio_pt: None,
1084 rid_ext_id: None,
1085 simulcast_rids: Vec::new(),
1086 };
1087 resource.pump().await.unwrap();
1088
1089 assert!(engine.get_stream(&key).is_err());
1092 }
1093
1094 #[tokio::test]
1095 async fn pump_requests_keyframe_on_a_depacketize_gap() {
1096 let engine = crate::Engine::builder()
1097 .application(crate::AppSpec::new("live").gop_cache(4))
1098 .build();
1099 let ctx = IngestContext::new(engine);
1100
1101 let mut q = std::collections::VecDeque::new();
1103 q.push_back(rtp_packet(1, 0, false, &[0x7C, 0x05, 0x11])); let transport = Arc::new(FakeTransport::with_packets(q));
1105
1106 let resource = WhipResource {
1107 ctx,
1108 key: StreamKey::new("live", "web2"),
1109 transport: transport.clone(),
1110 video_pt: 96,
1111 audio_pt: None,
1112 rid_ext_id: None,
1113 simulcast_rids: Vec::new(),
1114 };
1115 resource.pump().await.unwrap();
1116 assert!(
1117 !transport.rtcp.lock().await.is_empty(),
1118 "a PLI was sent after the depacketize gap"
1119 );
1120 }
1121
1122 #[tokio::test]
1126 async fn pump_routes_opus_audio_onto_the_bus() {
1127 let engine = crate::Engine::builder()
1128 .application(crate::AppSpec::new("live").gop_cache(8))
1129 .build();
1130 let key = StreamKey::new("live", "av");
1131 let ctx = IngestContext::new(engine.clone());
1132
1133 let handle = engine.get_stream(&key);
1135 assert!(handle.is_err(), "stream not live until pump opens publish");
1136
1137 let mut q = std::collections::VecDeque::new();
1138 q.push_back(rtp_packet_pt(111, 7, 4800, true, &[0xAA, 0xBB, 0xCC]));
1140 let transport = Arc::new(FakeTransport::with_packets(q));
1141
1142 let resource = WhipResource {
1143 ctx,
1144 key: key.clone(),
1145 transport,
1146 video_pt: 96,
1147 audio_pt: Some(111),
1148 rid_ext_id: None,
1149 simulcast_rids: Vec::new(),
1150 };
1151 let pump = tokio::spawn(async move { resource.pump().await });
1153 let _ = pump.await.unwrap();
1155 assert!(engine.get_stream(&key).is_err());
1158 }
1159
1160 #[tokio::test]
1161 async fn whep_egress_packetizes_published_frames_as_rtp() {
1162 use crate::FrameFlags;
1163 let engine = crate::Engine::builder()
1164 .application(crate::AppSpec::new("live").gop_cache(8))
1165 .build();
1166 let key = StreamKey::new("live", "show");
1167
1168 let ctx = IngestContext::new(engine.clone());
1170 let session = ctx.open_publish(key.clone()).await.unwrap();
1171 let mut cfg = MediaFrame::new_video(
1172 0,
1173 0,
1174 bytes::Bytes::from_static(&[0, 0, 0, 1, 0x67, 0x42]),
1175 CodecId::H264,
1176 false,
1177 );
1178 cfg.flags |= FrameFlags::CONFIG;
1179 session.publish_frame(cfg).unwrap();
1180 session
1181 .publish_frame(MediaFrame::new_video(
1182 10,
1183 10,
1184 bytes::Bytes::from_static(&[0, 0, 0, 1, 0x65, 0x88, 0x99]),
1185 CodecId::H264,
1186 true,
1187 ))
1188 .unwrap();
1189
1190 let whep = WhepEndpoint::new(engine.clone());
1192 let transport = Arc::new(FakeTransport::with_packets(Default::default()));
1193 let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n";
1194 let (resource, answer) = whep
1195 .accept_offer(offer, key.clone(), transport.clone())
1196 .unwrap();
1197 assert!(answer.contains("a=sendonly"), "WHEP answer is sendonly");
1198
1199 let pump = tokio::spawn(resource.pump());
1200 for _ in 0..32 {
1202 if !transport.sent_rtp.lock().await.is_empty() {
1203 break;
1204 }
1205 tokio::task::yield_now().await;
1206 }
1207 session.finish().await.unwrap();
1208 let _ = pump.await.unwrap();
1209
1210 let sent = transport.sent_rtp.lock().await;
1211 assert!(!sent.is_empty(), "egress sent RTP packets");
1212 let h = RtpHeader::parse(&sent[0]).unwrap();
1214 assert_eq!(h.payload_type, 96);
1215 }
1216
1217 #[tokio::test]
1221 async fn whep_egress_resends_keyframe_on_viewer_pli() {
1222 use crate::FrameFlags;
1223 let engine = crate::Engine::builder()
1224 .application(crate::AppSpec::new("live").gop_cache(8))
1225 .build();
1226 let key = StreamKey::new("live", "fb");
1227
1228 let ctx = IngestContext::new(engine.clone());
1229 let session = ctx.open_publish(key.clone()).await.unwrap();
1230 let mut cfg = MediaFrame::new_video(
1231 0,
1232 0,
1233 bytes::Bytes::from_static(&[0, 0, 0, 1, 0x67, 0x42]),
1234 CodecId::H264,
1235 false,
1236 );
1237 cfg.flags |= FrameFlags::CONFIG;
1238 session.publish_frame(cfg).unwrap();
1239 session
1240 .publish_frame(MediaFrame::new_video(
1241 10,
1242 10,
1243 bytes::Bytes::from_static(&[0, 0, 0, 1, 0x65, 0x88, 0x99]),
1244 CodecId::H264,
1245 true,
1246 ))
1247 .unwrap();
1248
1249 let mut script = std::collections::VecDeque::new();
1251 script.push_back(rtcp::build_pli(0, 0));
1252 let transport = Arc::new(FakeTransport::with_inbound_rtcp(script));
1253
1254 let whep = WhepEndpoint::new(engine.clone());
1255 let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n";
1256 let (resource, _) = whep
1257 .accept_offer(offer, key.clone(), transport.clone())
1258 .unwrap();
1259 let pump = tokio::spawn(resource.pump());
1260
1261 let count_keyframes = |pkts: &[Vec<u8>]| {
1264 pkts.iter()
1265 .filter(|p| {
1266 RtpHeader::parse(p)
1267 .map(|h| p[h.payload_offset..].windows(1).any(|w| w[0] == 0x65))
1268 .unwrap_or(false)
1269 })
1270 .count()
1271 };
1272 let mut refreshed = false;
1273 for _ in 0..500 {
1274 if count_keyframes(&transport.sent_rtp.lock().await) >= 2 {
1275 refreshed = true;
1276 break;
1277 }
1278 tokio::time::sleep(std::time::Duration::from_millis(2)).await;
1281 }
1282 session.finish().await.unwrap();
1283 let _ = pump.await.unwrap();
1284 assert!(refreshed, "viewer PLI re-sent the keyframe");
1285 }
1286
1287 #[test]
1288 fn select_layer_picks_fit_with_hysteresis() {
1289 let k = |s: &str| StreamKey::new("live", s);
1290 let layers = vec![
1292 (k("show"), 300_000u64),
1293 (k("show~h"), 800_000),
1294 (k("show~f"), 2_500_000),
1295 ];
1296
1297 assert_eq!(
1300 select_layer(&layers, Some(4_000_000), &k("show")),
1301 k("show~f")
1302 );
1303
1304 assert_eq!(select_layer(&layers, Some(820_000), &k("show")), k("show"));
1307
1308 assert_eq!(
1311 select_layer(&layers, Some(900_000), &k("show~f")),
1312 k("show~h")
1313 );
1314
1315 assert_eq!(select_layer(&layers, None, &k("show~h")), k("show~h"));
1317
1318 assert_eq!(
1320 select_layer(&layers, Some(100_000), &k("show~f")),
1321 k("show")
1322 );
1323
1324 let one = vec![(k("solo"), 0u64)];
1326 assert_eq!(select_layer(&one, Some(5_000_000), &k("solo")), k("solo"));
1327 }
1328
1329 #[tokio::test]
1332 async fn discover_layers_lists_base_and_simulcast_siblings() {
1333 let engine = crate::Engine::builder()
1334 .application(crate::AppSpec::new("live").gop_cache(4))
1335 .build();
1336 let ctx = IngestContext::new(engine.clone());
1337 for id in ["show", "show~h", "show~f", "other"] {
1339 let s = ctx.open_publish(StreamKey::new("live", id)).await.unwrap();
1340 s.publish_frame(MediaFrame::new_video(
1341 0,
1342 0,
1343 bytes::Bytes::from_static(&[0, 0, 0, 1, 0x65]),
1344 CodecId::H264,
1345 true,
1346 ))
1347 .unwrap();
1348 std::mem::forget(s); }
1350
1351 let whep = WhepEndpoint::new(engine.clone());
1352 let transport = Arc::new(FakeTransport::with_packets(Default::default()));
1353 let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n";
1354 let (resource, _) = whep
1355 .accept_offer(offer, StreamKey::new("live", "show"), transport)
1356 .unwrap();
1357
1358 let mut ids: Vec<String> = resource
1359 .discover_layers()
1360 .into_iter()
1361 .map(|(k, _)| k.stream_id.as_str().to_string())
1362 .collect();
1363 ids.sort();
1364 assert_eq!(ids, vec!["show", "show~f", "show~h"]);
1365 }
1366
1367 #[tokio::test]
1371 async fn request_keyframe_signals_the_publishers_handle() {
1372 let engine = crate::Engine::builder()
1373 .application(crate::AppSpec::new("live").gop_cache(8))
1374 .build();
1375 let key = StreamKey::new("live", "kf");
1376
1377 let ctx = IngestContext::new(engine.clone());
1378 let session = ctx.open_publish(key.clone()).await.unwrap();
1379 let pub_handle = session.handle().clone();
1381 let waiter = tokio::spawn(async move {
1382 tokio::time::timeout(
1383 std::time::Duration::from_secs(2),
1384 pub_handle.keyframe_requested(),
1385 )
1386 .await
1387 });
1388
1389 let view_handle = engine.get_stream(&key).unwrap();
1391 tokio::task::yield_now().await;
1393 view_handle.request_keyframe();
1394
1395 assert!(waiter.await.unwrap().is_ok(), "publisher saw the request");
1396 }
1397
1398 #[tokio::test]
1401 async fn whep_egress_packetizes_opus_audio() {
1402 let engine = crate::Engine::builder()
1403 .application(crate::AppSpec::new("live").gop_cache(8))
1404 .build();
1405 let key = StreamKey::new("live", "aud");
1406
1407 let ctx = IngestContext::new(engine.clone());
1408 let session = ctx.open_publish(key.clone()).await.unwrap();
1409 session
1411 .publish_frame(MediaFrame::new_video(
1412 0,
1413 0,
1414 bytes::Bytes::from_static(&[0, 0, 0, 1, 0x65, 0x88]),
1415 CodecId::H264,
1416 true,
1417 ))
1418 .unwrap();
1419 session
1420 .publish_frame(MediaFrame::new_audio(
1421 20,
1422 bytes::Bytes::from_static(&[0xDE, 0xAD, 0xBE, 0xEF]),
1423 CodecId::Opus,
1424 ))
1425 .unwrap();
1426
1427 let whep = WhepEndpoint::new(engine.clone());
1428 let transport = Arc::new(FakeTransport::with_packets(Default::default()));
1429 let offer = "v=0\r\n\
1431m=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n\
1432m=audio 9 UDP/TLS/RTP/SAVPF 111\r\na=rtpmap:111 opus/48000/2\r\n";
1433 let (resource, _answer) = whep
1434 .accept_offer(offer, key.clone(), transport.clone())
1435 .unwrap();
1436
1437 let pump = tokio::spawn(resource.pump());
1438 for _ in 0..64 {
1439 if transport
1440 .sent_rtp
1441 .lock()
1442 .await
1443 .iter()
1444 .any(|p| RtpHeader::parse(p).is_some_and(|h| h.payload_type == 111))
1445 {
1446 break;
1447 }
1448 tokio::task::yield_now().await;
1449 }
1450 session.finish().await.unwrap();
1451 let _ = pump.await.unwrap();
1452
1453 let sent = transport.sent_rtp.lock().await;
1454 assert!(
1455 sent.iter()
1456 .any(|p| RtpHeader::parse(p).is_some_and(|h| h.payload_type == 111)),
1457 "egress sent an Opus audio RTP packet on PT 111"
1458 );
1459 }
1460
1461 #[tokio::test]
1462 async fn whep_egress_packetizes_vp9_frames() {
1463 let engine = crate::Engine::builder()
1464 .application(crate::AppSpec::new("live").gop_cache(8))
1465 .build();
1466 let key = StreamKey::new("live", "vp9");
1467
1468 let ctx = IngestContext::new(engine.clone());
1470 let session = ctx.open_publish(key.clone()).await.unwrap();
1471 let frame_data = bytes::Bytes::from_static(&[0xAA, 0xBB, 0xCC, 0xDD, 0xEE]);
1472 session
1473 .publish_frame(MediaFrame::new_video(
1474 0,
1475 0,
1476 frame_data.clone(),
1477 CodecId::VP9,
1478 true,
1479 ))
1480 .unwrap();
1481
1482 let whep = WhepEndpoint::new(engine.clone());
1483 let transport = Arc::new(FakeTransport::with_packets(Default::default()));
1484 let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 VP9/90000\r\n";
1485 let (resource, _answer) = whep
1486 .accept_offer(offer, key.clone(), transport.clone())
1487 .unwrap();
1488
1489 let pump = tokio::spawn(resource.pump());
1490 for _ in 0..32 {
1491 if !transport.sent_rtp.lock().await.is_empty() {
1492 break;
1493 }
1494 tokio::task::yield_now().await;
1495 }
1496 session.finish().await.unwrap();
1497 let _ = pump.await.unwrap();
1498
1499 let sent = transport.sent_rtp.lock().await;
1501 assert!(!sent.is_empty(), "VP9 egress sent RTP packets");
1502 let mut depack = crate::protocol::rtp::Vp9Depacketizer::new();
1503 let mut out = None;
1504 for p in sent.iter() {
1505 let h = RtpHeader::parse(p).unwrap();
1506 if let Some(f) = depack
1507 .push(&p[h.payload_offset..], h.marker, h.timestamp)
1508 .unwrap()
1509 {
1510 out = Some(f);
1511 }
1512 }
1513 let out = out.expect("VP9 frame completed");
1514 assert_eq!(&out.data[..], &frame_data[..], "VP9 frame reconstructed");
1515 assert!(out.keyframe);
1516 }
1517}