From c4a7e622b7abf80179543a18c00eefe9a92676dd Mon Sep 17 00:00:00 2001 From: Doloro1978 Date: Sat, 29 Aug 2026 15:15:28 +0100 Subject: [PATCH] fmt all --- crates/server/src/audio.rs | 4 +-- crates/server/src/codec/av1.rs | 41 ++++++++++++++++++--------- crates/server/src/codec/h264.rs | 31 +++++++++++++------- crates/server/src/codec/h265.rs | 30 +++++++++++++------- crates/server/src/http.rs | 4 ++- crates/server/src/http_error.rs | 2 +- crates/server/src/main.rs | 5 +--- crates/server/src/webrtc_ingest.rs | 8 +++++- crates/server/src/webrtc_proxy.rs | 28 ++++++++++-------- crates/server/tests/whip_sdp_probe.rs | 2 +- 10 files changed, 99 insertions(+), 56 deletions(-) diff --git a/crates/server/src/audio.rs b/crates/server/src/audio.rs index ba01814..4267934 100644 --- a/crates/server/src/audio.rs +++ b/crates/server/src/audio.rs @@ -1,10 +1,10 @@ use std::error::Error; use bytes::Bytes; -use rubato::{audioadapter_buffers::direct::InterleavedSlice, Fft, Resampler}; +use rubato::{Fft, Resampler, audioadapter_buffers::direct::InterleavedSlice}; use symphonia::core::{ audio::SampleBuffer, - codecs::{CodecParameters, Decoder, DecoderOptions, CODEC_TYPE_AAC}, + codecs::{CODEC_TYPE_AAC, CodecParameters, Decoder, DecoderOptions}, formats::Packet, }; diff --git a/crates/server/src/codec/av1.rs b/crates/server/src/codec/av1.rs index 158fa73..120d4ca 100644 --- a/crates/server/src/codec/av1.rs +++ b/crates/server/src/codec/av1.rs @@ -50,7 +50,9 @@ impl Av1CodecParser { let has_size = (header >> 1) & 1 != 0; i += 1; if has_extension { - if i >= data.len() { return false; } + if i >= data.len() { + return false; + } i += 1; } if obu_type == 1 { @@ -61,13 +63,19 @@ impl Av1CodecParser { let mut size: usize = 0; let mut shift = 0; loop { - if i >= data.len() { return false; } + if i >= data.len() { + return false; + } let b = data[i] as usize; i += 1; size |= (b & 0x7F) << shift; shift += 7; - if b & 0x80 == 0 { break; } - if shift > 32 { return false; } + if b & 0x80 == 0 { + break; + } + if shift > 32 { + return false; + } } i += size; } else { @@ -78,7 +86,11 @@ impl Av1CodecParser { false } - fn obus_for_frame(&mut self, payload: &[u8], rtmp_is_keyframe: bool) -> Option<(Vec, bool)> { + fn obus_for_frame( + &mut self, + payload: &[u8], + rtmp_is_keyframe: bool, + ) -> Option<(Vec, bool)> { if payload.is_empty() { return None; } @@ -95,13 +107,12 @@ impl Av1CodecParser { self.first_coded_frame = false; - if is_keyframe - && let Some(config) = &self.config_obus { - let mut out = Vec::with_capacity(config.len() + payload.len()); - out.extend_from_slice(config); - out.extend_from_slice(payload); - return Some((out, true)); - } + if is_keyframe && let Some(config) = &self.config_obus { + let mut out = Vec::with_capacity(config.len() + payload.len()); + out.extend_from_slice(config); + out.extend_from_slice(payload); + return Some((out, true)); + } Some((payload.to_vec(), is_keyframe)) } @@ -129,7 +140,11 @@ impl CodecParser for Av1CodecParser { // FourCC — enhanced RTMP only defines CTS for hvc1 CodedFrames. let payload = data.get(5..)?; let (obus, is_keyframe) = self.obus_for_frame(payload, rtmp_is_keyframe)?; - Some(VideoFrame { data: Bytes::from(obus), is_keyframe, timestamp_ms }) + Some(VideoFrame { + data: Bytes::from(obus), + is_keyframe, + timestamp_ms, + }) } _ => None, } diff --git a/crates/server/src/codec/h264.rs b/crates/server/src/codec/h264.rs index 8088e93..d05e6b9 100644 --- a/crates/server/src/codec/h264.rs +++ b/crates/server/src/codec/h264.rs @@ -53,12 +53,20 @@ impl H264CodecParser { } let pts_ms = Self::pts_ms(timestamp_ms, &bytes[5..8]); let data = self.avcc_to_annexb(bytes.get(8..)?, is_keyframe)?; - Some(VideoFrame { data: Bytes::from(data), is_keyframe, timestamp_ms: pts_ms }) + Some(VideoFrame { + data: Bytes::from(data), + is_keyframe, + timestamp_ms: pts_ms, + }) } 3 => { // CodedFramesX: no CTS field, bytes 5+ = AVCC NALUs let data = self.avcc_to_annexb(bytes.get(5..)?, is_keyframe)?; - Some(VideoFrame { data: Bytes::from(data), is_keyframe, timestamp_ms }) + Some(VideoFrame { + data: Bytes::from(data), + is_keyframe, + timestamp_ms, + }) } _ => None, } @@ -77,7 +85,11 @@ impl H264CodecParser { // PTS = DTS + CTS — see note on CodedFrames above. let pts_ms = Self::pts_ms(timestamp_ms, &bytes[2..5]); let data = self.avcc_to_annexb(&bytes[5..], is_keyframe)?; - Some(VideoFrame { data: Bytes::from(data), is_keyframe, timestamp_ms: pts_ms }) + Some(VideoFrame { + data: Bytes::from(data), + is_keyframe, + timestamp_ms: pts_ms, + }) } _ => None, } @@ -148,13 +160,12 @@ impl H264CodecParser { // Prepend SPS+PPS before every keyframe so str0m's packetizer // can bundle them into a STAP-A alongside the IDR NALU. - if is_keyframe - && let (Some(sps), Some(pps)) = (&self.sps, &self.pps) { - out.extend_from_slice(&[0, 0, 0, 1]); - out.extend_from_slice(sps); - out.extend_from_slice(&[0, 0, 0, 1]); - out.extend_from_slice(pps); - } + if is_keyframe && let (Some(sps), Some(pps)) = (&self.sps, &self.pps) { + out.extend_from_slice(&[0, 0, 0, 1]); + out.extend_from_slice(sps); + out.extend_from_slice(&[0, 0, 0, 1]); + out.extend_from_slice(pps); + } // Convert each length-prefixed NALU to an Annex B start-code NALU. let mut i = 0; diff --git a/crates/server/src/codec/h265.rs b/crates/server/src/codec/h265.rs index 1dab529..e442994 100644 --- a/crates/server/src/codec/h265.rs +++ b/crates/server/src/codec/h265.rs @@ -74,15 +74,15 @@ impl H265CodecParser { fn hvcc_to_annexb(&self, payload: &[u8], is_keyframe: bool) -> Option> { let mut out = Vec::with_capacity(payload.len()); - if is_keyframe - && let (Some(vps), Some(sps), Some(pps)) = (&self.vps, &self.sps, &self.pps) { - out.extend_from_slice(&[0, 0, 0, 1]); - out.extend_from_slice(vps); - out.extend_from_slice(&[0, 0, 0, 1]); - out.extend_from_slice(sps); - out.extend_from_slice(&[0, 0, 0, 1]); - out.extend_from_slice(pps); - } + if is_keyframe && let (Some(vps), Some(sps), Some(pps)) = (&self.vps, &self.sps, &self.pps) + { + out.extend_from_slice(&[0, 0, 0, 1]); + out.extend_from_slice(vps); + out.extend_from_slice(&[0, 0, 0, 1]); + out.extend_from_slice(sps); + out.extend_from_slice(&[0, 0, 0, 1]); + out.extend_from_slice(pps); + } let mut i = 0; while i + 4 <= payload.len() { @@ -128,12 +128,20 @@ impl CodecParser for H265CodecParser { let cts = i32::from_be_bytes([0, data[5], data[6], data[7]]) << 8 >> 8; let pts_ms = (timestamp_ms as i64 + cts as i64).max(0) as u32; let annexb = self.hvcc_to_annexb(data.get(8..)?, is_keyframe)?; - Some(VideoFrame { data: Bytes::from(annexb), is_keyframe, timestamp_ms: pts_ms }) + Some(VideoFrame { + data: Bytes::from(annexb), + is_keyframe, + timestamp_ms: pts_ms, + }) } 3 => { // CodedFramesX: no CTS, bytes 5+ = HVCC let annexb = self.hvcc_to_annexb(data.get(5..)?, is_keyframe)?; - Some(VideoFrame { data: Bytes::from(annexb), is_keyframe, timestamp_ms }) + Some(VideoFrame { + data: Bytes::from(annexb), + is_keyframe, + timestamp_ms, + }) } _ => None, } diff --git a/crates/server/src/http.rs b/crates/server/src/http.rs index 8d44e2f..1defb4f 100644 --- a/crates/server/src/http.rs +++ b/crates/server/src/http.rs @@ -13,7 +13,9 @@ use axum_extra::extract::{CookieJar, cookie::Cookie}; /// dashes and apostrophes — no spaces. Mirrors the frontend /// `/^[A-Za-z0-9'-]{0,67}$/` used by the keys-page popups. fn valid_charset(s: &str) -> bool { - s.len() <= 67 && s.chars().all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '\'') + s.len() <= 67 + && s.chars() + .all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '\'') } use axum::{ diff --git a/crates/server/src/http_error.rs b/crates/server/src/http_error.rs index 188cf03..b1a6c59 100644 --- a/crates/server/src/http_error.rs +++ b/crates/server/src/http_error.rs @@ -42,7 +42,7 @@ impl HttpError { Self::BadRequest(_) => StatusCode::BAD_REQUEST, Self::Unprocessable(_) => StatusCode::UNPROCESSABLE_ENTITY, Self::NotAcceptable(_) => StatusCode::NOT_ACCEPTABLE, - Self::WhepCodecError(_) => StatusCode::UNSUPPORTED_MEDIA_TYPE, + Self::WhepCodecError(_) => StatusCode::UNSUPPORTED_MEDIA_TYPE, } } } diff --git a/crates/server/src/main.rs b/crates/server/src/main.rs index aeec6e2..7c6d671 100644 --- a/crates/server/src/main.rs +++ b/crates/server/src/main.rs @@ -16,10 +16,7 @@ use tracing_subscriber::EnvFilter; use uuid::Uuid; use crate::{ - audio::OpusAudioFrame, - codec::VideoFrame, - http::HttpServer, - webrtc_proxy::WebrtcProxy, + audio::OpusAudioFrame, codec::VideoFrame, http::HttpServer, webrtc_proxy::WebrtcProxy, }; const SERVER_VERSION: &str = env!("CARGO_PKG_VERSION"); diff --git a/crates/server/src/webrtc_ingest.rs b/crates/server/src/webrtc_ingest.rs index 0bbe726..dd7c54f 100644 --- a/crates/server/src/webrtc_ingest.rs +++ b/crates/server/src/webrtc_ingest.rs @@ -71,7 +71,13 @@ pub async fn handle_whip_injest_delete( state.appstate.lock().await.stream_sessions.remove(&key.id); // Drop the trickle channel too, so PATCHes to a deleted session 404 and // the lingering detach task's own cleanup can't reach a newer session. - state.appstate.lock().await.webrtc_proxy.trickle_tx.remove(&key.id); + state + .appstate + .lock() + .await + .webrtc_proxy + .trickle_tx + .remove(&key.id); Ok(StatusCode::OK) } diff --git a/crates/server/src/webrtc_proxy.rs b/crates/server/src/webrtc_proxy.rs index 9f5db81..005713e 100644 --- a/crates/server/src/webrtc_proxy.rs +++ b/crates/server/src/webrtc_proxy.rs @@ -34,10 +34,10 @@ impl WebrtcProxy { let sock = UdpSocket::bind(format!("0.0.0.0:{}", proxy_port)).await?; let port = sock.local_addr()?.port(); - let public_ip = match env::var("PUBLIC_DOMAIN") - .ok() - .filter(|s| !s.trim().is_empty()) - { + 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"); @@ -51,12 +51,14 @@ impl WebrtcProxy { // 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 + 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]) + warn!( + "no non-loopback IPv4 interface found, falling back to 127.0.0.1" + ); + IpAddr::from([127, 0, 0, 1]) } } } else { @@ -200,9 +202,9 @@ fn stun_attributes(data: &[u8]) -> impl Iterator { /// 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()) + 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> { @@ -254,7 +256,9 @@ fn parse_xor_mapped_address(data: &[u8]) -> Option { // byte 0: reserved, byte 1: family (0x01=IPv4, 0x02=IPv6) if value[1] == 0x01 { let x_addr = u32::from_be_bytes(value[4..8].try_into().ok()?); - Some(std::net::IpAddr::V4(std::net::Ipv4Addr::from(x_addr ^ magic))) + Some(std::net::IpAddr::V4(std::net::Ipv4Addr::from( + x_addr ^ magic, + ))) } else { None } diff --git a/crates/server/tests/whip_sdp_probe.rs b/crates/server/tests/whip_sdp_probe.rs index 5fdee14..2cb42bb 100644 --- a/crates/server/tests/whip_sdp_probe.rs +++ b/crates/server/tests/whip_sdp_probe.rs @@ -3,7 +3,7 @@ //! whether the string-fixups in webrtc_ingest.rs actually match. use std::{net::SocketAddr, time::Instant}; -use str0m::{change::SdpOffer, net::Protocol, Candidate, Rtc}; +use str0m::{Candidate, Rtc, change::SdpOffer, net::Protocol}; fn obs_like_offer() -> String { let mut s = String::new();