Something went wrong. Try again.
dope ass quickshell bar
Something went wrong. Try again.
7.8 kB · 217 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218use std::collections::HashMap;use std::sync::Arc;use std::sync::atomic::{AtomicBool, Ordering};use std::time::Duration;
use anyhow::{Context, Result, anyhow};use tokio::sync::mpsc;use tokio::task::JoinHandle;use tokio_util::sync::CancellationToken;use tracing::{info, warn};use zbus::Connection;
use crate::capture;use crate::cast_session::CastSession;use crate::config::RuntimeConfig;use crate::model::{CaptureStatus, WindowEvent, WindowMeta};use crate::publish::Publisher;use crate::window_registry::WindowRegistry;
struct Worker { cancel: CancellationToken, task: JoinHandle<()>,}
pub async fn run(config: RuntimeConfig, publisher: Publisher) -> Result<()> { let (niri_tx, mut niri_rx) = mpsc::channel(128); let mut niri_task = tokio::spawn(crate::niri_ipc::run_event_stream(niri_tx, config.reconnect)); let shutdown = CancellationToken::new(); let mut registry = WindowRegistry::default(); let mut sigterm = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) .context("installing SIGTERM handler")?; let mut workers: HashMap<u64, Worker> = HashMap::new(); let mut event_stream_error = None;
loop { tokio::select! { result = &mut niri_task => { let result = result .context("niri EventStream task failed") .and_then(|result| result.context("niri EventStream stopped")); if let Err(error) = result { warn!(?error, "niri EventStream stopped"); event_stream_error = Some(error); } break; } Some(event) = niri_rx.recv() => { for change in registry.apply(event) { reconcile(change, &config, &publisher, &shutdown, &mut workers).await; } } signal = tokio::signal::ctrl_c() => { signal.context("waiting for Ctrl-C")?; break; } _ = sigterm.recv() => { info!("received SIGTERM"); break; } } }
shutdown.cancel(); for worker in workers.values() { worker.cancel.cancel(); } for (_, worker) in workers { let _ = worker.task.await; } info!("daemon stopped"); match event_stream_error { Some(error) => Err(error), None => Ok(()), }}
async fn reconcile( change: WindowEvent, config: &RuntimeConfig, publisher: &Publisher, shutdown: &CancellationToken, workers: &mut HashMap<u64, Worker>,) { match change { WindowEvent::Upsert(meta) => { if !config.eligible(meta.app_id.as_deref()) { stop_worker(meta.id, publisher, workers).await; return; } publisher.upsert(meta.clone()); if workers.contains_key(&meta.id) { return; } let cancel = shutdown.child_token(); let worker_cancel = cancel.clone(); let worker_publisher = publisher.clone(); let worker_config = config.clone(); let worker_meta = meta.clone(); let task = tokio::spawn(async move { run_window(worker_meta, worker_config, worker_publisher, worker_cancel).await; }); workers.insert(meta.id, Worker { cancel, task }); } WindowEvent::Remove(id) => stop_worker(id, publisher, workers).await, }}
async fn stop_worker(id: u64, publisher: &Publisher, workers: &mut HashMap<u64, Worker>) { if let Some(worker) = workers.remove(&id) { worker.cancel.cancel(); let _ = worker.task.await; } publisher.remove(id);}
async fn run_window( meta: WindowMeta, config: RuntimeConfig, publisher: Publisher, cancel: CancellationToken,) { let id = meta.id; let mut backoff = Duration::from_millis(250); loop { if cancel.is_cancelled() { publisher.set_status(id, CaptureStatus::Stopped, None); return; } publisher.set_status(id, CaptureStatus::Starting, None); match run_capture_attempt(id, &config, &publisher, &cancel).await { Ok(()) if cancel.is_cancelled() => { publisher.set_status(id, CaptureStatus::Stopped, None); return; } Ok(()) => { publisher.set_status(id, CaptureStatus::Error, Some("capture ended".into())); } Err(error) => { warn!(window_id = id, ?error, "window capture failed"); publisher.set_status(id, CaptureStatus::Error, Some(error.to_string())); } } if cancel.is_cancelled() { publisher.set_status(id, CaptureStatus::Stopped, None); return; } publisher.set_status(id, CaptureStatus::Backoff, None); tokio::select! { _ = cancel.cancelled() => { publisher.set_status(id, CaptureStatus::Stopped, None); return; } _ = tokio::time::sleep(backoff) => {} } backoff = (backoff * 2).min(Duration::from_secs(5)); }}
async fn run_capture_attempt( id: u64, config: &RuntimeConfig, publisher: &Publisher, cancel: &CancellationToken,) -> Result<()> { let connection = Connection::session() .await .map_err(|error| anyhow!("connecting to session D-Bus: {error}"))?; let cast = CastSession::start(&connection, id).await?; let (frame_tx, mut frame_rx) = tokio::sync::watch::channel(None); let stop = Arc::new(AtomicBool::new(false)); let capture_thread = capture::start(cast.ready.clone(), config.clone(), frame_tx, stop.clone()); let mut capture_done = tokio::task::spawn_blocking(move || { capture_thread .join() .map_err(|_| anyhow!("capture thread panicked"))? }); let mut got_frame = false; let mut capture_finished = false; let result = loop { tokio::select! { biased; result = &mut capture_done => { capture_finished = true; break result.map_err(|error| anyhow!("capture task join failed: {error}"))?; } _ = cancel.cancelled() => break Ok(()), changed = frame_rx.changed() => { if let Err(error) = changed { capture_finished = true; let capture_result = (&mut capture_done) .await .map_err(|join_error| anyhow!("capture task join failed: {join_error}"))?; break capture_result .map_err(|capture_error| anyhow!("capture frame channel closed: {error}; capture task: {capture_error}")); } if let Some(frame) = frame_rx.borrow_and_update().clone() { if let Err(error) = publisher.publish(frame) { break Err(error); } if !got_frame { got_frame = true; // A successful frame proves the cast and conversion path recovered. publisher.set_status(id, CaptureStatus::Streaming, None); } } } } }; stop.store(true, Ordering::Release); // Stop the compositor stream before joining the capture thread. Closing PipeWire // unblocks a backend that is waiting for the next buffer during shutdown. cast.stop().await; if !capture_finished { let _ = capture_done.await; } result}