From 5f152e9db7e32164d18605ecee0f9f32ef85d3c9 Mon Sep 17 00:00:00 2001 From: Doloro1978 Date: Fri, 7 Aug 2026 11:48:31 +0100 Subject: [PATCH] feat: WHIP support --- crates/server/src/webrtc_ingest.rs | 308 +++++++++++++------------- crates/server/src/webrtc_proxy.rs | 112 ++++++---- crates/server/tests/whip_sdp_probe.rs | 118 ++++++++++ 3 files changed, 349 insertions(+), 189 deletions(-) create mode 100644 crates/server/tests/whip_sdp_probe.rs diff --git a/crates/server/src/webrtc_ingest.rs b/crates/server/src/webrtc_ingest.rs index 627426a..fe1978a 100644 --- a/crates/server/src/webrtc_ingest.rs +++ b/crates/server/src/webrtc_ingest.rs @@ -64,29 +64,19 @@ pub async fn handle_whip_injest_patch( body: String, ) -> Result { // The slug is the stream_key_id (set by the POST handler's Location header). - let stream_key_id: i32 = slug - .parse() - .map_err(|e| { -warn!(%slug, "Whip PATCH: bad slug: {:?}", e); -HttpError::NotFound - })?; + let stream_key_id: i32 = slug.parse().map_err(|e| { + warn!(%slug, "Whip PATCH: bad slug: {:?}", e); + HttpError::NotFound + })?; info!(stream_key_id, ct = ?headers.get(header::CONTENT_TYPE), "Whip PATCH: trickle candidate"); // Look up the trickle sender for this session. - let trickle_map = state - .appstate - .lock() - .await - .webrtc_proxy - .trickle_tx - .clone(); - let tx = trickle_map - .get(&stream_key_id) - .ok_or_else(|| { -warn!(stream_key_id, "Whip PATCH: no trickle channel for session"); -HttpError::NotFound - })?; + let trickle_map = state.appstate.lock().await.webrtc_proxy.trickle_tx.clone(); + let tx = trickle_map.get(&stream_key_id).ok_or_else(|| { + warn!(stream_key_id, "Whip PATCH: no trickle channel for session"); + HttpError::NotFound + })?; debug!(stream_key_id, %body, "Whip PATCH: forwarding trickle candidate"); tx.send(body).ok(); @@ -145,113 +135,117 @@ pub async fn handle_whip_injest( }; info!(stream_key_id = key.id, label = %key.label, "Whip authenticated stream key from {}", remote); - // Parse the SDP offer first so we can detect the codec and - // configure the Rtc before accepting. - let sdp_offer = match SdpOffer::from_sdp_string(&offer) { - Ok(o) => o, - Err(e) => { - error!("Whip: cant parse offer from {}: {:?}", remote, e); - return err(StatusCode::BAD_REQUEST, "invalid SDP offer"); - } - }; - info!( - "Whip offer from {} — {} media line(s):", - remote, - sdp_offer.media_lines.len() - ); - for x in &sdp_offer.media_lines { - info!(" {}", x); + // Parse the SDP offer first so we can detect the codec and + // configure the Rtc before accepting. + let sdp_offer = match SdpOffer::from_sdp_string(&offer) { + Ok(o) => o, + Err(e) => { + error!("Whip: cant parse offer from {}: {:?}", remote, e); + return err(StatusCode::BAD_REQUEST, "invalid SDP offer"); } + }; + info!( + "Whip offer from {} — {} media line(s):", + remote, + sdp_offer.media_lines.len() + ); + for x in &sdp_offer.media_lines { + info!(" {}", x); + } - // Detect codec from the offer so we can enable matching codecs. - let stream_codec = video_codec_from_sdp_offer(&sdp_offer); - info!(?stream_codec, "detected codec from WHIP offer"); + // Detect codec from the offer so we can enable matching codecs. + let stream_codec = video_codec_from_sdp_offer(&sdp_offer); + info!(?stream_codec, "detected codec from WHIP offer"); - // Build Rtc with ICE-Lite (required by WHIP RFC 9728 §4.1) and - // matching codecs enabled. - // Not using ICE-Lite: full ICE lets the server initiate checks - // when OBS hasn't sent its candidates yet (Trickle ICE without PATCH). - let mut builder = Rtc::builder(); - { - let cc = builder.codec_config(); - cc.clear(); - cc.enable_opus(true); - match &stream_codec { - Some(crate::StreamCodec::H264) => { - info!("Whip: enabling H.264 codec"); - cc.enable_h264(true); - } - Some(crate::StreamCodec::H265) => { - info!("Whip: enabling H.265 codec"); - cc.enable_h265(true); - } - Some(crate::StreamCodec::AV1) => { - info!("Whip: enabling AV1 codec"); - cc.enable_av1(true); - } - None => { - warn!("Whip: no video codec detected in offer, enabling H.264 as fallback"); - cc.enable_h264(true); - } + // Build Rtc with ICE-Lite (required by WHIP RFC 9728 §4.1) and + // matching codecs enabled. + // Not using ICE-Lite: full ICE lets the server initiate checks + // when OBS hasn't sent its candidates yet (Trickle ICE without PATCH). + let mut builder = Rtc::builder(); + { + let cc = builder.codec_config(); + cc.clear(); + cc.enable_opus(true); + match &stream_codec { + Some(crate::StreamCodec::H264) => { + info!("Whip: enabling H.264 codec"); + cc.enable_h264(true); + } + Some(crate::StreamCodec::H265) => { + info!("Whip: enabling H.265 codec"); + cc.enable_h265(true); + } + Some(crate::StreamCodec::AV1) => { + info!("Whip: enabling AV1 codec"); + cc.enable_av1(true); + } + None => { + warn!("Whip: no video codec detected in offer, enabling H.264 as fallback"); + cc.enable_h264(true); } } - let mut rtc = builder.build(Instant::now()); - let candidate = Candidate::host(public_addr, Protocol::Udp).unwrap(); - rtc.add_local_candidate(candidate); - info!(%public_addr, "Whip: added local ICE candidate, accepting offer…"); + } + let mut rtc = builder.build(Instant::now()); + let candidate = Candidate::host(public_addr, Protocol::Udp).unwrap(); + rtc.add_local_candidate(candidate); + info!(%public_addr, "Whip: added local ICE candidate, accepting offer…"); - let offer_answe = match rtc.sdp_api().accept_offer(sdp_offer) { - Ok(a) => a, - Err(e) => { - error!("cant accept inject offer: {:?}", e); - return err(StatusCode::BAD_REQUEST, "could not accept offer"); - } - }; - - // OBS sends no a=candidate: lines (disableAutoGathering). Derive a - // remote host candidate from the HTTP source address so the server - // has somewhere to send STUN checks. - if let Ok(c) = Candidate::host(remote, Protocol::Udp) { - info!(%remote, "Whip: no candidates in offer, adding HTTP-derived remote host candidate"); - rtc.add_remote_candidate(c); + let offer_answe = match rtc.sdp_api().accept_offer(sdp_offer) { + Ok(a) => a, + Err(e) => { + error!("cant accept inject offer: {:?}", e); + return err(StatusCode::BAD_REQUEST, "could not accept offer"); } + }; - // Extract negotiated codec, PT, and profile from the answer - // so WHEP viewers can use the exact same codec config. - let (negotiated_codec, video_pt, video_profile) = - extract_negotiated_codec_info(&offer_answe); - info!( - ?negotiated_codec, - video_pt, - ?video_profile, - "negotiated codec from WHIP answer" - ); - if negotiated_codec.is_none() || video_pt.is_none() { - warn!("Whip: no common video codec negotiated, rejecting"); - return err(StatusCode::NOT_ACCEPTABLE, "no common video codec"); - } + // OBS sends no a=candidate: lines (disableAutoGathering). Derive a + // remote host candidate from the HTTP source address so the server + // has somewhere to send STUN checks. + if let Ok(c) = Candidate::host(remote, Protocol::Udp) { + info!(%remote, "Whip: no candidates in offer, adding HTTP-derived remote host candidate"); + rtc.add_remote_candidate(c); + } - let answer_sdp = offer_answe.to_sdp_string() - // Strip a=ice-options:trickle so OBS starts ICE immediately. - .replace("a=ice-options:trickle\r\n", "") - .replace("a=ice-options:trickle\n", ""); - // Fix up the answer for libdatachannel (OBS WHIP): - // 1. Strip a=group:BUNDLE — OBS doesn't negotiate it. - // 2. Add the host candidate to the video m= line. - let answer_sdp = answer_sdp - .replace("a=group:BUNDLE 0 1\r\n", "") - .replace("a=group:BUNDLE 0 1\n", ""); - let answer_sdp = if let Some(cand_line) = answer_sdp - .lines() - .find(|l| l.starts_with("a=candidate:")) - { + // Extract negotiated codec, PT, and profile from the answer + // so WHEP viewers can use the exact same codec config. + let (negotiated_codec, video_pt, video_profile) = extract_negotiated_codec_info(&offer_answe); + info!( + ?negotiated_codec, + video_pt, + ?video_profile, + "negotiated codec from WHIP answer" + ); + if negotiated_codec.is_none() || video_pt.is_none() { + warn!("Whip: no common video codec negotiated, rejecting"); + return err(StatusCode::NOT_ACCEPTABLE, "no common video codec"); + } + + let answer_sdp = offer_answe + .to_sdp_string() + // Strip a=ice-options:trickle so OBS starts ICE immediately. + .replace("a=ice-options:trickle\r\n", "") + .replace("a=ice-options:trickle\n", ""); + // Fix up the answer for libdatachannel (OBS WHIP): + // 1. Strip a=group:BUNDLE — OBS doesn't negotiate it. + // 2. Add the host candidate to the video m= line. + let answer_sdp = answer_sdp + .replace("a=group:BUNDLE 0 1\r\n", "") + .replace("a=group:BUNDLE 0 1\n", ""); + let answer_sdp = + if let Some(cand_line) = answer_sdp.lines().find(|l| l.starts_with("a=candidate:")) { // Insert the candidate line after the video m= line. - let cand_replacement = format!("m=video 9 UDP/TLS/RTP/SAVPF 96\r\nc=IN IP4 127.0.0.1\r\n{}\r\n", cand_line); - answer_sdp.replace("m=video 9 UDP/TLS/RTP/SAVPF 96\r\nc=IN IP4 0.0.0.0\r\n", &cand_replacement) + let cand_replacement = format!( + "m=video 9 UDP/TLS/RTP/SAVPF 96\r\nc=IN IP4 {}\r\n{}\r\n", + public_addr.ip(), cand_line + ); + answer_sdp.replace( + "m=video 9 UDP/TLS/RTP/SAVPF 96\r\nc=IN IP4 0.0.0.0\r\n", + &cand_replacement, + ) } else { answer_sdp }; - info!("Serving Whip SDP answer to {}:\n{}", remote, answer_sdp); + info!("Serving Whip SDP answer to {}:\n{}", remote, answer_sdp); let ufrag = answer_sdp .lines() @@ -267,12 +261,21 @@ pub async fn handle_whip_injest( .find(|l| l.starts_with("a=ice-pwd:")) .and_then(|l| l.strip_prefix("a=ice-pwd:")) .map(|s| s.trim().to_string()); - info!(?ufrag, ?ice_pwd, "Whip: registering ICE credentials with proxy"); + info!( + ?ufrag, + ?ice_pwd, + "Whip: registering ICE credentials with proxy" + ); let (socket, rx) = state.appstate.lock().await.webrtc_proxy.add_client(ufrag); let stream_sessions = state.appstate.lock().await.stream_sessions.clone(); - let (video_tx, video_rx) = broadcast::>(4); - let (audio_tx, audio_rx) = broadcast::>(4); + let (mut video_tx, video_rx) = broadcast::>(32); + let (mut audio_tx, audio_rx) = broadcast::>(32); + // Never block the ingest loop on slow/missing viewers: overwrite the + // oldest frame instead (same as the RTMP path; the viewer re-waits for a + // keyframe on Overflowed). The undrained rx below must not stall sends. + video_tx.set_overflow(true); + audio_tx.set_overflow(true); stream_sessions.insert( key.id, @@ -288,7 +291,10 @@ pub async fn handle_whip_injest( video_profile_level_id: video_profile, }, ); - info!(stream_key_id = key.id, "Whip: StreamSession inserted, spawning detach task"); + info!( + stream_key_id = key.id, + "Whip: StreamSession inserted, spawning detach task" + ); // Trickle-ICE channel: OBS can send candidates via PATCH after the // initial offer. We forward them to the Rtc task for add_remote_candidate. @@ -326,7 +332,10 @@ pub async fn handle_whip_injest( .header(header::LOCATION, &location) // Provide a STUN server via Link header so OBS can gather // ICE candidates even without explicit STUN configuration. - .header(header::LINK, "; rel=\"ice-server\"") + .header( + header::LINK, + "; rel=\"ice-server\"", + ) .body(axum::body::Body::from(answer_sdp)) .unwrap() } @@ -372,14 +381,14 @@ async fn detach_inject_rtc( let deadline = loop { match rtc.poll_output() { Ok(Output::Timeout(t)) => break t, - Ok(Output::Transmit(t)) => { - debug!( - "Whip TX: {} bytes → {}:{}", - t.contents.len(), - t.destination.ip(), - t.destination.port() - ); - if let Err(e) = socket.send_to(&t.contents, t.destination).await { + Ok(Output::Transmit(t)) => { + debug!( + "Whip TX: {} bytes → {}:{}", + t.contents.len(), + t.destination.ip(), + t.destination.port() + ); + if let Err(e) = socket.send_to(&t.contents, t.destination).await { warn!("Whip UDP send error: {:?}", e); cleanup().await; return; @@ -430,7 +439,10 @@ async fn detach_inject_rtc( } } Event::Connected => { - info!(stream_key_id, "Whip DTLS+ICE connected, wiring broadcast channels"); + info!( + stream_key_id, + "Whip DTLS+ICE connected, wiring broadcast channels" + ); if let Some(session) = sessions_ref.get(&stream_key_id) { video_tx = Some(session.frame_channel.clone()); audio_tx = Some(session.audio_channel.clone()); @@ -448,7 +460,7 @@ async fn detach_inject_rtc( }; let sleep = tokio::time::sleep_until(deadline.max(Instant::now()).into()); - tokio::select! { + tokio::select! { _ = sleep => { if let Err(e) = rtc.handle_input(Input::Timeout(Instant::now())) { error!(stream_key_id, "Whip handle_input(Timeout) error: {:?}", e); @@ -510,28 +522,28 @@ async fn detach_inject_rtc( /// Walks the answer's media lines, finds the first video m-line with /// negotiated rtp_params, and returns the codec + PT + H.264 profile. pub fn extract_negotiated_codec_info( -answer: &str0m::change::SdpAnswer, + answer: &str0m::change::SdpAnswer, ) -> (Option, Option, Option) { -use str0m::format::Codec; + use str0m::format::Codec; -for line in answer.media_lines.iter() { -// Iterate rtp_params on every m-line; video check via -// codec.is_video() avoids needing str0m's private MediaType. -for p in line.rtp_params() { -if p.spec().codec.is_video() { -let pt = Some(*p.pt()); -let profile = p.spec().format.profile_level_id; -let codec = match p.spec().codec { -Codec::H264 => Some(crate::StreamCodec::H264), -Codec::H265 => Some(crate::StreamCodec::H265), -Codec::Av1 => Some(crate::StreamCodec::AV1), -_ => None, -}; -return (codec, pt, profile); -} -} -} -(None, None, None) + for line in answer.media_lines.iter() { + // Iterate rtp_params on every m-line; video check via + // codec.is_video() avoids needing str0m's private MediaType. + for p in line.rtp_params() { + if p.spec().codec.is_video() { + let pt = Some(*p.pt()); + let profile = p.spec().format.profile_level_id; + let codec = match p.spec().codec { + Codec::H264 => Some(crate::StreamCodec::H264), + Codec::H265 => Some(crate::StreamCodec::H265), + Codec::Av1 => Some(crate::StreamCodec::AV1), + _ => None, + }; + return (codec, pt, profile); + } + } + } + (None, None, None) } /// Extract the video codec from an SDP offer's media lines. diff --git a/crates/server/src/webrtc_proxy.rs b/crates/server/src/webrtc_proxy.rs index 6b1250c..9585e4e 100644 --- a/crates/server/src/webrtc_proxy.rs +++ b/crates/server/src/webrtc_proxy.rs @@ -36,19 +36,31 @@ impl WebrtcProxy { let sock = UdpSocket::bind(format!("0.0.0.0:{}", config.proxy_port)).await?; let port = sock.local_addr()?.port(); - let public_ip = match env::var("PUBLIC_DOMAIN") { - Ok(domain) => { + let public_ip = match env::var("PUBLIC_DOMAIN") + .ok() + .filter(|s| !s.trim().is_empty()) + { + Some(domain) => { let ip = resolve_domain(&domain).await?; info!(%domain, %ip, "resolved PUBLIC_DOMAIN for WebRTC candidates"); ip } - Err(_) => { + None => { if cfg!(debug_assertions) { - // For testing - info!( - "using ip 0.0.0.0 for WebRTC candidates (because we are in a debug build)" - ); - IpAddr::from([127, 0, 0, 1]) + // For testing — advertise a real interface IP. Loopback is + // unreachable from clients whose ICE stack pins its UDP + // sockets to a specific interface (OBS sets + // IP_UNICAST_IF), which makes checks to 127.0.0.1 vanish. + match default_iface_ipv4() { + Some(ip) => { + info!(%ip, "using default interface IP for WebRTC candidates (debug build)"); + ip + } + None => { + warn!("no non-loopback IPv4 interface found, falling back to 127.0.0.1"); + IpAddr::from([127, 0, 0, 1]) + } + } } else { let ip = stun_public_ip().await?; info!(%ip, "discovered public IP via STUN for WebRTC candidates"); @@ -83,9 +95,14 @@ impl WebrtcProxy { }; let data = Bytes::copy_from_slice(&buf[..b]); - // By addr - if let Some(tx) = by_addr.get(&from) { - debug!("proxy: routing {} bytes by addr {}:{} → channel", b, from.ip(), from.port()); + // By addr + if let Some(tx) = by_addr.get(&from) { + debug!( + "proxy: routing {} bytes by addr {}:{} → channel", + b, + from.ip(), + from.port() + ); match tx.try_send((data, from)) { Ok(_) => continue, Err(e) => { @@ -110,16 +127,22 @@ impl WebrtcProxy { // might be a response to OUR STUN request (remote:local) // or an incoming request from the remote peer (local:remote). let part2_lookup = part2.clone(); - let entry = by_ufrag.remove(&part1).or_else(|| { - part2_lookup.and_then(|p2| by_ufrag.remove(&p2)) - }); + let entry = by_ufrag + .remove(&part1) + .or_else(|| part2_lookup.and_then(|p2| by_ufrag.remove(&p2))); let Some((_, tx)) = entry else { warn!("STUN packet ({}/{:?}), isnt registored", part1, part2); continue; }; by_addr.insert(from, tx.clone()); - info!("proxy: STUN match → promoted {} → ufrag={} (match was {}/{})", from, part1, part1, part2.as_deref().unwrap_or("-")); + info!( + "proxy: STUN match → promoted {} → ufrag={} (match was {}/{})", + from, + part1, + part1, + part2.as_deref().unwrap_or("-") + ); debug!("sending data"); if let Err(e) = tx.try_send((data, from)) { match e { @@ -148,33 +171,40 @@ impl WebrtcProxy { self.public_addr } pub fn ufrag_pair(b: &Bytes) -> Option<(String, Option)> { - if b.len() <= 20 { - return None; - } - let magic = u32::from_be_bytes(b[4..8].try_into().ok()?); - if magic != STUN_MAGIC { - return None; - } - // attribies start at 20 - let mut pos = 20usize; - while (pos + 4) <= b.len() { - let attr_type: u16 = u16::from_be_bytes(b[pos..pos + 2].try_into().ok()?); - let attr_len: u16 = u16::from_be_bytes(b[pos + 2..pos + 4].try_into().ok()?); - pos += 4; - if attr_type == 0x0006 { - let value = std::str::from_utf8( - b[pos..pos + (attr_len as usize)].try_into().ok()?, - ) - .ok()?; - let mut parts = value.split(':'); - let first = parts.next()?.to_string(); - let second = parts.next().map(|s| s.to_string()); - return Some((first, second)); - } - pos += (attr_len as usize + 3) & !3; - } - None + if b.len() <= 20 { + return None; } + let magic = u32::from_be_bytes(b[4..8].try_into().ok()?); + if magic != STUN_MAGIC { + return None; + } + // attribies start at 20 + let mut pos = 20usize; + while (pos + 4) <= b.len() { + let attr_type: u16 = u16::from_be_bytes(b[pos..pos + 2].try_into().ok()?); + let attr_len: u16 = u16::from_be_bytes(b[pos + 2..pos + 4].try_into().ok()?); + pos += 4; + if attr_type == 0x0006 { + let value = + std::str::from_utf8(b[pos..pos + (attr_len as usize)].try_into().ok()?).ok()?; + let mut parts = value.split(':'); + let first = parts.next()?.to_string(); + let second = parts.next().map(|s| s.to_string()); + return Some((first, second)); + } + pos += (attr_len as usize + 3) & !3; + } + None + } +} + +/// IP of the interface holding the default route, via a UDP connect() trick: +/// connect() only does a route lookup (no packets sent), so the kernel binds +/// the source address the OS would use for outbound traffic. +fn default_iface_ipv4() -> Option { +let sock = std::net::UdpSocket::bind("0.0.0.0:0").ok()?; +sock.connect("8.8.8.8:9").ok()?; +sock.local_addr().ok().map(|a| a.ip()) } async fn resolve_domain(domain: &str) -> Result> { diff --git a/crates/server/tests/whip_sdp_probe.rs b/crates/server/tests/whip_sdp_probe.rs new file mode 100644 index 0000000..5fdee14 --- /dev/null +++ b/crates/server/tests/whip_sdp_probe.rs @@ -0,0 +1,118 @@ +//! Temporary probe: replicate handle_whip_injest's Rtc setup, feed it a +//! realistic OBS/libdatachannel WHIP offer, print the answer SDP and check +//! whether the string-fixups in webrtc_ingest.rs actually match. +use std::{net::SocketAddr, time::Instant}; + +use str0m::{change::SdpOffer, net::Protocol, Candidate, Rtc}; + +fn obs_like_offer() -> String { + let mut s = String::new(); + s.push_str("v=0\r\n"); + s.push_str("o=- 4527835755568137757 2 IN IP4 127.0.0.1\r\n"); + s.push_str("s=-\r\n"); + s.push_str("t=0 0\r\n"); + s.push_str("a=group:BUNDLE 0 1\r\n"); + s.push_str("a=msid-semantic: WMS *\r\n"); + s.push_str("m=audio 9 UDP/TLS/RTP/SAVPF 111\r\n"); + s.push_str("c=IN IP4 0.0.0.0\r\n"); + s.push_str("a=rtcp:9 IN IP4 0.0.0.0\r\n"); + s.push_str("a=ice-ufrag:obs_ufrag_audio\r\n"); + s.push_str("a=ice-pwd:obs_pwd_audio\r\n"); + s.push_str("a=ice-options:trickle\r\n"); + s.push_str("a=fingerprint:sha-256 5A:6B:7C:8D:9E:AF:B0:C1:D2:E3:F4:05:16:27:38:49:5A:6B:7C:8D:9E:AF:B0:C1:D2:E3:F4:05:16:27:38:49:5A\r\n"); + s.push_str("a=setup:actpass\r\n"); + s.push_str("a=mid:0\r\n"); + s.push_str("a=sendrecv\r\n"); + s.push_str("a=rtcp-mux\r\n"); + s.push_str("a=rtpmap:111 opus/48000/2\r\n"); + s.push_str("a=rtcp-fb:111 transport-cc\r\n"); + s.push_str("a=fmtp:111 minptime=10;useinbandfec=1\r\n"); + s.push_str("m=video 9 UDP/TLS/RTP/SAVPF 96\r\n"); + s.push_str("c=IN IP4 0.0.0.0\r\n"); + s.push_str("a=rtcp:9 IN IP4 0.0.0.0\r\n"); + s.push_str("a=ice-ufrag:obs_ufrag_video\r\n"); + s.push_str("a=ice-pwd:obs_pwd_video\r\n"); + s.push_str("a=ice-options:trickle\r\n"); + s.push_str("a=fingerprint:sha-256 5A:6B:7C:8D:9E:AF:B0:C1:D2:E3:F4:05:16:27:38:49:5A:6B:7C:8D:9E:AF:B0:C1:D2:E3:F4:05:16:27:38:49:5A\r\n"); + s.push_str("a=setup:actpass\r\n"); + s.push_str("a=mid:1\r\n"); + s.push_str("a=sendrecv\r\n"); + s.push_str("a=rtcp-mux\r\n"); + s.push_str("a=rtpmap:96 H264/90000\r\n"); + s.push_str("a=rtcp-fb:96 nack\r\n"); + s.push_str("a=rtcp-fb:96 nack pli\r\n"); + s.push_str("a=rtcp-fb:96 transport-cc\r\n"); + s.push_str( + "a=fmtp:96 level-asymmetry-allowed=1;packetization-mode=1;profile-level-id=42e01f\r\n", + ); + s +} + +#[test] +fn probe_whip_answer() { + let offer_sdp = obs_like_offer(); + let sdp_offer = SdpOffer::from_sdp_string(&offer_sdp).expect("parse offer"); + + // Mirror handle_whip_injest's builder setup. + let mut builder = Rtc::builder(); + { + let cc = builder.codec_config(); + cc.clear(); + cc.enable_opus(true); + cc.enable_h264(true); + } + let mut rtc = builder.build(Instant::now()); + let public_addr: SocketAddr = "203.0.113.7:6969".parse().unwrap(); + let candidate = Candidate::host(public_addr, Protocol::Udp).unwrap(); + rtc.add_local_candidate(candidate); + + let answer = rtc.sdp_api().accept_offer(sdp_offer).expect("accept offer"); + let answer_sdp = answer.to_sdp_string(); + + println!("=== RAW ANSWER ==="); + println!("{answer_sdp}"); + println!("=== END RAW ANSWER ==="); + + // Now replicate the fixups from webrtc_ingest.rs verbatim. + let answer_sdp = answer_sdp + .replace("a=ice-options:trickle\r\n", "") + .replace("a=ice-options:trickle\n", ""); + let answer_sdp = answer_sdp + .replace("a=group:BUNDLE 0 1\r\n", "") + .replace("a=group:BUNDLE 0 1\n", ""); + let answer_sdp = + if let Some(cand_line) = answer_sdp.lines().find(|l| l.starts_with("a=candidate:")) { + let cand_replacement = format!( + "m=video 9 UDP/TLS/RTP/SAVPF 96\r\nc=IN IP4 127.0.0.1\r\n{}\r\n", + cand_line + ); + answer_sdp.replace( + "m=video 9 UDP/TLS/RTP/SAVPF 96\r\nc=IN IP4 0.0.0.0\r\n", + &cand_replacement, + ) + } else { + answer_sdp + }; + + println!("=== FIXED ANSWER ==="); + println!("{answer_sdp}"); + println!("=== END FIXED ANSWER ==="); + + // Diagnostics + let video_mline = answer_sdp + .lines() + .find(|l| l.starts_with("m=video")) + .unwrap(); + println!("video m-line after fixup: {video_mline}"); + println!( + "has a=candidate after fixup: {}", + answer_sdp.lines().any(|l| l.starts_with("a=candidate:")) + ); + println!( + "a=candidate lines: {:?}", + answer_sdp + .lines() + .filter(|l| l.starts_with("a=candidate:")) + .collect::>() + ); +}