1pub mod ice;
41pub mod rtcp;
42pub mod sdp;
43
44pub use ice::{parse_trickle, IceCandidate};
45pub use sdp::{MediaDirection, SdpAnswerParams, SdpOffer};
46
47use crate::bus::PlaybackRegistry;
48use crate::inbound::{IngestContext, PublishSession};
49#[cfg(feature = "codec-av1")]
50use crate::protocol::rtp::Av1Packetizer;
51use crate::protocol::rtp::{
52 H264Depacketizer, OpusPacketizer, RtpHeader, RtpPacketizer, Vp9Packetizer,
53};
54use crate::{CodecId, MediaFrame, Result, StreamKey};
55use async_trait::async_trait;
56use std::sync::Arc;
57
58#[async_trait]
63pub trait DtlsSrtpTransport: Send + Sync {
64 fn fingerprint(&self) -> String;
67
68 fn ice_credentials(&self) -> (String, String);
75
76 async fn recv_rtp(&self) -> Option<Vec<u8>> {
80 None
81 }
82
83 async fn send_rtp(&self, _packet: &[u8]) -> Result<()> {
86 Ok(())
87 }
88
89 async fn send_rtcp(&self, packet: &[u8]) -> Result<()>;
91
92 async fn add_remote_candidate(&self, _candidate: &str) -> Result<()> {
97 Ok(())
98 }
99
100 async fn recv_data(&self) -> Option<(String, Vec<u8>)> {
104 None
105 }
106
107 async fn send_data(&self, _label: &str, _data: &[u8]) -> Result<()> {
110 Ok(())
111 }
112
113 fn answer(&self, offer_sdp: &str, direction: MediaDirection) -> String {
123 let Some(offer) = SdpOffer::parse(offer_sdp) else {
124 return String::new();
125 };
126 let (ice_ufrag, ice_pwd) = self.ice_credentials();
127 sdp::build_answer_directed(
128 &offer,
129 &SdpAnswerParams {
130 fingerprint: self.fingerprint(),
131 ice_ufrag,
132 ice_pwd,
133 },
134 direction,
135 )
136 }
137}
138
139#[derive(Clone)]
146pub struct WhipEndpoint {
147 ctx: IngestContext,
148}
149
150impl WhipEndpoint {
151 pub fn new(ctx: IngestContext) -> Self {
153 Self { ctx }
154 }
155
156 pub fn accept_offer(
161 &self,
162 offer_sdp: &str,
163 key: StreamKey,
164 transport: std::sync::Arc<dyn DtlsSrtpTransport>,
165 ) -> Result<(WhipResource, String)> {
166 let offer = SdpOffer::parse(offer_sdp)
169 .ok_or_else(|| crate::StreamError::protocol("malformed SDP offer"))?;
170 let answer = transport.answer(offer_sdp, MediaDirection::RecvOnly);
172 let resource = WhipResource {
173 ctx: self.ctx.clone(),
174 key,
175 transport,
176 video_pt: offer.payload_type,
177 audio_pt: offer.audio_payload_type,
178 rid_ext_id: offer.rid_ext_id,
179 simulcast_rids: offer.simulcast_rids,
180 };
181 Ok((resource, answer))
182 }
183}
184
185pub struct WhipResource {
188 ctx: IngestContext,
189 key: StreamKey,
190 transport: std::sync::Arc<dyn DtlsSrtpTransport>,
191 video_pt: u8,
194 audio_pt: Option<u8>,
197 rid_ext_id: Option<u8>,
200 simulcast_rids: Vec<String>,
204}
205
206impl WhipResource {
207 pub async fn pump(self) -> Result<()> {
214 match self.rid_ext_id {
215 Some(ext) if self.simulcast_rids.len() > 1 => self.pump_simulcast(ext).await,
216 _ => self.pump_single().await,
217 }
218 }
219
220 async fn pump_single(self) -> Result<()> {
222 let session: PublishSession = self.ctx.open_publish(self.key.clone()).await?;
223 let mut depack = H264Depacketizer::new();
224 let mut needs_keyframe = true;
225
226 while let Some(pkt) = self.transport.recv_rtp().await {
227 let Some(header) = RtpHeader::parse(&pkt) else {
228 continue;
229 };
230 let payload = &pkt[header.payload_offset..];
231
232 if self.audio_pt == Some(header.payload_type) {
235 if !payload.is_empty() {
236 let pts = (header.timestamp / 48) as i64;
237 let data = bytes::Bytes::copy_from_slice(payload);
238 let frame = MediaFrame::new_audio(pts, data, CodecId::Opus);
239 let _ = session.publish_frame(frame)?;
240 }
241 continue;
242 }
243
244 let _ = self.video_pt; match depack.push(payload, header.marker, header.timestamp, header.sequence) {
247 Ok(Some(au)) => {
248 needs_keyframe = false;
249 let pts = (au.timestamp / 90) as i64;
250 let frame =
251 MediaFrame::new_video(pts, pts, au.data, CodecId::H264, au.keyframe);
252 let _ = session.publish_frame(frame)?;
253 }
254 Ok(None) => {}
255 Err(_) => {
256 needs_keyframe = true;
258 }
259 }
260 if needs_keyframe {
261 let pli = rtcp::build_pli(0, header.ssrc);
262 let _ = self.transport.send_rtcp(&pli).await;
263 }
264 }
265
266 session.finish().await
267 }
268
269 async fn pump_simulcast(self, rid_ext: u8) -> Result<()> {
275 use std::collections::HashMap;
276 struct Layer {
277 session: PublishSession,
278 depack: H264Depacketizer,
279 needs_keyframe: bool,
280 }
281 let base = self.simulcast_rids[0].clone();
282 let mut layers: HashMap<String, Layer> = HashMap::new();
283
284 while let Some(pkt) = self.transport.recv_rtp().await {
285 let Some(header) = RtpHeader::parse(&pkt) else {
286 continue;
287 };
288 let rid = crate::protocol::rtp::rtp_extension_value(&pkt, rid_ext)
291 .and_then(|b| std::str::from_utf8(b).ok())
292 .map(str::to_owned)
293 .unwrap_or_else(|| base.clone());
294 if !self.simulcast_rids.contains(&rid) {
295 continue; }
297
298 if !layers.contains_key(&rid) {
299 let key = self.layer_key(&rid, &base);
300 let session = self.ctx.open_publish(key).await?;
301 layers.insert(
302 rid.clone(),
303 Layer {
304 session,
305 depack: H264Depacketizer::new(),
306 needs_keyframe: true,
307 },
308 );
309 }
310 let layer = layers.get_mut(&rid).unwrap();
311 let payload = &pkt[header.payload_offset..];
312 match layer
313 .depack
314 .push(payload, header.marker, header.timestamp, header.sequence)
315 {
316 Ok(Some(au)) => {
317 layer.needs_keyframe = false;
318 let pts = (au.timestamp / 90) as i64;
319 let frame =
320 MediaFrame::new_video(pts, pts, au.data, CodecId::H264, au.keyframe);
321 let _ = layer.session.publish_frame(frame)?;
322 }
323 Ok(None) => {}
324 Err(_) => layer.needs_keyframe = true,
325 }
326 if layer.needs_keyframe {
327 let pli = rtcp::build_pli(0, header.ssrc);
328 let _ = self.transport.send_rtcp(&pli).await;
329 }
330 }
331
332 for (_, layer) in layers {
333 layer.session.finish().await?;
334 }
335 Ok(())
336 }
337
338 fn layer_key(&self, rid: &str, base: &str) -> StreamKey {
341 if rid == base {
342 self.key.clone()
343 } else {
344 StreamKey::new(
345 self.key.app.as_str(),
346 format!("{}~{}", self.key.stream_id.as_str(), rid),
347 )
348 }
349 }
350
351 pub async fn close(self) -> Result<()> {
353 Ok(())
354 }
355}
356
357#[derive(Clone)]
366pub struct WhepEndpoint {
367 playback: Arc<dyn PlaybackRegistry>,
368}
369
370impl WhepEndpoint {
371 pub fn new(playback: Arc<dyn PlaybackRegistry>) -> Self {
373 Self { playback }
374 }
375
376 pub fn accept_offer(
379 &self,
380 offer_sdp: &str,
381 key: StreamKey,
382 transport: Arc<dyn DtlsSrtpTransport>,
383 ) -> Result<(WhepResource, String)> {
384 let offer = SdpOffer::parse(offer_sdp)
385 .ok_or_else(|| crate::StreamError::protocol("malformed SDP offer"))?;
386 let answer = transport.answer(offer_sdp, MediaDirection::SendOnly);
388 let resource = WhepResource {
389 playback: Arc::clone(&self.playback),
390 key,
391 transport,
392 payload_type: offer.payload_type,
393 audio_payload_type: offer.audio_payload_type,
394 warned_unsupported: std::sync::atomic::AtomicBool::new(false),
395 };
396 Ok((resource, answer))
397 }
398}
399
400pub struct WhepResource {
403 playback: Arc<dyn PlaybackRegistry>,
404 key: StreamKey,
405 transport: Arc<dyn DtlsSrtpTransport>,
406 payload_type: u8,
407 audio_payload_type: Option<u8>,
410 warned_unsupported: std::sync::atomic::AtomicBool,
413}
414
415enum EgressPacketizer {
420 Nal { p: RtpPacketizer, codec: CodecId },
422 Vp9(Vp9Packetizer),
424 #[cfg(feature = "codec-av1")]
426 Av1(Av1Packetizer),
427}
428
429impl EgressPacketizer {
430 fn for_codec(payload_type: u8, ssrc: u32, mtu: usize, codec: CodecId) -> Self {
434 match codec {
435 CodecId::H265 => EgressPacketizer::Nal {
436 p: RtpPacketizer::new_h265(payload_type, ssrc, mtu),
437 codec: CodecId::H265,
438 },
439 CodecId::VP9 => EgressPacketizer::Vp9(Vp9Packetizer::new(payload_type, ssrc, mtu)),
440 #[cfg(feature = "codec-av1")]
441 CodecId::AV1 => EgressPacketizer::Av1(Av1Packetizer::new(payload_type, ssrc, mtu)),
442 _ => EgressPacketizer::Nal {
443 p: RtpPacketizer::new(payload_type, ssrc, mtu),
444 codec: CodecId::H264,
445 },
446 }
447 }
448
449 fn packetize_into(&mut self, frame: &MediaFrame, out: &mut Vec<Vec<u8>>) -> bool {
453 let ts = (frame.pts.max(0) as u64).wrapping_mul(90) as u32; match self {
457 EgressPacketizer::Nal { p, codec } if frame.codec == *codec => {
458 p.packetize_into(&frame.data, ts, out);
459 true
460 }
461 EgressPacketizer::Vp9(p) if frame.codec == CodecId::VP9 => {
462 p.packetize_into(&frame.data, ts, frame.is_keyframe(), out);
463 true
464 }
465 #[cfg(feature = "codec-av1")]
466 EgressPacketizer::Av1(p) if frame.codec == CodecId::AV1 => {
467 p.packetize_into(&frame.data, ts, out);
468 true
469 }
470 _ => false,
471 }
472 }
473}
474
475impl WhepResource {
476 pub async fn pump(self) -> Result<()> {
485 let handle = self.playback.get_stream(&self.key)?;
486 let ssrc = 0x5745_4850; let mut sub = handle.subscribe_resilient();
490
491 let (vcfg, _) = handle.cached_configs();
493 let replay = handle.replay_buffer();
494 drop(handle);
498
499 let video_codec = vcfg
502 .as_ref()
503 .map(|c| c.codec)
504 .or_else(|| replay.iter().find(|f| f.is_video()).map(|f| f.codec))
505 .unwrap_or(CodecId::H264);
506 let mut packetizer =
507 EgressPacketizer::for_codec(self.payload_type, ssrc, 1200, video_codec);
508 let mut audio = self
512 .audio_payload_type
513 .map(|pt| OpusPacketizer::new(pt, 0x5745_4151)); let mut pkts: Vec<Vec<u8>> = Vec::new();
517
518 if let Some(cfg) = vcfg {
519 self.send_frame(&cfg, &mut packetizer, &mut audio, &mut pkts)
520 .await?;
521 }
522 for frame in replay {
523 self.send_frame(&frame, &mut packetizer, &mut audio, &mut pkts)
524 .await?;
525 }
526
527 while let Some(frame) = sub.recv().await {
528 self.send_frame(&frame, &mut packetizer, &mut audio, &mut pkts)
529 .await?;
530 }
531 Ok(())
532 }
533
534 async fn send_frame(
542 &self,
543 frame: &MediaFrame,
544 packetizer: &mut EgressPacketizer,
545 audio: &mut Option<OpusPacketizer>,
546 pkts: &mut Vec<Vec<u8>>,
547 ) -> Result<()> {
548 if frame.is_audio() {
549 if let Some(ap) = audio.as_mut() {
550 if frame.codec == CodecId::Opus {
551 let ts = (frame.pts.max(0) as u64).wrapping_mul(48) as u32; ap.packetize_into(&frame.data, ts, pkts);
553 for packet in pkts.iter() {
554 self.transport.send_rtp(packet).await?;
555 }
556 }
557 }
558 return Ok(());
559 }
560 if !frame.is_video() {
561 return Ok(());
562 }
563 if packetizer.packetize_into(frame, pkts) {
564 for packet in pkts.iter() {
565 self.transport.send_rtp(packet).await?;
566 }
567 } else {
568 use std::sync::atomic::Ordering;
569 if !self.warned_unsupported.swap(true, Ordering::Relaxed) {
570 tracing::warn!(
571 stream = %self.key,
572 codec = ?frame.codec,
573 "WHEP egress: unsupported video codec; frames skipped",
574 );
575 }
576 }
577 Ok(())
578 }
579}
580
581#[cfg(test)]
582mod tests {
583 use super::*;
584 use crate::bus::PlaybackRegistry;
585 use std::sync::Arc;
586 use tokio::sync::Mutex;
587
588 struct FakeTransport {
590 packets: Mutex<std::collections::VecDeque<Vec<u8>>>,
591 rtcp: Mutex<Vec<Vec<u8>>>,
592 sent_rtp: Mutex<Vec<Vec<u8>>>,
593 keep_open: bool,
596 }
597
598 impl FakeTransport {
599 fn with_packets(packets: std::collections::VecDeque<Vec<u8>>) -> Self {
600 Self {
601 packets: Mutex::new(packets),
602 rtcp: Mutex::new(Vec::new()),
603 sent_rtp: Mutex::new(Vec::new()),
604 keep_open: false,
605 }
606 }
607
608 fn with_packets_keep_open(packets: std::collections::VecDeque<Vec<u8>>) -> Self {
609 Self {
610 keep_open: true,
611 ..Self::with_packets(packets)
612 }
613 }
614 }
615
616 #[async_trait]
617 impl DtlsSrtpTransport for FakeTransport {
618 fn fingerprint(&self) -> String {
619 "sha-256 AA:BB".into()
620 }
621 fn ice_credentials(&self) -> (String, String) {
622 ("ufrag".into(), "pwd".into())
623 }
624 async fn recv_rtp(&self) -> Option<Vec<u8>> {
625 match self.packets.lock().await.pop_front() {
626 Some(p) => Some(p),
627 None if self.keep_open => std::future::pending().await,
628 None => None,
629 }
630 }
631 async fn send_rtp(&self, packet: &[u8]) -> Result<()> {
632 self.sent_rtp.lock().await.push(packet.to_vec());
633 Ok(())
634 }
635 async fn send_rtcp(&self, packet: &[u8]) -> Result<()> {
636 self.rtcp.lock().await.push(packet.to_vec());
637 Ok(())
638 }
639 }
640
641 fn rtp_packet(seq: u16, ts: u32, marker: bool, payload: &[u8]) -> Vec<u8> {
642 rtp_packet_pt(96, seq, ts, marker, payload)
643 }
644
645 fn rtp_packet_pt(pt: u8, seq: u16, ts: u32, marker: bool, payload: &[u8]) -> Vec<u8> {
646 let mut p = vec![0x80, if marker { 0x80 | pt } else { pt & 0x7F }];
647 p.extend_from_slice(&seq.to_be_bytes());
648 p.extend_from_slice(&ts.to_be_bytes());
649 p.extend_from_slice(&[0, 0, 0, 7]);
650 p.extend_from_slice(payload);
651 p
652 }
653
654 fn rtp_with_rid(ext_id: u8, rid: &str, seq: u16, marker: bool, payload: &[u8]) -> Vec<u8> {
656 let mut p = vec![0x90, if marker { 0x80 | 96 } else { 96 }]; p.extend_from_slice(&seq.to_be_bytes());
658 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)];
662 ext.extend_from_slice(rid.as_bytes());
663 while ext.len() % 4 != 0 {
664 ext.push(0);
665 }
666 p.extend_from_slice(&((ext.len() / 4) as u16).to_be_bytes());
667 p.extend_from_slice(&ext);
668 p.extend_from_slice(payload);
669 p
670 }
671
672 #[tokio::test]
675 async fn pump_routes_simulcast_layers_to_per_layer_streams() {
676 let engine = crate::Engine::builder()
677 .application(crate::AppSpec::new("live").gop_cache(4))
678 .build();
679 let ctx = IngestContext::new(engine.clone());
680 let offer = "v=0\r\n\
681o=- 0 0 IN IP4 0.0.0.0\r\n\
682m=video 9 UDP/TLS/RTP/SAVPF 96\r\n\
683a=mid:0\r\n\
684a=sendonly\r\n\
685a=rtpmap:96 H264/90000\r\n\
686a=extmap:4 urn:ietf:params:rtp-hdrext:sdes:rid\r\n\
687a=rid:q send\r\n\
688a=rid:h send\r\n\
689a=simulcast:send q;h\r\n";
690
691 let mut q = std::collections::VecDeque::new();
693 q.push_back(rtp_with_rid(4, "q", 1, true, &[0x65, 0x11]));
694 q.push_back(rtp_with_rid(4, "h", 2, true, &[0x65, 0x22]));
695 let transport = Arc::new(FakeTransport::with_packets_keep_open(q));
696
697 let endpoint = WhipEndpoint::new(ctx);
698 let (resource, _answer) = endpoint
699 .accept_offer(offer, StreamKey::new("live", "cam"), transport)
700 .unwrap();
701 let pump = tokio::spawn(resource.pump());
702
703 let base = wait_for_stream(&engine, &StreamKey::new("live", "cam")).await;
705 let high = wait_for_stream(&engine, &StreamKey::new("live", "cam~h")).await;
706 assert!(base, "base simulcast layer published to the requested key");
707 assert!(high, "second simulcast layer published to a per-rid key");
708
709 pump.abort();
710 }
711
712 async fn wait_for_stream(engine: &Arc<crate::Engine>, key: &StreamKey) -> bool {
713 for _ in 0..200 {
714 if engine.get_stream(key).is_ok() {
715 return true;
716 }
717 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
718 }
719 false
720 }
721
722 #[tokio::test]
723 async fn accept_offer_builds_answer_with_transport_credentials() {
724 let engine = crate::Engine::builder()
725 .application(crate::AppSpec::new("live"))
726 .build();
727 let endpoint = WhipEndpoint::new(IngestContext::new(engine));
728 let transport = Arc::new(FakeTransport::with_packets(Default::default()));
729 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";
730 let (_res, answer) = endpoint
731 .accept_offer(offer, StreamKey::new("live", "web"), transport)
732 .unwrap();
733 assert!(answer.contains("a=ice-ufrag:ufrag"));
734 assert!(answer.contains("a=fingerprint:sha-256 AA:BB"));
735 assert!(answer.contains("a=setup:passive"));
736 }
737
738 #[tokio::test]
739 async fn pump_publishes_idr_then_releases_slot() {
740 let engine = crate::Engine::builder()
741 .application(crate::AppSpec::new("live").gop_cache(4))
742 .build();
743 let key = StreamKey::new("live", "web");
744 let ctx = IngestContext::new(engine.clone());
745
746 let mut q = std::collections::VecDeque::new();
747 q.push_back(rtp_packet(1, 0, true, &[0x65, 0x11])); let transport = Arc::new(FakeTransport::with_packets(q));
749
750 let resource = WhipResource {
751 ctx,
752 key: key.clone(),
753 transport,
754 video_pt: 96,
755 audio_pt: None,
756 rid_ext_id: None,
757 simulcast_rids: Vec::new(),
758 };
759 resource.pump().await.unwrap();
760
761 assert!(engine.get_stream(&key).is_err());
764 }
765
766 #[tokio::test]
767 async fn pump_requests_keyframe_on_a_depacketize_gap() {
768 let engine = crate::Engine::builder()
769 .application(crate::AppSpec::new("live").gop_cache(4))
770 .build();
771 let ctx = IngestContext::new(engine);
772
773 let mut q = std::collections::VecDeque::new();
775 q.push_back(rtp_packet(1, 0, false, &[0x7C, 0x05, 0x11])); let transport = Arc::new(FakeTransport::with_packets(q));
777
778 let resource = WhipResource {
779 ctx,
780 key: StreamKey::new("live", "web2"),
781 transport: transport.clone(),
782 video_pt: 96,
783 audio_pt: None,
784 rid_ext_id: None,
785 simulcast_rids: Vec::new(),
786 };
787 resource.pump().await.unwrap();
788 assert!(
789 !transport.rtcp.lock().await.is_empty(),
790 "a PLI was sent after the depacketize gap"
791 );
792 }
793
794 #[tokio::test]
798 async fn pump_routes_opus_audio_onto_the_bus() {
799 let engine = crate::Engine::builder()
800 .application(crate::AppSpec::new("live").gop_cache(8))
801 .build();
802 let key = StreamKey::new("live", "av");
803 let ctx = IngestContext::new(engine.clone());
804
805 let handle = engine.get_stream(&key);
807 assert!(handle.is_err(), "stream not live until pump opens publish");
808
809 let mut q = std::collections::VecDeque::new();
810 q.push_back(rtp_packet_pt(111, 7, 4800, true, &[0xAA, 0xBB, 0xCC]));
812 let transport = Arc::new(FakeTransport::with_packets(q));
813
814 let resource = WhipResource {
815 ctx,
816 key: key.clone(),
817 transport,
818 video_pt: 96,
819 audio_pt: Some(111),
820 rid_ext_id: None,
821 simulcast_rids: Vec::new(),
822 };
823 let pump = tokio::spawn(async move { resource.pump().await });
825 let _ = pump.await.unwrap();
827 assert!(engine.get_stream(&key).is_err());
830 }
831
832 #[tokio::test]
833 async fn whep_egress_packetizes_published_frames_as_rtp() {
834 use crate::FrameFlags;
835 let engine = crate::Engine::builder()
836 .application(crate::AppSpec::new("live").gop_cache(8))
837 .build();
838 let key = StreamKey::new("live", "show");
839
840 let ctx = IngestContext::new(engine.clone());
842 let session = ctx.open_publish(key.clone()).await.unwrap();
843 let mut cfg = MediaFrame::new_video(
844 0,
845 0,
846 bytes::Bytes::from_static(&[0, 0, 0, 1, 0x67, 0x42]),
847 CodecId::H264,
848 false,
849 );
850 cfg.flags |= FrameFlags::CONFIG;
851 session.publish_frame(cfg).unwrap();
852 session
853 .publish_frame(MediaFrame::new_video(
854 10,
855 10,
856 bytes::Bytes::from_static(&[0, 0, 0, 1, 0x65, 0x88, 0x99]),
857 CodecId::H264,
858 true,
859 ))
860 .unwrap();
861
862 let whep = WhepEndpoint::new(engine.clone());
864 let transport = Arc::new(FakeTransport::with_packets(Default::default()));
865 let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n";
866 let (resource, answer) = whep
867 .accept_offer(offer, key.clone(), transport.clone())
868 .unwrap();
869 assert!(answer.contains("a=sendonly"), "WHEP answer is sendonly");
870
871 let pump = tokio::spawn(resource.pump());
872 for _ in 0..32 {
874 if !transport.sent_rtp.lock().await.is_empty() {
875 break;
876 }
877 tokio::task::yield_now().await;
878 }
879 session.finish().await.unwrap();
880 let _ = pump.await.unwrap();
881
882 let sent = transport.sent_rtp.lock().await;
883 assert!(!sent.is_empty(), "egress sent RTP packets");
884 let h = RtpHeader::parse(&sent[0]).unwrap();
886 assert_eq!(h.payload_type, 96);
887 }
888
889 #[tokio::test]
892 async fn whep_egress_packetizes_opus_audio() {
893 let engine = crate::Engine::builder()
894 .application(crate::AppSpec::new("live").gop_cache(8))
895 .build();
896 let key = StreamKey::new("live", "aud");
897
898 let ctx = IngestContext::new(engine.clone());
899 let session = ctx.open_publish(key.clone()).await.unwrap();
900 session
902 .publish_frame(MediaFrame::new_video(
903 0,
904 0,
905 bytes::Bytes::from_static(&[0, 0, 0, 1, 0x65, 0x88]),
906 CodecId::H264,
907 true,
908 ))
909 .unwrap();
910 session
911 .publish_frame(MediaFrame::new_audio(
912 20,
913 bytes::Bytes::from_static(&[0xDE, 0xAD, 0xBE, 0xEF]),
914 CodecId::Opus,
915 ))
916 .unwrap();
917
918 let whep = WhepEndpoint::new(engine.clone());
919 let transport = Arc::new(FakeTransport::with_packets(Default::default()));
920 let offer = "v=0\r\n\
922m=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 H264/90000\r\n\
923m=audio 9 UDP/TLS/RTP/SAVPF 111\r\na=rtpmap:111 opus/48000/2\r\n";
924 let (resource, _answer) = whep
925 .accept_offer(offer, key.clone(), transport.clone())
926 .unwrap();
927
928 let pump = tokio::spawn(resource.pump());
929 for _ in 0..64 {
930 if transport
931 .sent_rtp
932 .lock()
933 .await
934 .iter()
935 .any(|p| RtpHeader::parse(p).is_some_and(|h| h.payload_type == 111))
936 {
937 break;
938 }
939 tokio::task::yield_now().await;
940 }
941 session.finish().await.unwrap();
942 let _ = pump.await.unwrap();
943
944 let sent = transport.sent_rtp.lock().await;
945 assert!(
946 sent.iter()
947 .any(|p| RtpHeader::parse(p).is_some_and(|h| h.payload_type == 111)),
948 "egress sent an Opus audio RTP packet on PT 111"
949 );
950 }
951
952 #[tokio::test]
953 async fn whep_egress_packetizes_vp9_frames() {
954 let engine = crate::Engine::builder()
955 .application(crate::AppSpec::new("live").gop_cache(8))
956 .build();
957 let key = StreamKey::new("live", "vp9");
958
959 let ctx = IngestContext::new(engine.clone());
961 let session = ctx.open_publish(key.clone()).await.unwrap();
962 let frame_data = bytes::Bytes::from_static(&[0xAA, 0xBB, 0xCC, 0xDD, 0xEE]);
963 session
964 .publish_frame(MediaFrame::new_video(
965 0,
966 0,
967 frame_data.clone(),
968 CodecId::VP9,
969 true,
970 ))
971 .unwrap();
972
973 let whep = WhepEndpoint::new(engine.clone());
974 let transport = Arc::new(FakeTransport::with_packets(Default::default()));
975 let offer = "v=0\r\nm=video 9 UDP/TLS/RTP/SAVPF 96\r\na=rtpmap:96 VP9/90000\r\n";
976 let (resource, _answer) = whep
977 .accept_offer(offer, key.clone(), transport.clone())
978 .unwrap();
979
980 let pump = tokio::spawn(resource.pump());
981 for _ in 0..32 {
982 if !transport.sent_rtp.lock().await.is_empty() {
983 break;
984 }
985 tokio::task::yield_now().await;
986 }
987 session.finish().await.unwrap();
988 let _ = pump.await.unwrap();
989
990 let sent = transport.sent_rtp.lock().await;
992 assert!(!sent.is_empty(), "VP9 egress sent RTP packets");
993 let mut depack = crate::protocol::rtp::Vp9Depacketizer::new();
994 let mut out = None;
995 for p in sent.iter() {
996 let h = RtpHeader::parse(p).unwrap();
997 if let Some(f) = depack
998 .push(&p[h.payload_offset..], h.marker, h.timestamp)
999 .unwrap()
1000 {
1001 out = Some(f);
1002 }
1003 }
1004 let out = out.expect("VP9 frame completed");
1005 assert_eq!(&out.data[..], &frame_data[..], "VP9 frame reconstructed");
1006 assert!(out.keyframe);
1007 }
1008}