Something went wrong. Try again.
dope ass quickshell bar
Something went wrong. Try again.
21 kB · 578 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};use std::sync::{Arc, Mutex, OnceLock};use std::thread::{self, JoinHandle};use std::time::Instant;
use anyhow::{Context, Result, anyhow};use gstreamer as gst;use gstreamer::prelude::*;use gstreamer_app::{AppSink, AppSinkCallbacks};use gstreamer_video::VideoInfo;use tokio::sync::watch;use tracing::{debug, info, warn};
use crate::cast_session::CastReady;use crate::config::RuntimeConfig;use crate::model::{CpuFrame, PixelFormat};
#[derive(Clone, Debug)]struct RawFrame { width: u32, height: u32, stride: u32, pixels: Arc<[u8]>,}
struct FrameGate { next_frame_at: Instant, interval: std::time::Duration,}
impl FrameGate { fn new(interval: std::time::Duration) -> Self { Self { next_frame_at: Instant::now(), interval, } }
fn take(&mut self, now: Instant) -> bool { if now < self.next_frame_at { return false; } self.next_frame_at = now + self.interval; true }}
pub fn init() -> Result<()> { gst::init().context("initializing GStreamer")}
pub fn start( ready: CastReady, config: RuntimeConfig, frame_tx: watch::Sender<Option<CpuFrame>>, stop: Arc<AtomicBool>,) -> JoinHandle<Result<()>> { thread::Builder::new() .name(format!("cast-{}", ready.window_id)) .spawn(move || run(ready, config, frame_tx, stop)) .expect("failed to spawn capture thread")}
fn run( ready: CastReady, config: RuntimeConfig, frame_tx: watch::Sender<Option<CpuFrame>>, stop: Arc<AtomicBool>,) -> Result<()> { // Keep DMA-BUF negotiation on the source side, then download to stable BGRA // system memory. Some NVIDIA DMA-BUF paths expose dimensions/PAR that make // GStreamer videoscale overflow during caps fixation, so bounded downscale // happens after readback instead. let negotiated_rate = config.fps.negotiated(); let frame_interval = config.fps.interval(); let description = format!( "pipewiresrc path={} do-timestamp=true ! video/x-raw(memory:DMABuf),max-framerate={} ! glupload ! gldownload ! videoconvert ! video/x-raw,format=BGRA ! appsink name=sink emit-signals=true sync=false max-buffers=1 drop=true", ready.pw_node_id, negotiated_rate ); let max_width = config.max_width; info!( window_id = ready.window_id, session_path = %ready.session_path, stream_path = %ready.stream_path, requested_fps = %config.fps, negotiated_fps = %negotiated_rate, frame_interval = ?frame_interval, pipeline = %description, "starting PipeWire capture" ); let pipeline = gst::parse::launch(&description) .with_context(|| format!("building PipeWire pipeline for window {}", ready.window_id))? .downcast::<gst::Pipeline>() .map_err(|_| anyhow!("capture pipeline is not a GstPipeline"))?; let sink = pipeline .by_name("sink") .context("capture pipeline has no appsink")? .downcast::<AppSink>() .map_err(|_| anyhow!("capture sink is not an appsink"))?;
let sequence = Arc::new(AtomicU64::new(0)); let frame_gate = Arc::new(Mutex::new(FrameGate::new(frame_interval))); let frame_gate_callback = frame_gate.clone(); let frame_tx_callback = frame_tx.clone(); let sequence_callback = sequence.clone(); sink.set_callbacks( AppSinkCallbacks::builder() .new_sample(move |sink| { let sample = sink.pull_sample().map_err(|_| gst::FlowError::Eos)?; let mut frame_gate = frame_gate_callback .lock() .map_err(|_| gst::FlowError::Error)?; if !frame_gate.take(Instant::now()) { return Ok(gst::FlowSuccess::Ok); } drop(frame_gate); let raw = decode_sample(&sample, max_width).map_err(|error| { warn!(?error, "dropping malformed PipeWire frame"); gst::FlowError::Error })?; let frame = CpuFrame { window_id: ready.window_id, sequence: sequence_callback.fetch_add(1, Ordering::Relaxed) + 1, monotonic_ns: monotonic_now_ns(), width: raw.width, height: raw.height, stride: raw.stride, format: PixelFormat::Bgra8888, pixels: raw.pixels, }; let _ = frame_tx_callback.send(Some(frame)); Ok(gst::FlowSuccess::Ok) }) .build(), );
pipeline .set_state(gst::State::Playing) .map_err(|error| anyhow!("starting capture pipeline: {error:?}"))?; let bus = pipeline.bus().context("capture pipeline has no bus")?; while !stop.load(Ordering::Acquire) { if let Some(message) = bus.timed_pop(gst::ClockTime::from_mseconds(100)) { use gst::MessageView; match message.view() { MessageView::Error(error) => { let source = error.src().map(|source| source.path_string()); let detail = format!("{} ({:?})", error.error(), error.debug()); let _ = pipeline.set_state(gst::State::Null); return Err(anyhow!("PipeWire capture error from {source:?}: {detail}")); } MessageView::Eos(..) => { let _ = pipeline.set_state(gst::State::Null); return Err(anyhow!("PipeWire capture reached EOS")); } _ => {} } } } debug!(window_id = ready.window_id, "stopping PipeWire capture"); let _ = pipeline.set_state(gst::State::Null); Ok(())}
fn decode_sample(sample: &gst::Sample, max_width: u32) -> Result<RawFrame> { let caps = sample.caps().context("sample has no caps")?; let info = VideoInfo::from_caps(caps).context("decoding sample video caps")?; let buffer = sample.buffer().context("sample has no buffer")?; let map = buffer.map_readable().context("mapping sample buffer")?; let width = info.width(); let height = info.height(); if width == 0 || height == 0 { return Err(anyhow!("sample has empty dimensions")); } let stride = u32::try_from(info.stride()[0]).context("negative sample stride")?; let row_bytes = usize::try_from(width)? .checked_mul(4) .context("sample row size overflow")?; let stride_usize = usize::try_from(stride)?; let height_usize = usize::try_from(height)?; let required = stride_usize .checked_mul(height_usize) .context("sample size overflow")?; if map.size() < required || stride_usize < row_bytes { return Err(anyhow!( "sample buffer is smaller than its advertised video layout" )); }
// max_width=0 is the explicit native-resolution escape hatch. // Otherwise publish through BOX -> Lanczos3 -> weak unsharp downscaling. let target_width = if max_width == 0 { width } else { width.min(max_width) }; if target_width == width { let mut pixels = vec![0u8; required]; for row in 0..height_usize { let source = &map.as_slice()[row * stride_usize..row * stride_usize + stride_usize]; pixels[row * stride_usize..(row + 1) * stride_usize].copy_from_slice(source); } return Ok(RawFrame { width, height, stride, pixels: pixels.into(), }); }
let target_height = scaled_height(width, height, target_width)?; let packed_size = row_bytes .checked_mul(height_usize) .context("packed sample size overflow")?; let mut source = vec![0u8; packed_size]; for row in 0..height_usize { let source_row = &map.as_slice()[row * stride_usize..row * stride_usize + row_bytes]; source[row * row_bytes..(row + 1) * row_bytes].copy_from_slice(source_row); }
let pixels = resize_bgra(&source, width, height, target_width, target_height)?; let stride = target_width .checked_mul(4) .context("scaled sample stride overflow")?; Ok(RawFrame { width: target_width, height: target_height, stride, pixels: pixels.into(), })}
fn scaled_height(width: u32, height: u32, target_width: u32) -> Result<u32> { let height = (u64::from(height) * u64::from(target_width) + u64::from(width) / 2) / u64::from(width); u32::try_from(height) .context("scaled sample height overflow") .map(|height| height.max(1))}
fn resize_bgra( source: &[u8], source_width: u32, source_height: u32, target_width: u32, target_height: u32,) -> Result<Vec<u8>> { if source_width == 0 || source_height == 0 || target_width == 0 || target_height == 0 { return Err(anyhow!("cannot resize an empty BGRA image")); } let source_size = usize::try_from(source_width)? .checked_mul(usize::try_from(source_height)?) .context("source image dimensions overflow")? .checked_mul(4) .context("source image size overflow")?; if source.len() != source_size { return Err(anyhow!("source BGRA image has an unexpected size")); }
let mut pixels = source.to_vec(); let mut width = source_width; let mut height = source_height;
// Reduce large images in bounded area-filtered steps first. This avoids asking // the final convolution to span hundreds of source pixels per output pixel. while width > target_width.saturating_mul(2) || height > target_height.saturating_mul(2) { let next_width = ((width + 1) / 2).max(target_width); let next_height = ((height + 1) / 2).max(target_height); pixels = box_reduce_bgra(&pixels, width, height, next_width, next_height)?; width = next_width; height = next_height; }
// Finish at the exact aspect-preserving dimensions with a high-quality // Lanczos3 convolution, then restore a small amount of edge contrast. pixels = lanczos3_resize_bgra(&pixels, width, height, target_width, target_height)?; apply_weak_unsharp_bgra(&mut pixels, target_width, target_height); Ok(pixels)}
fn box_reduce_bgra( source: &[u8], source_width: u32, source_height: u32, target_width: u32, target_height: u32,) -> Result<Vec<u8>> { if target_width == 0 || target_height == 0 || target_width > source_width || target_height > source_height { return Err(anyhow!("invalid BOX reduction dimensions")); } let source_width = usize::try_from(source_width)?; let source_height = usize::try_from(source_height)?; let target_width = usize::try_from(target_width)?; let target_height = usize::try_from(target_height)?; let target_size = target_width .checked_mul(target_height) .context("BOX image dimensions overflow")? .checked_mul(4) .context("BOX image size overflow")?; let mut target = vec![0u8; target_size];
let total_area = u64::try_from(source_width)? .checked_mul(u64::try_from(source_height)?) .context("BOX source area overflow")?; for y in 0..target_height { let output_y0 = y * source_height; let output_y1 = (y + 1) * source_height; let source_y0 = output_y0 / target_height; let source_y1 = ((y + 1) * source_height).div_ceil(target_height); for x in 0..target_width { let output_x0 = x * source_width; let output_x1 = (x + 1) * source_width; let source_x0 = x * source_width / target_width; let source_x1 = ((x + 1) * source_width).div_ceil(target_width); let mut sum = [0u64; 4]; for source_y in source_y0..source_y1 { let overlap_y = (output_y1.min((source_y + 1) * target_height) - output_y0.max(source_y * target_height)) as u64; for source_x in source_x0..source_x1 { let overlap_x = (output_x1.min((source_x + 1) * target_width) - output_x0.max(source_x * target_width)) as u64; let weight = overlap_x * overlap_y; let offset = (source_y * source_width + source_x) * 4; for channel in 0..4 { sum[channel] += u64::from(source[offset + channel]) * weight; } } } let target_offset = (y * target_width + x) * 4; for channel in 0..4 { target[target_offset + channel] = ((sum[channel] + total_area / 2) / total_area) as u8; } } } Ok(target)}
fn lanczos3_resize_bgra( source: &[u8], source_width: u32, source_height: u32, target_width: u32, target_height: u32,) -> Result<Vec<u8>> { let source_width = usize::try_from(source_width)?; let source_height = usize::try_from(source_height)?; let target_width = usize::try_from(target_width)?; let target_height = usize::try_from(target_height)?; let horizontal = lanczos_weights(source_width, target_width); let vertical = lanczos_weights(source_height, target_height); let intermediate_size = target_width .checked_mul(source_height) .context("Lanczos intermediate dimensions overflow")? .checked_mul(4) .context("Lanczos intermediate image size overflow")?; let target_size = target_width .checked_mul(target_height) .context("Lanczos dimensions overflow")? .checked_mul(4) .context("Lanczos image size overflow")?; let mut intermediate = vec![0u8; intermediate_size]; let mut target = vec![0u8; target_size];
for y in 0..source_height { for x in 0..target_width { let target_offset = (y * target_width + x) * 4; let mut values = [0.0f32; 4]; for &(source_x, weight) in &horizontal[x] { let source_offset = (y * source_width + source_x) * 4; for channel in 0..4 { values[channel] += f32::from(source[source_offset + channel]) * weight; } } for channel in 0..4 { intermediate[target_offset + channel] = values[channel].clamp(0.0, 255.0).round() as u8; } } }
for y in 0..target_height { for x in 0..target_width { let target_offset = (y * target_width + x) * 4; let mut values = [0.0f32; 4]; for &(source_y, weight) in &vertical[y] { let source_offset = (source_y * target_width + x) * 4; for channel in 0..4 { values[channel] += f32::from(intermediate[source_offset + channel]) * weight; } } for channel in 0..4 { target[target_offset + channel] = values[channel].clamp(0.0, 255.0).round() as u8; } } } Ok(target)}
fn lanczos_weights(source_length: usize, target_length: usize) -> Vec<Vec<(usize, f32)>> { let scale = source_length as f32 / target_length as f32; let support = 3.0 * scale.max(1.0); (0..target_length) .map(|target_index| { let center = (target_index as f32 + 0.5) * scale - 0.5; let first = (center - support).floor().max(0.0) as usize; let last = (center + support) .ceil() .min(source_length.saturating_sub(1) as f32) as usize; let mut weights = Vec::with_capacity(last.saturating_sub(first) + 1); let mut total = 0.0; for source_index in first..=last { let distance = (source_index as f32 - center) / scale; let weight = lanczos_kernel(distance); if weight != 0.0 { weights.push((source_index, weight)); total += weight; } } if total.abs() > f32::EPSILON { for (_, weight) in &mut weights { *weight /= total; } } else { weights.clear(); weights.push(( center.round().clamp(0.0, source_length as f32 - 1.0) as usize, 1.0, )); } weights }) .collect()}
fn lanczos_kernel(value: f32) -> f32 { if value.abs() >= 3.0 { 0.0 } else { sinc(value) * sinc(value / 3.0) }}
fn sinc(value: f32) -> f32 { if value.abs() < 1.0e-5 { 1.0 } else { let pi_value = std::f32::consts::PI * value; pi_value.sin() / pi_value }}
fn apply_weak_unsharp_bgra(pixels: &mut [u8], width: u32, height: u32) { const AMOUNT: f32 = 0.20; let width = usize::try_from(width).unwrap_or(0); let height = usize::try_from(height).unwrap_or(0); if width < 2 || height < 2 { return; } let original = pixels.to_vec(); for y in 0..height { for x in 0..width { let offset = (y * width + x) * 4; for channel in 0..3 { let mut blur_sum = 0u32; for blur_y in y.saturating_sub(1)..=(y + 1).min(height - 1) { for blur_x in x.saturating_sub(1)..=(x + 1).min(width - 1) { blur_sum += u32::from(original[(blur_y * width + blur_x) * 4 + channel]); } } let blur_count = ((y + 1).min(height - 1) - y.saturating_sub(1) + 1) * ((x + 1).min(width - 1) - x.saturating_sub(1) + 1); let value = f32::from(original[offset + channel]); let blurred = blur_sum as f32 / blur_count as f32; pixels[offset + channel] = (value + AMOUNT * (value - blurred)) .clamp(0.0, 255.0) .round() as u8; } pixels[offset + 3] = original[offset + 3]; } }}
fn monotonic_now_ns() -> u64 { static START: OnceLock<Instant> = OnceLock::new(); START .get_or_init(Instant::now) .elapsed() .as_nanos() .try_into() .unwrap_or(u64::MAX)}
#[cfg(test)]mod tests { use super::*;
#[test] fn box_reduction_averages_each_area() { let mut source = vec![0u8; 4 * 4 * 4]; for y in 0..4 { for x in 0..4 { let value = (y * 4 + x) as u8; let offset = (y * 4 + x) * 4; source[offset..offset + 4].copy_from_slice(&[value, value, value, 255]); } }
let reduced = box_reduce_bgra(&source, 4, 4, 2, 2).expect("BOX reduction");
assert_eq!(&reduced[0..4], &[3, 3, 3, 255]); assert_eq!(&reduced[4..8], &[5, 5, 5, 255]); assert_eq!(&reduced[8..12], &[11, 11, 11, 255]); assert_eq!(&reduced[12..16], &[13, 13, 13, 255]); }
#[test] fn box_reduction_weights_fractional_edge_coverage() { let mut source = vec![0u8; 5 * 1 * 4]; for x in 0..5 { let offset = x * 4; source[offset..offset + 4].copy_from_slice(&[x as u8, x as u8, x as u8, 255]); }
let reduced = box_reduce_bgra(&source, 5, 1, 4, 1).expect("fractional BOX reduction");
assert_eq!( reduced .chunks_exact(4) .map(|pixel| pixel[0]) .collect::<Vec<_>>(), vec![0, 1, 3, 4] ); }
#[test] fn resize_reaches_exact_dimensions_after_repeated_reductions() { let source = vec![96u8; 64 * 32 * 4]; let resized = resize_bgra(&source, 64, 32, 7, 3).expect("resize");
assert_eq!(resized.len(), 7 * 3 * 4); assert!( resized .chunks_exact(4) .all(|pixel| pixel == [96, 96, 96, 96]) ); }
#[test] fn weak_unsharp_preserves_alpha() { let mut pixels = vec![0u8; 3 * 3 * 4]; for pixel in pixels.chunks_exact_mut(4) { pixel[0..3].copy_from_slice(&[32, 64, 96]); pixel[3] = 217; }
apply_weak_unsharp_bgra(&mut pixels, 3, 3);
assert!(pixels.chunks_exact(4).all(|pixel| pixel[3] == 217)); } #[test] fn frame_gate_holds_sub_one_rate_until_interval() { let start = Instant::now(); let mut gate = FrameGate { next_frame_at: start, interval: std::time::Duration::from_secs(5), };
assert!(gate.take(start)); assert!(!gate.take(start + std::time::Duration::from_secs(4))); assert!(gate.take(start + std::time::Duration::from_secs(5))); }}