From 5baeb8feb5c849ebf4f97de5d98e6fff288831de Mon Sep 17 00:00:00 2001 From: Mariano Abad Date: Tue, 21 Jul 2026 15:38:57 -0300 Subject: [PATCH] drm: bound the _drm body read, stream-scope cursor teardown, refresh a stale verdict, drop dead clear (review 5) - recv_msg_timeout2 only gated the wait for the first byte, so a peer that sent one byte then stalled pinned the task forever. The same budget now also bounds the body read; a body that overruns is a hard error that tears the stream down (recv_msg bodies are small JSON, so a healthy peer never trips it). - The cursor cache is keyed by display index, which a rebuilt stream reuses, so a predecessor exiting after its replacement published a fresh cursor erased it. Stamp each entry with a monotonic per-stream epoch and compare-and-remove on teardown. - ProbeState::Available had no TTL, so an idle hotplug left a phantom display in enumeration. Give it a timestamp and refresh the list off the hot path once it ages past POSITIVE_TTL. The verdict stays true across the refresh (never bounces a live session to the portal) and the probe runs on a background thread (never blocks the async enumeration). - Remove the dead clear(): it is unreferenced, and wiring it into teardown would force the blocking re-probe on the next enumeration that swap_available_displays exists to avoid. --- src/ipc.rs | 20 +++++-- src/server/drm_capturer.rs | 114 ++++++++++++++++++++++++++++--------- 2 files changed, 102 insertions(+), 32 deletions(-) diff --git a/src/ipc.rs b/src/ipc.rs index f5d529dc8..a0cc4e9fe 100644 --- a/src/ipc.rs +++ b/src/ipc.rs @@ -2753,9 +2753,14 @@ impl DrmConn { } /// Cancel-safe timeout wrapper around `recv_msg`, mirroring `ConnectionTmpl::next_timeout2`, so a - /// dropped consumer re-checks its `stop` flag between frames. `None` on timeout. The timeout gates - /// ONLY the wait for the first byte (`readable()` consumes nothing), so a fired timeout leaves the - /// stream at a clean frame boundary and never strands a partial frame or its fd. + /// dropped consumer re-checks its `stop` flag between frames. `None` when no frame started within + /// the window (a clean boundary: `readable()` consumes nothing, so re-polling is safe). Once a byte + /// is available the frame has started, so the SAME budget also bounds the body read: a peer that + /// sends one byte then stalls cannot pin this task forever (the `readable()` gate alone does not + /// cover the length prefix or payload). A body that overruns the budget is a hard error, not a + /// `None`, because the frame is partially consumed and cannot be safely resumed -- the caller tears + /// the stream down. `recv_msg` bodies are small length-prefixed JSON (<= MAX_DRM_JSON_BYTES), well + /// under any caller's budget over a local socket, so this never trips a healthy peer. pub async fn recv_msg_timeout2( &mut self, ms_timeout: u64, @@ -2764,9 +2769,14 @@ impl DrmConn { // at the `;` (releasing `&self.stream`) BEFORE `recv_msg()` takes `&mut self` in an arm. let ready = timeout(ms_timeout, self.stream.readable()).await; match ready { - Err(_) => None, // timed out at a frame boundary + Err(_) => None, // no frame started: clean boundary, caller re-checks `stop` Ok(Err(e)) => Some(Err(e.into())), - Ok(Ok(())) => Some(self.recv_msg().await), + Ok(Ok(())) => match timeout(ms_timeout, self.recv_msg()).await { + Ok(res) => Some(res), + Err(_) => Some(Err(anyhow::anyhow!( + "drm: frame body stalled past {ms_timeout}ms after first byte; closing" + ))), + }, } } diff --git a/src/server/drm_capturer.rs b/src/server/drm_capturer.rs index bf405a7f4..90a28440d 100644 --- a/src/server/drm_capturer.rs +++ b/src/server/drm_capturer.rs @@ -226,6 +226,9 @@ async fn recv_thread( stop: Arc, tx: std::sync::mpsc::Sender>>, ) { + // Unique tag for this stream's cursor entries so teardown only erases its own (see + // remove_drm_cursor); a rebuilt stream for the same display index gets a newer epoch. + let cursor_epoch = next_cursor_epoch(); // Handshake: connect, receive the display list, request the display. let mut conn = match connect_drm(1000).await { Ok(c) => c, @@ -393,6 +396,7 @@ async fn recv_thread( } set_drm_cursor( display, + cursor_epoch, DrmCursorData { id, width: width as i32, @@ -425,8 +429,8 @@ async fn recv_thread( // `IpcDrmCapturer::Drop` (which runs on the encoder thread). drop(converter); // Drop only THIS stream's cursor entry so a torn-down monitor does not erase the cursor state of - // other still-active streams. - remove_drm_cursor(display); + // other still-active streams, nor a replacement stream that already re-took this display index. + remove_drm_cursor(display, cursor_epoch); let mut slot = shared.slot.lock().unwrap(); slot.ended = Some(format!("drm stream ended ({end_reason})")); shared.cv.notify_one(); @@ -448,14 +452,28 @@ pub struct DrmCursorData { pub colors: Vec, } -static DRM_CURSOR: Mutex> = Mutex::new(BTreeMap::new()); +static DRM_CURSOR: Mutex> = Mutex::new(BTreeMap::new()); +// Monotonic per-stream tag. A display index can be served by successive recv_threads (a rebuilt +// stream reuses the index), so each cursor entry is stamped with the writing stream's epoch and a +// torn-down stream drops its entry ONLY if the epoch still matches. Without this, a predecessor that +// exits just after its replacement published a fresh cursor for the same index would erase it. +static DRM_CURSOR_EPOCH: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1); -fn set_drm_cursor(display: i32, c: DrmCursorData) { - DRM_CURSOR.lock().unwrap().insert(display, c); +fn next_cursor_epoch() -> u64 { + DRM_CURSOR_EPOCH.fetch_add(1, std::sync::atomic::Ordering::Relaxed) } -fn remove_drm_cursor(display: i32) { - DRM_CURSOR.lock().unwrap().remove(&display); +fn set_drm_cursor(display: i32, epoch: u64, c: DrmCursorData) { + DRM_CURSOR.lock().unwrap().insert(display, (epoch, c)); +} + +// Compare-and-remove: drop the entry only if THIS stream (epoch) still owns it. A replacement stream +// for the same index holds a newer epoch, so a late-exiting predecessor leaves the fresh cursor intact. +fn remove_drm_cursor(display: i32, epoch: u64) { + let mut map = DRM_CURSOR.lock().unwrap(); + if map.get(&display).map(|(e, _)| *e) == Some(epoch) { + map.remove(&display); + } } // Pick the cursor to present: prefer the visible one (the pointer is over exactly one captured CRTC @@ -464,8 +482,9 @@ fn remove_drm_cursor(display: i32) { fn pick_drm_cursor() -> Option { let map = DRM_CURSOR.lock().unwrap(); map.values() + .map(|(_, c)| c) .find(|c| c.id != scrap::drm_reader::HIDDEN_CURSOR_ID) - .or_else(|| map.values().next()) + .or_else(|| map.values().map(|(_, c)| c).next()) .cloned() } @@ -497,12 +516,19 @@ enum ProbeState { // is_available): displays that appear after startup (a headless boot settling, a monitor // hotplug, or a --service restart) can then re-enable it without restarting the --server. Unavailable(Instant), - Available(Vec), + // Timestamped with the instant the list was probed. A positive verdict is sticky (DRM stays + // selected), but once it ages past POSITIVE_TTL is_available refreshes the list off the hot path + // so an idle hotplug does not leave a phantom display in enumeration. The timestamp lives in the + // variant so every site that publishes an Available list is forced to stamp it. + Available(Instant, Vec), } static DRM_STATE: Mutex = Mutex::new(ProbeState::Unknown); // How long a negative availability verdict is trusted before is_available re-probes. const NEGATIVE_TTL: Duration = Duration::from_secs(30); +// How long a positive verdict is served before is_available kicks a background list refresh. The +// verdict stays true across the refresh; only the cached display list is renewed. +const POSITIVE_TTL: Duration = Duration::from_secs(15); /// Query the service for the current DRM display list without starting a stream: connect, read the /// list the service sends on connect, then drop the connection (the service closes it when we do @@ -547,7 +573,7 @@ pub(super) fn is_available() -> bool { // no displays at probe time can still enable DRM once displays appear (without a --server // restart). NEVER call the blocking probe while holding DRM_STATE: a cold or expired probe would // otherwise serialize every async caller for the whole query_displays() timeout (~4s). - { + let verdict = { let mut st = DRM_STATE.lock().unwrap(); if let ProbeState::Unavailable(since) = &*st { if since.elapsed() >= NEGATIVE_TTL { @@ -556,17 +582,28 @@ pub(super) fn is_available() -> bool { } } match &*st { - ProbeState::Available(_) => return true, - ProbeState::Unavailable(_) => return false, - ProbeState::Unknown => {} // fall through and probe with the lock released + // Sticky positive; refresh the list off the hot path if it has gone stale (see below). + ProbeState::Available(since, _) => Some((true, since.elapsed() >= POSITIVE_TTL)), + ProbeState::Unavailable(_) => Some((false, false)), + ProbeState::Unknown => None, // fall through and probe with the lock released } + }; + if let Some((available, stale)) = verdict { + // A stale positive verdict re-probes on a background thread and returns the still-valid + // verdict immediately: never block the (often async) caller for the probe timeout, and never + // flip a live positive to false, which would bounce a capturing session to the portal. The + // refresh only replaces the list when the probe returns a non-empty set. + if stale { + refresh_available_async(); + } + return available; } // Single-flight: exactly one caller probes at a time. While a probe is in flight, others return // the current cache-only verdict instead of stacking redundant `_drm` probes or blocking on the // mutex across the I/O. warm_availability normally seeds `Available` before clients connect, so // this cold path is rare. if DRM_PROBE_IN_FLIGHT.swap(true, Ordering::AcqRel) { - return matches!(&*DRM_STATE.lock().unwrap(), ProbeState::Available(_)); + return matches!(&*DRM_STATE.lock().unwrap(), ProbeState::Available(..)); } let t = Instant::now(); let result = query_displays(); @@ -578,7 +615,7 @@ pub(super) fn is_available() -> bool { list.len(), t.elapsed() ); - *st = ProbeState::Available(list); + *st = ProbeState::Available(Instant::now(), list); true } Ok(_) => { @@ -605,6 +642,34 @@ pub(super) fn is_available() -> bool { available } +/// Refresh a stale positive verdict off the hot path: probe on a background thread and, on a non-empty +/// result, renew the cached display list so enumeration stops advertising a display an idle hotplug +/// removed. Single-flight via the same guard as the cold probe so a refresh never stacks with a probe +/// or another refresh. It never demotes an `Available` verdict to `Unavailable` (a transient empty or +/// failed probe must not disable a working DRM session), and it re-stamps the verdict on any completed +/// probe so we re-probe at most once per POSITIVE_TTL even when the list is unchanged or the probe +/// fails. +fn refresh_available_async() { + if DRM_PROBE_IN_FLIGHT.swap(true, Ordering::AcqRel) { + return; + } + std::thread::spawn(|| { + let result = query_displays(); + { + let mut st = DRM_STATE.lock().unwrap(); + if let ProbeState::Available(since, list) = &mut *st { + *since = Instant::now(); + if let Ok(fresh) = result { + if !fresh.is_empty() { + *list = fresh; + } + } + } + } + DRM_PROBE_IN_FLIGHT.store(false, Ordering::Release); + }); +} + /// Warm the availability cache at `--server` startup so the first client connection does not race a /// cold `_drm` probe. A cold probe blocks display enumeration, and if it has not settled when the /// peer info is built the display list goes out empty and the client shows "No displays" and @@ -612,13 +677,13 @@ pub(super) fn is_available() -> bool { /// the positive result; a genuinely DRM-less host just falls through to the lazy `is_available()`. pub(super) fn warm_availability() { for _ in 0..10 { - if matches!(&*DRM_STATE.lock().unwrap(), ProbeState::Available(_)) { + if matches!(&*DRM_STATE.lock().unwrap(), ProbeState::Available(..)) { return; } match query_displays() { Ok(list) if !list.is_empty() => { log::info!("drm: consumer cache warmed ({} displays) at startup", list.len()); - *DRM_STATE.lock().unwrap() = ProbeState::Available(list); + *DRM_STATE.lock().unwrap() = ProbeState::Available(Instant::now(), list); return; } // Producer not ready yet (or no DRM): back off and retry; never cache a negative here. @@ -632,7 +697,7 @@ pub(super) fn warm_availability() { /// (per-monitor position + scale). `None` until probed/available. pub(super) fn get_display_infos() -> Option> { let list = match &*DRM_STATE.lock().unwrap() { - ProbeState::Available(list) => list.clone(), + ProbeState::Available(_, list) => list.clone(), _ => return None, }; let multi = list.len() > 1; @@ -665,7 +730,7 @@ pub(super) fn get_display_infos() -> Option> { /// display whenever the primary is not connector 0. pub(super) fn get_primary_index() -> usize { let list = match &*DRM_STATE.lock().unwrap() { - ProbeState::Available(list) => list.clone(), + ProbeState::Available(_, list) => list.clone(), _ => return 0, }; let wl = scrap::wayland::display::get_displays(); @@ -754,11 +819,6 @@ fn normalize_connector(name: &str) -> String { } } -/// Reset the probe cache so the next session re-probes (called on capture teardown). -pub(super) fn clear() { - *DRM_STATE.lock().unwrap() = ProbeState::Unknown; -} - /// Swap the sticky positive availability cache to a freshly-enumerated display list, driven by a /// service-pushed `DrmDisplaysChanged` hotplug signal on a live stream. This is the off-hot-path cache /// refresh that keeps mid-session hotplug geometry fresh WITHOUT the blocking `_drm` re-probe that @@ -768,9 +828,9 @@ pub(super) fn clear() { /// force DRM on; establishing availability stays the job of the probe path. fn swap_available_displays(list: Vec) { let mut st = DRM_STATE.lock().unwrap(); - if matches!(&*st, ProbeState::Available(_)) { + if matches!(&*st, ProbeState::Available(..)) { log::info!("drm: hotplug refresh -> {} display(s)", list.len()); - *st = ProbeState::Available(list); + *st = ProbeState::Available(Instant::now(), list); } } @@ -859,7 +919,7 @@ pub(super) fn get_capturer_info( .get(display_idx) .map(|di| (di.x, di.y)) .unwrap_or((d.x, d.y)); - *DRM_STATE.lock().unwrap() = ProbeState::Available(displays); + *DRM_STATE.lock().unwrap() = ProbeState::Available(Instant::now(), displays); Ok(super::video_service::CapturerInfo { origin, width: d.width as usize,