diff --git a/crates/rds-desktop/src/input/x11.rs b/crates/rds-desktop/src/input/x11.rs index d7c68ff..9b36d6d 100644 --- a/crates/rds-desktop/src/input/x11.rs +++ b/crates/rds-desktop/src/input/x11.rs @@ -4,8 +4,10 @@ //! live in `input/portal` and `input/wlr`. use std::collections::{BTreeMap, BTreeSet}; +use std::sync::Arc; mod keymap; +mod repeat; use keymap::KeyMap; use x11rb::protocol::xkb::ConnectionExt as _; @@ -17,8 +19,6 @@ use x11rb::rust_connection::RustConnection; use crate::{DesktopError, InputSink}; -const KEY_PRESS: u8 = 2; -const KEY_RELEASE: u8 = 3; const BUTTON_PRESS: u8 = 4; const BUTTON_RELEASE: u8 = 5; const MOTION_NOTIFY: u8 = 6; @@ -27,13 +27,14 @@ const MAX_SCROLL_CLICKS: f64 = 32.0; /// XTEST input bound to one X screen, resolving evdev through actual XKB key names. /// Extended evdev keys outside core X11's 8-bit keycodes are refused. pub struct XtestInput { - conn: RustConnection, + conn: Arc, root: x11rb::protocol::xproto::Window, screen: u32, width: u16, height: u16, keymap: KeyMap, - keys: BTreeMap, + server: String, + keys: BTreeMap, buttons: BTreeSet, scroll_x: f64, scroll_y: f64, @@ -77,17 +78,26 @@ impl XtestInput { )); } let keymap = load_keymap(&conn)?; + let display = + std::env::var("DISPLAY").map_err(|_| error("X11 display name is unavailable"))?; + // X screens share the same core keyboard. Normalize the screen suffix + // so parallel controllers retain one per-key repeat/hold reference. + let server = display + .rsplit_once('.') + .filter(|(_, suffix)| suffix.parse::().is_ok()) + .map_or(display.clone(), |(server, _)| server.to_owned()); tracing::info!( mapped_keys = keymap.len(), "X11 physical keyboard map ready" ); Ok(Self { - conn, + conn: Arc::new(conn), root, screen, width, height, keymap, + server, keys: BTreeMap::new(), buttons: BTreeSet::new(), scroll_x: 0.0, @@ -160,8 +170,8 @@ impl XtestInput { self.keymap = load_keymap(&self.conn)?; } } - let key = if let Some(key) = self.keys.get(&code) { - *key + let key = if let Some(hold) = self.keys.get(&code) { + hold.key } else { let key = self.keymap.resolve(code)?; if !pressed { @@ -171,12 +181,18 @@ impl XtestInput { }; if pressed { self.keyboard_on_screen()?; - } - self.fake(if pressed { KEY_PRESS } else { KEY_RELEASE }, key, 0, 0)?; - if pressed { - self.keys.insert(code, key); - } else { - self.keys.remove(&code); + let modifier = matches!( + code, + 29 | 42 | 54 | 56 | 97 | 100 | 125 | 126 | 58 | 69 | 70 + ); + if let Some(hold) = self.keys.get(&code) { + return hold.repeat(modifier); + } + let hold = + repeat::Hold::press(self.conn.clone(), self.root, &self.server, key, modifier)?; + self.keys.insert(code, hold); + } else if let Some(mut hold) = self.keys.remove(&code) { + hold.release()?; } Ok(()) } @@ -318,11 +334,7 @@ impl Drop for XtestInput { fn drop(&mut self) { // Best-effort release of only this sink's injected holds. The worker // drops the sink after its last in-flight call, outside Tokio workers. - for key in self.keys.values() { - let _ = self - .conn - .xtest_fake_input(KEY_RELEASE, *key, 0, self.root, 0, 0, 0); - } + self.keys.clear(); for button in &self.buttons { let _ = self .conn diff --git a/crates/rds-desktop/src/input/x11/repeat.rs b/crates/rds-desktop/src/input/x11/repeat.rs new file mode 100644 index 0000000..61eacc9 --- /dev/null +++ b/crates/rds-desktop/src/input/x11/repeat.rs @@ -0,0 +1,126 @@ +//! Client-paced repeats with real native holds, shared by keyboard controllers. +use crate::DesktopError; +use std::collections::BTreeMap; +use std::sync::{Arc, Mutex}; +use x11rb::protocol::xproto::{AutoRepeatMode, ChangeKeyboardControlAux, ConnectionExt as _}; +use x11rb::protocol::xtest::ConnectionExt as _; +use x11rb::rust_connection::RustConnection; + +struct State { + owners: u32, + repeat: bool, +} +static HOLDS: Mutex> = Mutex::new(BTreeMap::new()); + +pub(super) struct Hold { + conn: Arc, + root: u32, + server: String, + pub(super) key: u8, + released: bool, +} + +fn error() -> DesktopError { + DesktopError::Input("X11 keyboard repeat state failed".into()) +} +fn mode(conn: &RustConnection, key: u8, on: bool) -> Result<(), DesktopError> { + conn.change_keyboard_control( + &ChangeKeyboardControlAux::new() + .key(u32::from(key)) + .auto_repeat_mode(if on { + AutoRepeatMode::ON + } else { + AutoRepeatMode::OFF + }), + ) + .map_err(|_| error())? + .check() + .map_err(|_| error()) +} +fn fake(conn: &RustConnection, root: u32, key: u8, press: bool) -> Result<(), DesktopError> { + conn.xtest_fake_input(if press { 2 } else { 3 }, key, 0, root, 0, 0, 0) + .map_err(|_| error())? + .check() + .map_err(|_| error()) +} + +impl Hold { + pub(super) fn press( + conn: Arc, + root: u32, + server: &str, + key: u8, + modifier: bool, + ) -> Result { + let mut holds = HOLDS.lock().map_err(|_| error())?; + let id = (server.to_owned(), key); + if let Some(state) = holds.get_mut(&id) { + let next = state.owners.checked_add(1).ok_or_else(error)?; + if !modifier { + fake(&conn, root, key, false)?; + fake(&conn, root, key, true)?; + } + state.owners = next; + } else { + let keyboard = conn + .get_keyboard_control() + .map_err(|_| error())? + .reply() + .map_err(|_| error())?; + let repeat = keyboard.auto_repeats[usize::from(key) / 8] & (1 << (key % 8)) != 0; + mode(&conn, key, false)?; + if let Err(e) = fake(&conn, root, key, true) { + let _ = mode(&conn, key, repeat); + return Err(e); + } + holds.insert(id, State { owners: 1, repeat }); + } + Ok(Self { + conn, + root, + server: server.to_owned(), + key, + released: false, + }) + } + + pub(super) fn repeat(&self, modifier: bool) -> Result<(), DesktopError> { + if modifier { + return Ok(()); + } + // A repeated client KeyDown is an intentional native repeat. Core X + // suppresses duplicate presses with typematic disabled, so pulse once. + let _holds = HOLDS.lock().map_err(|_| error())?; + fake(&self.conn, self.root, self.key, false)?; + fake(&self.conn, self.root, self.key, true) + } + + pub(super) fn release(&mut self) -> Result<(), DesktopError> { + if self.released { + return Ok(()); + } + let mut holds = HOLDS.lock().map_err(|_| error())?; + let id = (self.server.clone(), self.key); + let Some(state) = holds.get_mut(&id) else { + self.released = true; + return Ok(()); + }; + if state.owners > 1 { + state.owners -= 1; + self.released = true; + return Ok(()); + } + let release = fake(&self.conn, self.root, self.key, false); + let restore = mode(&self.conn, self.key, state.repeat); + holds.remove(&id); + self.released = true; + release.and(restore) + } +} +impl Drop for Hold { + fn drop(&mut self) { + if self.release().is_err() { + tracing::warn!("X11 keyboard hold cleanup failed"); + } + } +} diff --git a/crates/rds-desktop/src/lib.rs b/crates/rds-desktop/src/lib.rs index b83ec86..2d9e792 100644 --- a/crates/rds-desktop/src/lib.rs +++ b/crates/rds-desktop/src/lib.rs @@ -28,6 +28,7 @@ mod decode_work; mod delivery_rate; pub mod input; pub mod mailbox; +mod media_repair; mod order; pub mod render; #[cfg(any(all(target_os = "linux", feature = "x11"), test))] diff --git a/crates/rds-desktop/src/media_repair.rs b/crates/rds-desktop/src/media_repair.rs new file mode 100644 index 0000000..f09343d --- /dev/null +++ b/crates/rds-desktop/src/media_repair.rs @@ -0,0 +1,198 @@ +//! Session-local repair epochs and ownership of the independent-picture gate. +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use tokio::sync::watch; + +const ACTIVE: u64 = 1; +const KEY_PENDING: u64 = 2; +const FLAGS: u64 = ACTIVE | KEY_PENDING; +const NEXT_EPOCH: u64 = FLAGS + 1; + +pub(crate) struct MediaRepair { + // Epoch and gate flags change atomically: an old producer/receipt must + // never install or release the gate belonging to a replacement epoch. + state: AtomicU64, + wake: watch::Sender<()>, + pub(crate) requests: AtomicU64, + pub(crate) coalesced: AtomicU64, +} + +#[derive(Debug, PartialEq, Eq)] +pub(crate) enum Request { + Accepted, + Coalesced, + Exhausted, +} + +impl MediaRepair { + pub(crate) fn new() -> (Arc, watch::Receiver<()>) { + let (wake, changes) = watch::channel(()); + ( + Arc::new(Self { + state: AtomicU64::new(0), + wake, + requests: AtomicU64::new(0), + coalesced: AtomicU64::new(0), + }), + changes, + ) + } + + pub(crate) fn generation(&self) -> u64 { + self.state.load(Ordering::Acquire) & !FLAGS + } + + pub(crate) fn key_pending(&self) -> bool { + self.state.load(Ordering::Acquire) & KEY_PENDING != 0 + } + + pub(crate) fn active(&self) -> bool { + self.state.load(Ordering::Acquire) & ACTIVE != 0 + } + + pub(crate) fn request(&self) -> Request { + match self + .state + .fetch_update(Ordering::AcqRel, Ordering::Acquire, |state| { + if state & ACTIVE != 0 { + return None; + } + (state & !FLAGS) + .checked_add(NEXT_EPOCH) + .map(|epoch| epoch | ACTIVE) + }) { + Ok(_) => { + self.requests.fetch_add(1, Ordering::Relaxed); + self.wake.send_replace(()); + Request::Accepted + } + Err(state) if state & ACTIVE != 0 => { + self.coalesced.fetch_add(1, Ordering::Relaxed); + Request::Coalesced + } + Err(_) => Request::Exhausted, + } + } + + pub(crate) async fn changed( + &self, + changes: &mut watch::Receiver<()>, + observed: u64, + ) -> Result<(), watch::error::RecvError> { + loop { + if self.generation() != observed { + return Ok(()); + } + // A sender may publish its wake after the writer already noticed + // the atomic epoch. Retire that wake without canceling fresh work. + changes.changed().await?; + } + } + + pub(crate) fn key( + self: &Arc, + generation: u64, + idr: Arc, + ) -> Option { + self.state + .fetch_update(Ordering::AcqRel, Ordering::Acquire, |state| { + (state & !FLAGS == generation && state & KEY_PENDING == 0) + .then_some(state | KEY_PENDING) + }) + .ok()?; + Some(KeyPermit { + repair: self.clone(), + generation, + idr, + confirmed: false, + }) + } + + fn release_key(&self, generation: u64, confirmed: bool) { + let _ = self + .state + .fetch_update(Ordering::AcqRel, Ordering::Acquire, |state| { + if state & !FLAGS != generation { + return None; + } + Some(state & !(KEY_PENDING | if confirmed { ACTIVE } else { 0 })) + }); + } +} + +pub(crate) struct KeyPermit { + repair: Arc, + generation: u64, + idr: Arc, + pub(crate) confirmed: bool, +} + +impl Drop for KeyPermit { + fn drop(&mut self) { + if self.repair.generation() == self.generation { + if !self.confirmed { + // A fresh key taken from the queue can be dropped with a + // canceled write. Rebuild it before admitting any successor. + self.idr.store(true, Ordering::Release); + } + self.repair.release_key(self.generation, self.confirmed); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn old_key_cannot_install_or_release_the_replacement_gate() { + let (repair, _) = MediaRepair::new(); + let idr = Arc::new(AtomicBool::new(false)); + let mut old = repair.key(0, idr.clone()).unwrap(); + assert_eq!(repair.request(), Request::Accepted); + let epoch = repair.generation(); + assert!(repair.key(0, idr.clone()).is_none()); + let mut fresh = repair.key(epoch, idr.clone()).unwrap(); + old.confirmed = true; + drop(old); + assert!(repair.key_pending()); + assert!(repair.active()); + assert!(!idr.load(Ordering::Acquire)); + fresh.confirmed = true; + drop(fresh); + assert!(!repair.key_pending()); + assert!(!repair.active()); + } + + #[test] + fn canceled_fresh_key_releases_credit_but_keeps_repair_active() { + let (repair, _) = MediaRepair::new(); + let idr = Arc::new(AtomicBool::new(false)); + assert_eq!(repair.request(), Request::Accepted); + drop(repair.key(repair.generation(), idr.clone()).unwrap()); + assert!(!repair.key_pending()); + assert!(repair.active()); + assert!(idr.load(Ordering::Acquire)); + assert_eq!(repair.request(), Request::Coalesced); + assert!(repair.key(repair.generation(), idr).is_some()); + } + + #[test] + fn repeated_requests_preserve_an_active_recovery_picture() { + let (repair, _) = MediaRepair::new(); + let idr = Arc::new(AtomicBool::new(false)); + assert_eq!(repair.request(), Request::Accepted); + let epoch = repair.generation(); + let mut key = repair.key(epoch, idr.clone()).unwrap(); + for _ in 0..100 { + assert_eq!(repair.request(), Request::Coalesced); + assert_eq!(repair.generation(), epoch); + } + assert!(repair.key_pending()); + assert!(!idr.load(Ordering::Acquire)); + key.confirmed = true; + drop(key); + assert_eq!(repair.request(), Request::Accepted); + assert!(repair.generation() > epoch); + } +} diff --git a/crates/rds-desktop/src/render/viewer.rs b/crates/rds-desktop/src/render/viewer.rs index 2c589c6..5717891 100644 --- a/crates/rds-desktop/src/render/viewer.rs +++ b/crates/rds-desktop/src/render/viewer.rs @@ -700,10 +700,24 @@ impl ApplicationHandler<()> for App { self.modifiers = modifiers; self.sync_modifiers(event_loop); } - WindowEvent::KeyboardInput { event, .. } if !event.repeat => { + WindowEvent::KeyboardInput { event, .. } => { if let PhysicalKey::Code(key) = event.physical_key && let Some(code) = evdev(key) { + if event.repeat { + // Repeat only a key actually forwarded and still held. + // Never re-run clipboard/local shortcut side effects. + if event.state == ElementState::Pressed + && self.keys.contains(&code) + && !matches!( + code, + 29 | 42 | 54 | 56 | 97 | 100 | 125 | 126 | 58 | 69 | 70 + ) + { + self.input(event_loop, InputKind::KeyDown { code }); + } + return; + } if !matches!(code, 29 | 42 | 54 | 56 | 97 | 100 | 125 | 126) { self.sync_modifiers(event_loop); } diff --git a/crates/rds-desktop/src/session.rs b/crates/rds-desktop/src/session.rs index 0e9c21d..36f8079 100644 --- a/crates/rds-desktop/src/session.rs +++ b/crates/rds-desktop/src/session.rs @@ -161,7 +161,7 @@ impl DeliveryPressure { #[derive(Clone)] struct FrameDelivery { - keyframe_pending: Arc, + repair: Arc, idr: Arc, feedback: Arc, latest_key_seq: Arc, @@ -225,6 +225,8 @@ impl Drop for CapturePermit { } struct AdmittedFrame { produced: Produced, + generation: u64, + key: Option, permit: CapturePermit, } impl Borrow for AdmittedFrame { @@ -632,7 +634,7 @@ pub async fn serve_desktop_with( // Capture+encode runs on a blocking thread; frames flow to the writer. let (tx, mut rx) = mpsc::channel::(2); - let keyframe_pending = Arc::new(AtomicBool::new(false)); + let (media_repair, mut media_changes) = crate::media_repair::MediaRepair::new(); let latest_key_seq = Arc::new(AtomicU64::new(u64::MAX)); let capture_admission = Arc::new(AtomicU64::new(0)); let delivery_feedback = Arc::new(DeliveryFeedback::default()); @@ -645,7 +647,7 @@ pub async fn serve_desktop_with( let input_refresh_until_ms = Arc::clone(&controls.input_refresh_until_ms); let input_refresh_pending = Arc::clone(&controls.input_refresh_pending); let mut producer = config.producer; - let keyframe_pending = keyframe_pending.clone(); + let repair = media_repair.clone(); let latest_key_seq = latest_key_seq.clone(); let capture_admission = capture_admission.clone(); let feedback = delivery_feedback.clone(); @@ -666,6 +668,7 @@ pub async fn serve_desktop_with( None => platform_producer(hello.display, frame_interval, config.output_height), }; let mut seq = 0u64; + let mut generation = 0; loop { // The channel is the session's lifecycle: when the // session ends the writer drops `rx` and this loop exits @@ -679,7 +682,7 @@ pub async fn serve_desktop_with( // force repeated expensive IDRs while QUIC was backlogged. let mut paused = false; let permit = loop { - while tx.capacity() == 0 || keyframe_pending.load(Ordering::Acquire) { + while tx.capacity() == 0 || repair.key_pending() { paused = true; if tx.is_closed() { return; @@ -703,11 +706,21 @@ pub async fn serve_desktop_with( if paused { source.resume_after_backpressure(); } + let frame_generation = repair.generation(); + if generation != frame_generation { + generation = frame_generation; + producer_controls.idr.store(true, Ordering::Release); + } feedback.producing.store(true, Ordering::Relaxed); let result = source.produce(seq, &producer_controls, &clock); feedback.producing.store(false, Ordering::Relaxed); match result { Some(p) => { + // A repair can race a blocking encode. Its mutated + // references are discarded; the next epoch forces IDR. + if repair.generation() != frame_generation { + continue; + } if p.payload.is_empty() && source.preserves_reference() { feedback.codec_skips.fetch_add(1, Ordering::Relaxed); continue; @@ -718,16 +731,25 @@ pub async fn serve_desktop_with( .store(clock.now_ms(), Ordering::Relaxed); feedback.produced.fetch_add(1, Ordering::Relaxed); } - if p.header.keyframe && !p.payload.is_empty() { + let key = if p.header.keyframe && !p.payload.is_empty() { + let Some(key) = + repair.key(frame_generation, producer_controls.idr.clone()) + else { + continue; + }; latest_key_seq.store(p.header.seq, Ordering::Release); - keyframe_pending.store(true, Ordering::Release); - } + Some(key) + } else { + None + }; // Losing any encoded reference breaks its successors, // not only losing an IDR. Keep the two-slot bound and // ask the producer for an independent replacement. if tx .try_send(AdmittedFrame { produced: p, + generation: frame_generation, + key, permit, }) .is_err() @@ -754,7 +776,7 @@ pub async fn serve_desktop_with( let misses = Arc::clone(&controls.deadline_misses); let feedback = delivery_feedback.clone(); let admission = capture_admission.clone(); - let pending_key = keyframe_pending.clone(); + let repair = media_repair.clone(); let progress_clock = clock.clone(); let mut controller = BitrateController::new(4_000_000, ceiling); let mut delivery_rate = crate::delivery_rate::DeliveryRate::default(); @@ -865,7 +887,11 @@ pub async fn serve_desktop_with( codec_skips = feedback.codec_skips.load(Ordering::Relaxed), last_produced_age_ms = ?(produced > 0).then(|| progress_clock.now_ms().saturating_sub(feedback.last_produced_ms.load(Ordering::Relaxed))), pending_media_frames = admission.load(Ordering::Acquire), - keyframe_pending = pending_key.load(Ordering::Acquire), + keyframe_pending = repair.key_pending(), + media_repair_generation = repair.generation(), + media_repair_active = repair.active(), + media_repair_requests = repair.requests.load(Ordering::Relaxed), + media_repair_coalesced = repair.coalesced.load(Ordering::Relaxed), bitrate_bps = bps, delayed_delivery = delayed, timely_delivery_floor_bps = delivery_floor, @@ -909,6 +935,7 @@ pub async fn serve_desktop_with( let writer_idr = Arc::clone(&controls.idr); let writer_feedback = delivery_feedback.clone(); let frame_route = config.frame_route.unwrap_or(rds_core::UniHello::Desktop); + let repair = media_repair.clone(); workers.spawn(async move { // Token bucket on the paced bitrate: offering faster than the // path sustains only backlogs QUIC's send buffer with frames @@ -925,7 +952,31 @@ pub async fn serve_desktop_with( let mut obsolete = 0u64; let mut failed_delivery = 0u64; let mut health = Instant::now(); + let mut generation = 0; 'writer: loop { + let current_generation = repair.generation(); + if current_generation != generation { + let canceled_receipts = acknowledgements.len(); + // Only this desktop writer owns these tasks. Their reset-on- + // drop streams release obsolete media without closing control + // or any other connection service. + acknowledgements.shutdown().await; + if pending + .as_ref() + .is_some_and(|p| p.generation != current_generation) + { + pending = None; + } + chain = FrameChain::default(); + budget = 0.0; + last = Instant::now(); + generation = current_generation; + tracing::info!( + media_repair_generation = generation, + canceled_receipts, + "desktop obsolete media retired for explicit repair" + ); + } drain_frame_receipts( &mut acknowledgements, &mut chain, @@ -935,7 +986,15 @@ pub async fn serve_desktop_with( &mut failed_delivery, ); if acknowledgements.len() >= MAX_PENDING_FRAME_ACKS { - match acknowledgements.join_next().await { + let receipt = tokio::select! { + biased; + changed = repair.changed(&mut media_changes, generation) => { + if changed.is_err() { break 'writer; } + continue 'writer; + } + receipt = acknowledgements.join_next() => receipt, + }; + match receipt { Some(Ok((_, FrameReceipt::Delivered))) => acknowledged += 1, Some(Ok((_, FrameReceipt::Obsolete))) => obsolete += 1, Some(Ok((seq, FrameReceipt::Failed))) => { @@ -952,11 +1011,24 @@ pub async fn serve_desktop_with( } let mut produced = match pending.take() { Some(p) => p, - None => match rx.recv().await { - Some(p) => p, - None => break, - }, + None => { + let next = tokio::select! { + biased; + changed = repair.changed(&mut media_changes, generation) => { + if changed.is_err() { break 'writer; } + continue 'writer; + } + frame = rx.recv() => frame, + }; + match next { + Some(p) => p, + None => break 'writer, + } + } }; + if produced.generation != generation { + continue; + } if produced.payload.is_empty() { chain.next = None; request_frame_repair(&latest_key_seq, &writer_idr, produced.header.seq); @@ -998,7 +1070,14 @@ pub async fn serve_desktop_with( ); } if !wait.is_zero() { - tokio::time::sleep(wait).await; + tokio::select! { + biased; + changed = repair.changed(&mut media_changes, generation) => { + if changed.is_err() { break 'writer; } + continue 'writer; + } + _ = tokio::time::sleep(wait) => {}, + } } // Admission follows the final selection: advancing this before // the pacing wait would lose track of frames collapsed afterward. @@ -1022,21 +1101,21 @@ pub async fn serve_desktop_with( .saturating_sub(produced.header.encode_done_ts_ms), "desktop frame sending" ); - match send_frame( - &writer_conn, - frame_route, - produced, - &mut rx, - &mut acknowledgements, - FrameDelivery { - keyframe_pending: keyframe_pending.clone(), - idr: writer_idr.clone(), - feedback: writer_feedback.clone(), - latest_key_seq: latest_key_seq.clone(), - }, - ) - .await - { + let outcome = tokio::select! { + biased; + changed = repair.changed(&mut media_changes, generation) => { + if changed.is_err() { break 'writer; } + continue 'writer; + } + outcome = send_frame( + &writer_conn, frame_route, produced, &mut rx, &mut acknowledgements, + FrameDelivery { + repair: repair.clone(), idr: writer_idr.clone(), + feedback: writer_feedback.clone(), latest_key_seq: latest_key_seq.clone(), + }, + ) => outcome, + }; + match outcome { SendOutcome::Sent => { sent += 1; } @@ -1065,7 +1144,7 @@ pub async fn serve_desktop_with( failed_delivery, pending_acknowledgements = acknowledgements.len(), pending_media_frames = capture_admission.load(Ordering::Acquire), - keyframe_pending = keyframe_pending.load(Ordering::Acquire), + keyframe_pending = repair.key_pending(), bitrate_bps = writer_bitrate.load(Ordering::Relaxed), delayed_delivery = writer_feedback.delayed.load(Ordering::Relaxed), last_ack_ms = writer_feedback.last_ack_ms.load(Ordering::Relaxed), @@ -1156,9 +1235,21 @@ pub async fn serve_desktop_with( } } } - Ok(DesktopControl::RequestIdr) => { - controls.idr.store(true, Ordering::Relaxed); - } + Ok(DesktopControl::RequestIdr) => match media_repair.request() { + crate::media_repair::Request::Accepted => { + tracing::info!( + media_repair_generation = media_repair.generation(), + "desktop explicit media repair accepted" + ); + } + crate::media_repair::Request::Coalesced => { + tracing::debug!("desktop repair already delivering an independent picture"); + } + crate::media_repair::Request::Exhausted => { + tracing::warn!("desktop repair generation exhausted"); + break; + } + }, Ok(DesktopControl::SetBitrate(bps)) => { let bps = u64::from(bps.max(50_000)).min(ceiling); controls.requested.store(bps, Ordering::Relaxed); @@ -1663,9 +1754,6 @@ fn obsolete_send( ) -> SendOutcome { sending.finished = true; delivery.feedback.obsolete.fetch_add(1, Ordering::Relaxed); - if produced.header.keyframe { - delivery.keyframe_pending.store(false, Ordering::Release); - } tracing::info!( frame_seq = produced.header.seq, payload_bytes = produced.payload.len(), @@ -1733,7 +1821,12 @@ async fn send_frame_inner( delivery: FrameDelivery, ) -> SendOutcome { let transfer_started = Instant::now(); - let AdmittedFrame { produced, permit } = produced; + let AdmittedFrame { + produced, + permit, + generation, + mut key, + } = produced; let mut sending = match conn.open_uni().await { Ok(stream) => FrameSend { stream, @@ -1798,7 +1891,7 @@ async fn send_frame_inner( let payload_bytes = produced.payload.len(); let delay_budget = delivery_delay_budget(conn.current_path_stats()); let FrameDelivery { - keyframe_pending, + repair, idr, feedback, latest_key_seq, @@ -1824,9 +1917,13 @@ async fn send_frame_inner( } }; feedback.last_ack_ms.store(started.elapsed().as_millis() as u64, Ordering::Relaxed); - let acknowledged = match result { + let acknowledged = if repair.generation() != generation { + feedback.obsolete.fetch_add(1, Ordering::Relaxed); + FrameReceipt::Obsolete + } else { match result { Ok(Ok(None)) => { sending.finished = true; + if let Some(key) = key.as_mut() { key.confirmed = true; } feedback.acknowledged(payload_bytes, transfer_started.elapsed(), delay_budget); if late.is_some() { // Complete the soft-delay record at ordinary diagnostic @@ -1856,15 +1953,13 @@ async fn send_frame_inner( tracing::warn!(frame_seq=seq,keyframe,payload_bytes,ack_ms=started.elapsed().as_millis(),outcome=?result,"desktop frame delivery unconfirmed"); FrameReceipt::Failed } - }; + }}; drop(late); // Capture the complete reset owner, including its Drop implementation, // rather than allowing disjoint field captures in the async closure. drop(sending); + drop(key); drop(permit); - if keyframe { - keyframe_pending.store(false, Ordering::Release); - } (seq, acknowledged) }); outcome @@ -2416,12 +2511,14 @@ mod tests { }, payload: Bytes::from(vec![1; 32 * 1024 * 1024]), }, + generation: 0, + key: None, permit: CapturePermit(admission.clone()), }; let (_tx, mut rx) = mpsc::channel(2); let mut receipts = JoinSet::new(); let delivery = FrameDelivery { - keyframe_pending: Arc::new(AtomicBool::new(false)), + repair: crate::media_repair::MediaRepair::new().0, idr: idr.clone(), feedback: feedback.clone(), latest_key_seq: Arc::new(AtomicU64::new(u64::MAX)), diff --git a/crates/rds-desktop/tests/media_repair.rs b/crates/rds-desktop/tests/media_repair.rs new file mode 100644 index 0000000..d340c78 --- /dev/null +++ b/crates/rds-desktop/tests/media_repair.rs @@ -0,0 +1,338 @@ +//! Explicit media repair must bypass an obsolete key's delivery barrier. +//! Synthetic transport/input only: never touches a host display or input seat. +use bytes::Bytes; +use rds_core::{ + Codec, DesktopControl, DesktopEvent, DesktopHello, FrameHeader, InputEvent, InputKind, +}; +use rds_desktop::{ + Decoder, DesktopError, EncodedFrame, Encoder, FrameProducer, InputSink, Produced, + ProducerControls, RawFrame, SessionClock, SessionConfig, serve_desktop_with, +}; +use rds_net::{Backend, EndpointConfig, read_frame, write_frame}; +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; + +struct BlockedKey { + stop: std::sync::mpsc::Receiver<()>, + deltas: bool, + produced: Arc, + old_bytes: usize, + encoder: Option>, +} + +impl FrameProducer for BlockedKey { + fn produce( + &mut self, + seq: u64, + controls: &ProducerControls, + clock: &SessionClock, + ) -> Option { + let blocked_end = if self.deltas { 3 } else { 0 }; + if seq > blocked_end + 1 { + let _ = self.stop.recv_timeout(Duration::from_secs(3)); + return None; + } + self.produced.fetch_add(1, Ordering::Release); + let requested = controls.idr.swap(false, Ordering::AcqRel); + let (keyframe, payload, side) = if let Some(encoder) = self.encoder.as_mut() { + if requested { + encoder.request_idr(); + } + let shade = (50 + seq * 43 % 150) as u8; + let raw = RawFrame { + width: 64, + height: 64, + stride: 256, + data: Bytes::from([shade, shade, shade, 255].repeat(64 * 64)), + }; + let encoded = encoder.encode(&raw).unwrap(); + (encoded.keyframe, encoded.data, 64) + } else { + (requested, Bytes::from_static(&[7; 128]), 16) + }; + Some(Produced { + header: FrameHeader { + seq, + capture_ts_ms: clock.now_ms(), + encode_done_ts_ms: clock.now_ms(), + send_ts_ms: 0, + keyframe, + codec: Codec::H264, + width: side, + height: side, + }, + // Exceed the receiver's per-stream credit without draining it. + // The following independent picture is deliberately small. + payload: if seq <= blocked_end && (!self.deltas || seq > 0) { + Bytes::from(vec![1; self.old_bytes]) + } else { + payload + }, + }) + } +} + +struct Input(Arc); +impl InputSink for Input { + fn inject(&mut self, _: &InputEvent) -> Result<(), DesktopError> { + self.0.fetch_add(1, Ordering::Relaxed); + Ok(()) + } +} + +type Case = ( + (Backend, bool, bool), + Option>, + Option>, +); + +#[cfg(feature = "x11")] +fn native_cases() -> Vec { + [Backend::Iroh, Backend::Noq] + .into_iter() + .map(|backend| { + ( + (backend, false, false), + Some( + Box::new(rds_desktop::H264Encoder::new(4_000_000, 60.0).unwrap()) + as Box, + ), + Some(Box::new(rds_desktop::H264Decoder::new().unwrap()) as Box), + ) + }) + .collect() +} +#[cfg(not(feature = "x11"))] +fn native_cases() -> Vec { + Vec::new() +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn explicit_repair_supersedes_blocked_media_without_closing_controls() { + let cases = [Backend::Iroh, Backend::Noq] + .into_iter() + .flat_map(|backend| [false, true].map(move |deltas| (backend, deltas, false))) + .chain([(Backend::Iroh, false, true)]) + .map(|case| (case, None, None)) + .chain(native_cases()); + for ((backend, deltas, delayed_receipt), encoder, mut decoder) in cases { + tokio::time::timeout(Duration::from_secs(20), async { + let config = EndpointConfig { + backend, + discovery: false, + bind_addrs: vec!["127.0.0.1:0".parse().unwrap()], + max_multipath_paths: delayed_receipt.then_some(1), + ..Default::default() + }; + let client = rds_net::bind_endpoint(config.clone()).await.unwrap(); + let server = rds_net::bind_endpoint(config).await.unwrap(); + let mut target = server.addr(); + let proxy = if delayed_receipt { + let upstream = target + .addrs + .iter() + .find_map(|addr| { + if let rds_core::TransportAddr::Ip(ip) = addr { + Some(*ip) + } else { + None + } + }) + .unwrap(); + let proxy = rds_bench::impair::spawn( + upstream, + rds_bench::impair::Impairment { + delay_ms: 80, + rate_mbps: Some(1.0), + ..Default::default() + }, + ) + .await + .unwrap(); + target.addrs = [rds_core::TransportAddr::Ip(proxy.listen)] + .into_iter() + .collect(); + Some(proxy) + } else { + None + }; + let (a, b) = tokio::join!(client.connect(target, rds_core::ALPN), async { + server.accept().await.unwrap().await + }); + let (a, b) = (a.unwrap(), b.unwrap()); + let mut frames = a.uni_streams(rds_core::UniHello::Desktop).unwrap(); + let (mut control, mut replies) = a.open_bi().await.unwrap(); + write_frame( + &mut control, + &DesktopControl::Heartbeat { seq: 0, ts_ms: 0 }, + ) + .await + .unwrap(); + let (send, recv) = b.accept_bi().await.unwrap(); + let (stop, stopped) = std::sync::mpsc::channel(); + let injected = Arc::new(AtomicU64::new(0)); + let produced = Arc::new(AtomicU64::new(0)); + let serving = tokio::spawn(serve_desktop_with( + b.clone(), + send, + recv, + DesktopHello { + display: 0, + max_fps: 60, + codec: Codec::H264, + input_acks: true, + }, + SessionConfig { + producer: Some(Box::new(BlockedKey { + stop: stopped, + deltas, + produced: produced.clone(), + // Below stream credit, enqueue/FIN complete while + // the rate-capped proxy keeps the receipt pending. + old_bytes: if delayed_receipt { + 512 * 1024 + } else { + 16 * 1024 * 1024 + }, + encoder, + })), + input_sink: Some(Box::new(Input(injected.clone()))), + ..Default::default() + }, + )); + let mut old = frames.recv().await.unwrap(); + let mut old_header: FrameHeader = read_frame(&mut old).await.unwrap(); + assert_eq!(old_header.seq, 0); + assert!(old_header.keyframe); + let mut held = Vec::new(); + if deltas { + assert_eq!(old.read_to_end(512).await.unwrap(), &[7; 128]); + // One partial send and two queued references fill all three + // admission slots. Do not require their blocked stream tags + // to arrive before requesting the repair that releases them. + let mut stream = frames.recv().await.unwrap(); + let header: FrameHeader = read_frame(&mut stream).await.unwrap(); + assert_eq!(header.seq, 1); + assert!(!header.keyframe); + held.push(stream); + old_header = header; + tokio::time::timeout(Duration::from_secs(2), async { + while produced.load(Ordering::Acquire) < 4 { + tokio::task::yield_now().await; + } + }) + .await + .expect("fixture did not fill its media admission slots"); + } else { + held.push(old); + } + // Keep bodies unread: FIN cannot be delivered within stream credit. + let observation_started = std::time::Instant::now(); + write_frame(&mut control, &DesktopControl::RequestIdr) + .await + .unwrap(); + write_frame( + &mut control, + &DesktopControl::Input(InputEvent { + seq: 17, + event_ts_ms: 0, + display_id: 0, + kind: InputKind::KeyDown { code: 30 }, + }), + ) + .await + .unwrap(); + write_frame( + &mut control, + &DesktopControl::Heartbeat { seq: 77, ts_ms: 11 }, + ) + .await + .unwrap(); + let mut input_acked = false; + tokio::time::timeout(Duration::from_secs(2), async { + loop { + match read_frame::<_, DesktopEvent>(&mut replies).await.unwrap() { + DesktopEvent::InputAck { seq: 17, .. } => input_acked = true, + DesktopEvent::Heartbeat { seq: 77, ts_ms: 11 } => break, + _ => {} + } + } + }) + .await + .expect("control was blocked by old media"); + assert!(input_acked); + assert_eq!(injected.load(Ordering::Relaxed), 1); + // Echo is a fence: the preceding RequestIdr was handled by the peer. + let fresh = tokio::time::timeout(Duration::from_secs(1), async { + let mut stream = frames.recv().await.unwrap(); + let header: FrameHeader = read_frame(&mut stream).await.unwrap(); + let body = stream.read_to_end(512).await.unwrap(); + (header, body) + }) + .await; + println!( + "backend={backend:?} deltas={deltas} delayed_receipt={delayed_receipt} request_to_observation_ms={}", + observation_started.elapsed().as_millis() + ); + if let Some(proxy) = &proxy { + assert!( + proxy.stats().forwarded > 10, + "receipt fixture escaped its proxy" + ); + assert!(a.current_path_stats().unwrap().rtt >= Duration::from_millis(100)); + } + if fresh.is_ok() { + for mut old in held { + let reset = tokio::time::timeout(Duration::from_secs(2), async { + let mut chunk = [0; 16 * 1024]; + loop { + match old.read(&mut chunk).await { + Ok(Some(_)) => {} + Err(rds_net::ReadError::Reset(code)) => break code, + other => panic!("old media was not reset: {other:?}"), + } + } + }) + .await + .expect("old frame was retained after recovery"); + assert_eq!(reset, 1u32.into()); + } + // An unrelated stream still works after media-only cancellation. + let (mut other, mut echo) = a.open_bi().await.unwrap(); + other.write_all(b"other").await.unwrap(); + let (mut other_send, mut other_recv) = b.accept_bi().await.unwrap(); + let mut bytes = [0; 5]; + other_recv.read_exact(&mut bytes).await.unwrap(); + assert_eq!(&bytes, b"other"); + other_send.write_all(b"usable").await.unwrap(); + let mut bytes = [0; 6]; + echo.read_exact(&mut bytes).await.unwrap(); + assert_eq!(&bytes, b"usable"); + } + let _ = stop.send(()); + serving.abort(); + let _ = serving.await; + a.close(0u32.into(), b"fixture done"); + b.close(0u32.into(), b"fixture done"); + tokio::join!(client.close(), server.close()); + let (header, body) = + fresh.expect("repair waited for the obsolete key receipt deadline"); + assert!(header.seq > old_header.seq); + assert!( + header.keyframe, + "repair must start an independent reference chain" + ); + if let Some(decoder) = decoder.as_mut() { + let raw = decoder.decode(&EncodedFrame { codec: header.codec, + keyframe: header.keyframe, data: Bytes::from(body) }).unwrap().unwrap(); + assert_eq!((raw.width, raw.height), (64, 64)); + let expected = (50 + header.seq * 43 % 150) as i16; + assert!((i16::from(raw.data[0]) - expected).abs() < 12, + "recovered native pixel did not match the new independent picture"); + } else { assert_eq!(body, &[7; 128]); } + }) + .await + .expect("media repair fixture did not release its resources"); + } +} diff --git a/crates/rds-desktop/tests/x11_native.rs b/crates/rds-desktop/tests/x11_native.rs index d3b3c38..f6a4541 100644 --- a/crates/rds-desktop/tests/x11_native.rs +++ b/crates/rds-desktop/tests/x11_native.rs @@ -22,6 +22,158 @@ fn observer() -> (RustConnection, u32) { (conn, root) } +#[test] +#[ignore = "requires a dedicated Xvfb server; injects native input"] +fn delayed_key_release_does_not_generate_remote_typematic_repeats() { + use std::time::{Duration, Instant}; + use x11rb::protocol::{ + Event, + xkb::{BoolCtrl, ConnectionExt as _, Control, ID}, + xproto::{ + AutoRepeatMode, ChangeKeyboardControlAux, ChangeWindowAttributesAux, EventMask, + InputFocus, + }, + }; + let (observer, root) = observer(); + observer.xkb_use_extension(1, 0).unwrap().reply().unwrap(); + let original = observer + .xkb_get_controls(ID::USE_CORE_KBD.into()) + .unwrap() + .reply() + .unwrap(); + let rate = |delay, interval| { + observer + .xkb_set_controls( + ID::USE_CORE_KBD.into(), + 0u16.into(), + 0u16.into(), + 0u16.into(), + 0u16.into(), + 0u16.into(), + 0u16.into(), + 0u16.into(), + 0u16.into(), + original.mouse_keys_dflt_btn, + original.groups_wrap, + original.access_x_option, + BoolCtrl::REPEAT_KEYS, + BoolCtrl::REPEAT_KEYS, + Control::from(u32::from(BoolCtrl::REPEAT_KEYS)), + delay, + interval, + original.slow_keys_delay, + original.debounce_delay, + original.mouse_keys_delay, + original.mouse_keys_interval, + original.mouse_keys_time_to_max, + original.mouse_keys_max_speed, + original.mouse_keys_curve, + original.access_x_timeout, + original.access_x_timeout_mask, + original.access_x_timeout_values, + original.access_x_timeout_options_mask, + original.access_x_timeout_options_values, + &original.per_key_repeat, + ) + .unwrap() + .check() + .unwrap(); + }; + rate(40, 20); + observer + .change_keyboard_control( + &ChangeKeyboardControlAux::new() + .key(38u32) + .auto_repeat_mode(AutoRepeatMode::ON), + ) + .unwrap() + .check() + .unwrap(); + observer + .change_window_attributes( + root, + &ChangeWindowAttributesAux::new() + .event_mask(EventMask::KEY_PRESS | EventMask::KEY_RELEASE), + ) + .unwrap() + .check() + .unwrap(); + observer + .set_input_focus(InputFocus::POINTER_ROOT, root, x11rb::CURRENT_TIME) + .unwrap() + .check() + .unwrap(); + let mut sink = XtestInput::new().unwrap(); + sink.inject(&event(InputKind::KeyDown { code: 30 })) + .unwrap(); + std::thread::sleep(Duration::from_millis(160)); // delayed network KeyUp + let keys = observer.query_keymap().unwrap().reply().unwrap().keys; + assert_ne!( + keys[38 / 8] & (1 << (38 % 8)), + 0, + "genuine key hold was released early" + ); + sink.inject(&event(InputKind::KeyUp { code: 30 })).unwrap(); + let mut presses = 0; + let deadline = Instant::now() + Duration::from_millis(50); + while Instant::now() < deadline { + if let Some(Event::KeyPress(key)) = observer.poll_for_event().unwrap() { + if key.detail == 38 { + presses += 1; + } + } else { + std::thread::sleep(Duration::from_millis(1)); + } + } + assert_eq!( + presses, 1, + "a delayed release manufactured extra letters on the server" + ); + for _ in 0..3 { + sink.inject(&event(InputKind::KeyDown { code: 30 })) + .unwrap(); + } + sink.inject(&event(InputKind::KeyUp { code: 30 })).unwrap(); + let mut repeats = 0; + while let Some(e) = observer.poll_for_event().unwrap() { + if matches!(e, Event::KeyPress(key) if key.detail == 38) { + repeats += 1; + } + } + assert_eq!( + repeats, 3, + "intentional client repeats were lost or duplicated" + ); + let mut other = XtestInput::new().unwrap(); + sink.inject(&event(InputKind::KeyDown { code: 30 })) + .unwrap(); + other + .inject(&event(InputKind::KeyDown { code: 30 })) + .unwrap(); + sink.inject(&event(InputKind::KeyUp { code: 30 })).unwrap(); + let keys = observer.query_keymap().unwrap().reply().unwrap().keys; + assert_ne!( + keys[38 / 8] & (1 << (38 % 8)), + 0, + "one controller released another's hold" + ); + drop(other); + let keys = observer.query_keymap().unwrap().reply().unwrap().keys; + assert_eq!( + keys[38 / 8] & (1 << (38 % 8)), + 0, + "last controller retained the key after drop" + ); + let keyboard = observer.get_keyboard_control().unwrap().reply().unwrap(); + assert_ne!( + keyboard.auto_repeats[38 / 8] & (1 << (38 % 8)), + 0, + "original native repeat setting was not restored" + ); + drop(sink); + rate(original.repeat_delay, original.repeat_interval); +} + #[test] #[ignore = "requires a dedicated Xvfb server; injects native input"] fn native_evdev_keyboard_mapping() { diff --git a/crates/rds-net/tests/metrics.rs b/crates/rds-net/tests/metrics.rs index 8d36f25..d52f3ef 100644 --- a/crates/rds-net/tests/metrics.rs +++ b/crates/rds-net/tests/metrics.rs @@ -28,12 +28,32 @@ async fn echo_once(ep: &rds_net::Endpoint) { let n = recv.read(&mut buf).await.unwrap().unwrap(); send.write_all(&buf[..n]).await.unwrap(); send.finish().unwrap(); + // finish() queues the echo; it does not prove the UDP driver sent it. + // Fence on the client's observation before sampling server traffic. + let (mut fence_send, mut fence_recv) = conn.accept_bi().await.unwrap(); + let mut observed = [0]; + fence_recv.read_exact(&mut observed).await.unwrap(); + assert_eq!(observed, [1]); sampler.sample(); + fence_send.write_all(&observed).await.unwrap(); + fence_send.finish().unwrap(); // Hold the connection until the peer closes it — dropping the last // handle now could cut the echo reply before the client reads it. let _ = conn.accept_bi().await; } +async fn reply_observed(conn: &rds_net::Connection) { + let (mut send, mut recv) = conn.open_bi().await.unwrap(); + send.write_all(&[1]).await.unwrap(); + send.finish().unwrap(); + let mut sampled = [0]; + tokio::time::timeout(Duration::from_secs(5), recv.read_exact(&mut sampled)) + .await + .expect("server did not sample observed echo") + .unwrap(); + assert_eq!(sampled, [1]); +} + fn counter(reg: &Registry, name: &str) -> u64 { *reg.snapshot().get(name).unwrap_or(&0) } @@ -72,6 +92,7 @@ async fn direct_traffic_counts_direct_not_relay() { let mut buf = [0u8; 9]; recv.read_exact(&mut buf).await.unwrap(); assert_eq!(&buf, b"ping-pong"); + reply_observed(&conn).await; let mut sampler = client.metrics().sampler(conn.clone()); sampler.sample(); @@ -157,6 +178,7 @@ async fn relay_only_traffic_counts_relay_not_direct() { .await .expect("relay echo timed out") .unwrap(); + reply_observed(&conn).await; let mut sampler = client.metrics().sampler(conn.clone()); sampler.sample(); diff --git a/docs/architecture.md b/docs/architecture.md index f6f537b..9b77a16 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -89,6 +89,14 @@ input ACK/heartbeat observation. Both sockets share desktop termination, without closing the manager's unrelated peer streams. Existing combined APIs and remote wire formats remain; see [the native contract](native-viewer.md). +Explicit video repair advances a session-local media epoch and interrupts only +obsolete media work. An independent-picture lease holds admission until receipt; +repeated requests coalesce during that recovery. X11 keyboard holds similarly +own per-key native repeat suppression: the viewer forwards local OS repeats, +so network delay cannot manufacture additional letters. Both native endpoints +need the coordinated update. See [the regression receipt](reports/rds-interactive-repair-20261004.md) +for tested boundaries and remaining installed-device qualification. + Owned path policy now uses [validated eligibility](path-selection.md): only the handshake path is seeded; application-opened candidates stay Backup until an Established event. An owned bounded queue retries temporary path-credit diff --git a/docs/native-viewer.md b/docs/native-viewer.md index 4851000..8bdb8e4 100644 --- a/docs/native-viewer.md +++ b/docs/native-viewer.md @@ -518,3 +518,17 @@ Enqueue spans stream opening/writes and scheduling of the receipt worker; ACK spans that worker's transport receipt wait. Neither is an application input ACK or a physical presentation measurement. Timely frame completions remain trace level. The diagnostic addition changes no pacing, admission or deadlines. + +Explicit receiver repair now advances a session-owned media epoch. It retires +obsolete writes/receipt waits without closing control or other connection +services, and forces an independent producer reference. Repeated requests +coalesce until that independent picture is transport-confirmed. Health logs +expose repair generation, accepted/coalesced requests and the active repair gate. + +Native keyboard repeat is client-paced: the viewer forwards native OS repeats +only for keys it already forwarded and still holds, without repeating clipboard +side effects. The X11 sink suppresses autonomous per-key repeat during those +holds, pulses an intentional client repeat and restores the original setting +when its last controller releases. This prevents a late KeyUp from manufacturing +letters. Update both native viewer and serving agent together. See the +[regression receipt](reports/rds-interactive-repair-20261004.md) for scope limits. diff --git a/docs/reports/rds-interactive-repair-20261004.md b/docs/reports/rds-interactive-repair-20261004.md new file mode 100644 index 0000000..fb42cb1 --- /dev/null +++ b/docs/reports/rds-interactive-repair-20261004.md @@ -0,0 +1,56 @@ +# Interactive repair and client-paced repeat — 2026-10-04 + +W6.6/W6.7 follow-up; native two-device latency/stability acceptance remains open. + +An explicit RequestIdr previously set a producer flag while the producer could +be waiting for an obsolete independent picture or all three media permits. +The isolated transport regression kept input/heartbeat working but failed its +one-second replacement bound before the change. Session-local repair epochs +now cancel only that writer's obsolete media writes/receipt tasks. Raced encodes +and stale queued references are discarded; the next producer epoch forces IDR. + +An owned key lease closes the admission gate until transport confirmation. An +old producer/receipt cannot install or release a newer epoch's gate. Dropping a +fresh key during cancellation requests another independent picture before its +successors. Requests coalesce while a recovery picture is in flight, preventing +repeated watchdog observations from repeatedly canceling its valid progress. +Epoch and gate flags update atomically. Existing deadlines, queue/reader budgets, +connection services and wire types remain unchanged. + +Real-loopback cases cover Iroh/Noq partial independent writes and full admission, +an Iroh receipt wait through a pinned 80 ms/1 Mbps UDP proxy, native H.264 +replacement decode/pixel checks, input ACKs and an unrelated bidirectional +stream. The reported request-to-observation bound includes its control fence; +it is not physical presentation or production WAN latency. + +A second native Xvfb regression reproduced eight keypress events from one +physical-style KeyDown with its KeyUp delayed 160 ms. Server typematic generated +extra letters independently of the client's intent. The native viewer previously +ignored local OS repeat events. X11 injection now leases per-key native repeat +to off while retaining a real held key. Local repeats for forwarded, still-held +non-modifier keys pulse exactly one release/press. Repeated modifiers/locks do +not toggle, and repeat events do not repeat clipboard side effects. + +Keyboard leases share counts across controllers and X screens. The last release +restores the original per-key native repeat setting. The native regression now +observes one keypress, three intentional client repeats, a hold surviving one +controller's release and complete key/repeat cleanup after the last owner drops. +It runs only on disposable Xvfb, never a developer's active desktop. + +The native viewer and serving agent must be updated together for client-driven +repeat behavior. Hardware/composited input, sustained WAN loss/jitter, and native +causal click-to-visible timing require installed qualification. This report is +not a completed stability gate; final checks belong to the published candidate. + +Local validation: formatting; default and desktop workspace Clippy with warnings +denied on macOS; the complete desktop suite with X11/viewer enabled; Linux X11 +workspace Clippy; all eight native X11 input cases on disposable Xvfb. Default +workspace and published-candidate CI results are recorded by their actual run, +not inferred from these feature-specific checks. No real desktop input was +injected by the native regression fixture. + +The first Linux CI run exposed an existing metrics-fixture race: server sampling +immediately after send.finish() could precede actual UDP delivery. The echo test +now uses a peer-observation fence before sampling, retaining the traffic +assertions without arbitrary sleep. This changes test synchronization only; +it does not change production transport accounting. diff --git a/docs/research.md b/docs/research.md index 778a9c7..2eef671 100644 --- a/docs/research.md +++ b/docs/research.md @@ -706,3 +706,19 @@ exclude sparse tiny updates, not the WebRTC ALR algorithm or a capacity estimate Actual stalled outstanding delivery, hard failure, path feedback and existing bounded growth/deadlines remain. The [regression report](reports/rds-payload-pressure-20261003.md) separates the reproduced quality collapse from installed network acceptance. + +## 2026-10-04 client-paced keyboard repeat + +The [X11 keyboard control contract](https://xorg.freedesktop.org/archive/current/doc/libX11/libX11/libX11.html#Manipulating_the_Keyboard_and_Pointer_Settings) +defines independent per-key repeat bits and alternating press/release events +generated by server typematic. [Winit's repeat flag](https://docs.rs/winit/0.30.13/winit/event/struct.KeyEvent.html#structfield.repeat) +identifies intentional local OS repeats. These contracts support retaining a +real remote hold while letting the local keyboard determine repeat count. + +The native regression reproduced eight presses from one down event with a +160 ms delayed release. Per-key repeat leases now suppress that amplification; +the viewer forwards intentional repeats without re-running clipboard effects. +Release/drop restores original native repeat state after the last owner. +Abnormal process/X-server failure and competing local clients remain outside +that cleanup guarantee. See [the combined repair receipt](reports/rds-interactive-repair-20261004.md) +for the isolated fixture and coordinated-update requirement. diff --git a/docs/x11-input.md b/docs/x11-input.md index aae8c28..db66ce7 100644 --- a/docs/x11-input.md +++ b/docs/x11-input.md @@ -16,6 +16,16 @@ Buttons map explicitly: left/right/middle to X buttons 1/3/2 and evdev side/extra/forward/back/task to 8–12. Server pointer remapping still applies; unsupported native buttons produce an error rather than a successful ACK. +Native repeat is disabled only for an injected held key. The viewer supplies +intentional local OS repeats as repeated KeyDown events; non-modifier repeats +pulse one release/press while retaining the hold. Network-delayed KeyUp does +not generate server-side typematic letters. Shared per-key leases across +controllers/screens restore the original repeat bit after the final release. +Viewer and serving agent must be updated together. Native state restoration is +best effort on teardown; a killed process or failed X connection cannot guarantee +cleanup. Other local X clients and distinct agent processes do not participate +in this process-local ownership. Concurrent local-seat ownership remains open. + Absolute motion uses the selected root window and checked display coordinates. Relative motion uses XTEST's relative flag, with checked signed 16-bit deltas; it never interprets offsets as absolute coordinates. Absolute positions round @@ -68,5 +78,6 @@ Xvfb is a test fixture only. Physical/composited X11, Wayland and macOS native qualification and a graphical viewer still have their own gates. Primary contracts: [XTEST](https://www.x.org/releases/X11R7.7/doc/xextproto/xtest.pdf), +[per-key repeat control](https://xorg.freedesktop.org/archive/current/doc/libX11/libX11/libX11.html#Manipulating_the_Keyboard_and_Pointer_Settings), [Xorg evdev mapping](https://cgit.freedesktop.org/xorg/driver/xf86-input-evdev/tree/src/evdev.c), [Linux event codes](https://docs.kernel.org/input/event-codes.html).