diff --git a/crates/rds-cli/src/desktop.rs b/crates/rds-cli/src/desktop.rs index 1971e05..4d7384c 100644 --- a/crates/rds-cli/src/desktop.rs +++ b/crates/rds-cli/src/desktop.rs @@ -81,6 +81,9 @@ pub async fn read_grant( } } +#[cfg(feature = "desktop")] +mod control; + #[cfg(feature = "desktop")] mod native { use super::*; @@ -263,8 +266,22 @@ mod native { failures = 0; } failures = failures.saturating_add(1); + let snapshot = view.snapshot(); + tracing::warn!( + error = ?result.as_ref().err(), attempt = failures, + network_stage = %snapshot.network_stage, + render_stage = %snapshot.render_stage, + decoded_frame_age_ms = snapshot.decoded_frame_age_ms, + submission_age_ms = snapshot.submission_age_ms, + ui_event_age_ms = snapshot.ui_event_age_ms, + pending_input_acks = snapshot.report.pending_input_acks, + oldest_input_ack_age_ms = snapshot.report.oldest_input_ack_age_ms, + control_rtt_ms = snapshot.report.control_rtt_ms, + occluded = snapshot.occluded, + "desktop reconnecting" + ); view.status("Reconnecting"); - tracing::warn!(error = ?result.as_ref().err(), attempt=failures,"desktop reconnecting"); + view.stage("reconnecting"); let wait = Duration::from_millis((250u64 << failures.min(5)).min(8000)); if !retry_pause(wait, &stop, input).await { return Ok(()); @@ -339,46 +356,68 @@ mod native { extent(view, &channel.caps, options.display)?; view.status("Waiting for screen"); let control = channel.control_handle(); - let mut decoder = RelayDecoder::new(); - let mut last_frame = Instant::now(); - let mut tick = tokio::time::interval(Duration::from_secs(1)); - tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); - let result = loop { - tokio::select! { - _ = stop.cancelled() => break Ok(true), - message = input.recv() => match message { - Some(ViewerInput::Control(message)) => { - view.input_sent(&message); - tokio::time::timeout(Duration::from_secs(2), control.control(message)).await.map_err(|_|anyhow::anyhow!("desktop control write stalled"))??; - }, - Some(ViewerInput::Close)|None => break Ok(true), - }, - _ = tick.tick() => { - anyhow::ensure!(last_frame.elapsed() < Duration::from_secs(15),"remote video stopped making progress"); - tokio::time::timeout(Duration::from_secs(2), control.control(rds_core::DesktopControl::Heartbeat { seq: 0,ts_ms: started.elapsed().as_millis() as u64 })).await.map_err(|_|anyhow::anyhow!("desktop heartbeat write stalled"))??; - }, - message = channel.recv() => match message? { + let (progress, last_frame) = tokio::sync::watch::channel(tokio::time::Instant::now()); + // Keep the entire receive/decode future alive while controls progress. + // Selecting individual recv calls and awaiting decode in their handler + // prevents input, heartbeat and close from being polled during decode. + let media = async { + let mut decoder = RelayDecoder::new(); + loop { + match channel.recv().await? { None => break Ok(false), - Some(ManagedMessage::Event(rds_core::DesktopEvent::Heartbeat { ts_ms,.. })) => view.control_rtt((started.elapsed().as_millis() as u64).saturating_sub(ts_ms)), - Some(ManagedMessage::Event(rds_core::DesktopEvent::InputAck { seq,.. })) => view.input_ack(seq), - Some(ManagedMessage::Event(rds_core::DesktopEvent::ClipboardReady { bytes,.. })) => {view.clipboard_ready(bytes);tracing::info!(bytes,"remote clipboard ready");}, + Some(ManagedMessage::Event(rds_core::DesktopEvent::Heartbeat { + ts_ms, + .. + })) => view + .control_rtt((started.elapsed().as_millis() as u64).saturating_sub(ts_ms)), + Some(ManagedMessage::Event(rds_core::DesktopEvent::InputAck { + seq, .. + })) => view.input_ack(seq), + Some(ManagedMessage::Event(rds_core::DesktopEvent::ClipboardReady { + bytes, + .. + })) => { + view.clipboard_ready(bytes); + tracing::info!(bytes, "remote clipboard ready"); + } Some(ManagedMessage::Frame(frame)) => { view.stage("decoding"); let received = Instant::now(); view.media_timing(&frame.header); - let (next,outcome) = tokio::time::timeout(Duration::from_secs(5), decoder.push_bounded(frame.header,frame.payload)).await.map_err(|_|anyhow::anyhow!("desktop decode stalled"))??; + let (next, outcome) = tokio::time::timeout( + Duration::from_secs(5), + decoder.push_bounded(frame.header, frame.payload), + ) + .await + .map_err(|_| anyhow::anyhow!("desktop decode stalled"))??; decoder = next; match outcome { - RelayOutcome::Frame(raw) => { last_frame = Instant::now(); view.frame(raw,received); }, - RelayOutcome::NeedIdr => control.request_idr().await?, - RelayOutcome::Pending => {}, + RelayOutcome::Frame(raw) => { + progress.send_replace(tokio::time::Instant::now()); + view.frame(raw, received); + } + RelayOutcome::NeedIdr => { + tokio::time::timeout(Duration::from_secs(2), control.request_idr()) + .await + .map_err(|_| { + anyhow::anyhow!("desktop repair write stalled") + })?? + } + RelayOutcome::Pending => {} } view.stage("receiving"); } } } }; - let _ = tokio::time::timeout(Duration::from_secs(1), channel.finish()).await; + let controls = control::pump(input, &control, &last_frame, started, |message| { + view.input_sent(message); + }); + let result = control::run(controls, media, stop).await; + // A winning leg may cancel a partially written control on the other + // leg. EOF closes the manager's desktop; never append Finished to a + // potentially incomplete frame. Unrelated manager sessions survive. + drop(channel); result } diff --git a/crates/rds-cli/src/desktop/control.rs b/crates/rds-cli/src/desktop/control.rs new file mode 100644 index 0000000..9835dd3 --- /dev/null +++ b/crates/rds-cli/src/desktop/control.rs @@ -0,0 +1,317 @@ +//! Managed native control dispatch progresses independently of decode waits. +use rds_core::DesktopControl; +use rds_desktop::render::{InputReceiver, ViewerInput}; +use std::{ + future::Future, + time::{Duration, Instant}, +}; +use tokio::sync::watch; +use tokio_util::sync::CancellationToken; + +pub(super) trait Input { + fn recv(&mut self) -> impl Future> + Send; +} +impl Input for InputReceiver { + async fn recv(&mut self) -> Option { + self.recv().await + } +} +pub(super) trait Sender { + fn send(&self, message: DesktopControl) -> impl Future> + Send; +} +impl Sender for rds_client::local::ManagedControl { + async fn send(&self, message: DesktopControl) -> anyhow::Result<()> { + Ok(self.control(message).await?) + } +} + +pub(super) async fn pump( + input: &mut impl Input, + sender: &impl Sender, + progress: &watch::Receiver, + started: Instant, + sent: impl Fn(&DesktopControl), +) -> anyhow::Result { + let mut tick = tokio::time::interval(Duration::from_secs(1)); + tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + let message = tokio::select! { + message = input.recv() => match message { + Some(ViewerInput::Control(message)) => message, + Some(ViewerInput::Close) | None => return Ok(true), + }, + _ = tick.tick() => { + let age = progress.borrow().elapsed(); + if age >= Duration::from_secs(15) { + tracing::warn!(last_decoded_age_ms = age.as_millis() as u64, + "desktop progress watchdog expired"); + anyhow::bail!("remote video stopped making progress"); + } + DesktopControl::Heartbeat { + seq: 0, ts_ms: started.elapsed().as_millis() as u64, + } + }, + }; + sent(&message); + // A partially written control is never reused after timeout: the + // owner ends both legs and drops the channel before reconnecting. + tokio::time::timeout(Duration::from_secs(2), sender.send(message)) + .await + .map_err(|_| anyhow::anyhow!("desktop control write stalled"))??; + } +} + +pub(super) async fn run( + controls: impl Future>, + media: impl Future>, + stop: &CancellationToken, +) -> anyhow::Result { + // Neither leg is respawned or reconstructed after the other wakes. On + // completion all pending async work is dropped with the owning session. + tokio::select! { + biased; + _ = stop.cancelled() => Ok(true), + result = controls => result, + result = media => result, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use rds_core::{InputEvent, InputKind, local::DesktopUp}; + use std::sync::{ + Arc, + atomic::{AtomicBool, Ordering}, + }; + use tokio::{ + io::DuplexStream, + sync::{Mutex, mpsc, oneshot}, + }; + + impl Input for mpsc::Receiver { + async fn recv(&mut self) -> Option { + self.recv().await + } + } + struct Wire(Mutex); + impl Sender for Wire { + async fn send(&self, message: DesktopControl) -> anyhow::Result<()> { + Ok( + rds_net::write_frame(&mut *self.0.lock().await, &DesktopUp::Control(message)) + .await?, + ) + } + } + struct Dropped(Arc); + impl Drop for Dropped { + fn drop(&mut self) { + self.0.store(true, Ordering::SeqCst); + } + } + + // Read the actual framed controls, rather than only observing that an + // input was dequeued. The media leg remains held until the session ends. + #[tokio::test(start_paused = true)] + async fn blocked_decode_does_not_hold_input_or_heartbeat_and_close_drops_it() { + let (wire, mut remote) = tokio::io::duplex(256); + let (tx, mut input) = mpsc::channel(8); + let stop = CancellationToken::new(); + let dropped = Arc::new(AtomicBool::new(false)); + let held = dropped.clone(); + let (decoding, started_decode) = oneshot::channel(); + let (_release, blocked) = oneshot::channel::<()>(); + let task = tokio::spawn(async move { + let (_progress, last_frame) = watch::channel(tokio::time::Instant::now()); + let media = async { + let _drop = Dropped(held); + decoding.send(()).unwrap(); + blocked.await.unwrap(); + Ok(false) + }; + run( + pump( + &mut input, + &Wire(Mutex::new(wire)), + &last_frame, + Instant::now(), + |_| {}, + ), + media, + &stop, + ) + .await + }); + started_decode.await.unwrap(); + let started = tokio::time::Instant::now(); + for (seq, kind) in [ + InputKind::KeyDown { code: 56 }, + InputKind::KeyDown { code: 105 }, + InputKind::KeyUp { code: 105 }, + InputKind::KeyUp { code: 56 }, + ] + .into_iter() + .enumerate() + { + tx.send(ViewerInput::Control(DesktopControl::Input(InputEvent { + seq: seq as u64, + event_ts_ms: 10, + display_id: 0, + kind, + }))) + .await + .unwrap(); + } + let mut observed = Vec::new(); + let mut heartbeats = 0; + while observed.len() < 4 || heartbeats < 2 { + let message: DesktopUp = tokio::time::timeout( + Duration::from_millis(1200), + rds_net::read_frame(&mut remote), + ) + .await + .unwrap() + .unwrap(); + match message { + DesktopUp::Control(DesktopControl::Input(event)) => observed.push(event), + DesktopUp::Control(DesktopControl::Heartbeat { .. }) => heartbeats += 1, + other => panic!("unexpected control: {other:?}"), + } + } + assert_eq!( + observed.iter().map(|e| e.seq).collect::>(), + [0, 1, 2, 3] + ); + assert!(started.elapsed() < Duration::from_millis(1200)); + assert!( + !dropped.load(Ordering::SeqCst), + "decoder must stay held during dispatch" + ); + tx.send(ViewerInput::Close).await.unwrap(); + assert!( + tokio::time::timeout(Duration::from_millis(1), task) + .await + .unwrap() + .unwrap() + .unwrap() + ); + assert!(dropped.load(Ordering::SeqCst)); + } + + #[tokio::test(start_paused = true)] + async fn blocked_decode_cannot_disable_progress_watchdog() { + let (wire, mut remote) = tokio::io::duplex(256); + let (_tx, mut input) = mpsc::channel(1); + let stop = CancellationToken::new(); + let task = tokio::spawn(async move { + let (_progress, last_frame) = watch::channel(tokio::time::Instant::now()); + run( + pump( + &mut input, + &Wire(Mutex::new(wire)), + &last_frame, + Instant::now(), + |_| {}, + ), + std::future::pending(), + &stop, + ) + .await + }); + let started = tokio::time::Instant::now(); + for _ in 0..15 { + let _: DesktopUp = rds_net::read_frame(&mut remote).await.unwrap(); + } + let error = task.await.unwrap().unwrap_err(); + assert!(error.to_string().contains("video stopped making progress")); + assert_eq!(started.elapsed(), Duration::from_secs(15)); + } + + #[tokio::test(start_paused = true)] + async fn stalled_control_write_is_bounded_and_cancel_releases_both_legs() { + for cancel in [false, true] { + let (wire, _remote) = tokio::io::duplex(1); + let (tx, mut input) = mpsc::channel(1); + tx.send(ViewerInput::Control(DesktopControl::RequestIdr)) + .await + .unwrap(); + let stop = CancellationToken::new(); + let canceled = stop.clone(); + let dropped = Arc::new(AtomicBool::new(false)); + let held = dropped.clone(); + let (decoding, started_decode) = oneshot::channel(); + let task = tokio::spawn(async move { + let (_progress, last_frame) = watch::channel(tokio::time::Instant::now()); + let media = async { + let _drop = Dropped(held); + decoding.send(()).unwrap(); + std::future::pending().await + }; + run( + pump( + &mut input, + &Wire(Mutex::new(wire)), + &last_frame, + Instant::now(), + |_| {}, + ), + media, + &stop, + ) + .await + }); + started_decode.await.unwrap(); + let started = tokio::time::Instant::now(); + if cancel { + canceled.cancel(); + } + let result = task.await.unwrap(); + if cancel { + assert!(result.unwrap()); + assert_eq!(started.elapsed(), Duration::ZERO); + } else { + assert!( + result + .unwrap_err() + .to_string() + .contains("control write stalled") + ); + assert_eq!(started.elapsed(), Duration::from_secs(2)); + } + assert!(dropped.load(Ordering::SeqCst)); + } + } + + #[tokio::test(start_paused = true)] + async fn fresh_decoded_progress_renews_watchdog_without_restarting_control_leg() { + let (wire, mut remote) = tokio::io::duplex(256); + let (tx, mut input) = mpsc::channel(1); + let (progress, last_frame) = watch::channel(tokio::time::Instant::now()); + let stop = CancellationToken::new(); + let task = tokio::spawn(async move { + run( + pump( + &mut input, + &Wire(Mutex::new(wire)), + &last_frame, + Instant::now(), + |_| {}, + ), + std::future::pending(), + &stop, + ) + .await + }); + let started = tokio::time::Instant::now(); + for second in 0..31 { + let _: DesktopUp = rds_net::read_frame(&mut remote).await.unwrap(); + if second == 10 || second == 20 { + progress.send_replace(tokio::time::Instant::now()); + } + } + assert_eq!(started.elapsed(), Duration::from_secs(30)); + assert!(!task.is_finished()); + tx.send(ViewerInput::Close).await.unwrap(); + assert!(task.await.unwrap().unwrap()); + } +} diff --git a/docs/architecture.md b/docs/architecture.md index 6d471b2..2151925 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -81,6 +81,12 @@ a wgpu surface and one pending BGRA image. The CLI owns the cancelable network worker and reopens a desktop session on the same authenticated peer after loss. Per-session wire routing and broader native media acceptance remain separate. +Managed native input/heartbeat dispatch and its decoded-progress watchdog are +retained async legs independent of receive/decode waits. They share the existing +serialized control writer, with bounded writes and EOF teardown after either +leg ends. Incoming event observation still shares the bounded media IPC route; +see [the native contract](native-viewer.md) for the tested boundary. + 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 b97f416..ac301b8 100644 --- a/docs/native-viewer.md +++ b/docs/native-viewer.md @@ -6,6 +6,17 @@ decoding stays on bounded blocking workers. The serving device needs a real capture/input backend; the current Linux implementation uses X11/XTEST. macOS capture, VideoToolbox and Wayland serving remain separate work. +Managed native control dispatch has its own retained async future, polled +concurrently with the receive/decode future. Input writes, heartbeat and the +15-second decoded-progress watchdog continue while a blocking codec worker is +pending. Each control write remains bounded to two seconds; codec calls retain +the existing five-second caller bound and global worker limits. Decoder state +and encoded ordering are unchanged. Either leg's completion or cancellation +ends the desktop IPC by EOF; it never appends a control frame after canceling a +potentially partial write. Incoming acknowledgement observation can still wait +behind decode in the bounded IPC receive path. See the +[dispatch regression evidence](reports/rds-managed-control-20261002.md). + The viewer uses ordinary OS window stacking and can move behind other applications. An occluded Metal surface may pause presentation. Returning focus or uncovering the window requests an immediate redraw of its latest image. @@ -73,6 +84,11 @@ automatic grant issuance/renewal remains GDS work. Fifteen seconds without a decoded frame triggers a new desktop session. Closing the window cancels and joins its network worker, leaving unrelated managed streams intact. +Each reconnect warning records the last network/render stage, decoded and +submitted frame ages, UI dispatch age, outstanding input age/count, control RTT +and occlusion. These metadata explain the interrupted state without logging +keys, typed text, pointer coordinates, clipboard contents or screen pixels. + Discarding input during a reconnect pause keeps the original retry deadline; pointer/key activity can neither shorten nor restart it. Window close and cancellation still interrupt the pause immediately. diff --git a/docs/reports/rds-managed-control-20261002.md b/docs/reports/rds-managed-control-20261002.md new file mode 100644 index 0000000..a84541b --- /dev/null +++ b/docs/reports/rds-managed-control-20261002.md @@ -0,0 +1,62 @@ +# Managed native control progress during decode + +Scope: W6.3/W6.4 and diagnostic W6.7 remediation; native latency/stability +acceptance remains open. + +The native managed viewer previously awaited `RelayDecoder::push_bounded` +inside a selected receive handler. While that future was pending, the same +loop could not dispatch input, send heartbeat or observe close. The outer +five-second decode timeout bounded the wait but still coupled input latency +to codec scheduling. This source-level defect is separate from any physical +network outage. + +The corrected owner concurrently polls two whole, retained futures: outbound +input/heartbeat/watchdog and inbound event/decode. It adds no frame queue, +blocking task or codec permit. H.264 reference order and the one pending raw +presentation frame remain unchanged. Two-second control writes and the +15-second decoded-progress watchdog retain their bounds. A decode failure, +write failure, remote end, close or cancellation ends both futures and the +desktop IPC by EOF. No Finished frame follows a potentially canceled partial +control write. The manager's authenticated connection and unrelated streams +remain owned by the manager. + +Every reconnect warning now captures metadata about frame/UI ages, network +and render stages, outstanding input acknowledgements, control RTT and +occlusion before changing status. It contains no screen/text/input content. + +## Functional regression boundary + +`desktop::control::tests` drives the production dispatch/ownership functions +with a blocked synthetic media future and real bounded, framed duplex I/O. +It verifies: + +- Ordered Alt/arrow press and release controls reach the reader while decode + is still held; two heartbeat frames arrive within 1.2 seconds of simulated + time. Close drops the held media future within one simulated millisecond. +- A permanently held media future cannot suppress the 15-second progress + watchdog. +- A blocked control write fails at two seconds, and cancellation immediately + drops both pending legs. +- Actual decoded-progress notifications renew the watchdog through 30 seconds + without reconstructing the control future. + +These paused-clock fixtures exercise dispatch, framing and ownership. They +do not use a real OS input sink, blocking codec, GPU or physical display. +Incoming ACK/heartbeat observation can still wait behind decode in the existing +bounded IPC receive route; this change does not claim otherwise. + +## Reproduction + +```sh +cargo fmt --check +cargo test --locked -p rds-cli --all-features +cargo clippy --locked --workspace --all-targets --all-features -- -D warnings +``` + +The macOS all-feature CLI lane passed 42 tests across 12 groups, with zero +failures and one explicitly ignored external OpenSSH/account fixture. Its +desktop library group passed nine tests (four new dispatch fixtures, three +existing reconnect fixtures and two logging fixtures), with zero failures or +ignored cases. Full-workspace/all-target/all-feature strict Clippy and formatting +also passed. Linux/build results are recorded with the candidate's PR. +No physical latency, quality or long-soak gate is closed by these checks. diff --git a/docs/research.md b/docs/research.md index 7230823..c344018 100644 --- a/docs/research.md +++ b/docs/research.md @@ -8,6 +8,16 @@ decisions and a build order. ## 0. Executive summary — what changed vs v0.1 +The [managed native dispatch regression](reports/rds-managed-control-20261002.md) +separates the whole control future from receive/decode. Awaiting a codec inside +a selected receive handler stops polling the other branches until that handler +finishes, even when the codec uses `spawn_blocking`. Tokio documents +[retaining and concurrently polling futures](https://docs.rs/tokio/latest/tokio/macro.select.html) +and [blocking-work cancellation limits](https://docs.rs/tokio/latest/tokio/task/fn.spawn_blocking.html). +The correction retains serialized control writes, bounded codec permits and +encoded references. It does not establish native physical-pixel latency or +remove network outages, and incoming ACK observation still shares media IPC. + | Area | v0.1 assumption | Deep-research correction | | --- | --- | --- | | QUIC impl | "quinn via iroh" | iroh 1.2 runs on **noq** — a real fork with **QUIC Multipath + QNT + QAD merged**. Relay and direct are *simultaneous first-class paths* with per-path RTT/congestion, not magic-socket trickery. |